This is an automated email from the ASF dual-hosted git repository. krishvishal pushed a commit to branch sim-workload-faults in repository https://gitbox.apache.org/repos/asf/iggy.git
commit 72e58db678d348477516d46b621471173a0e2a42 Author: Krishna Vishal <[email protected]> AuthorDate: Fri Aug 14 14:40:41 2026 +0530 feat(simulator): resend timed-out client requests A client held one in-flight slot and nothing ever timed out, so the first dropped request or reply stranded that slot for the rest of the run. Any packet loss therefore wedged the workload rather than exercising recovery: at 2% loss a 5000-tick run drained four replies where a perfect network drained 439. It also forced the crash injector to spare primaries, since a request lost to a dead primary could never be answered. The workload now retains each outstanding request and resubmits it once the reply is overdue, rotating to the next replica so a client whose primary died finds the new one. The retained message is the original byte-for-byte, because the request id is what the metadata client table dedups on; rebuilding it would draw a fresh id and commit a second time instead of replaying the cached reply. Outstanding requests live in a BTreeMap, as the resulting submit order is observable and hash order would break replay. The two setup handshakes had the same defect and gave up after 100 and 200 steps; both now resubmit on the same reasoning, which is what let seed 1 fail during registration before the workload even started. A failed drain becomes a hard failure carrying which requests are outstanding, how many attempts each has had, and what every live replica believes about view, commit offset and primary. It was a warning only because a stall was previously expected and unactionable. --- core/simulator/src/bin/workload-fuzz.rs | 25 +++-- core/simulator/src/lib.rs | 155 +++++++++++++++++++++++------ core/simulator/src/workload/mod.rs | 171 +++++++++++++++++++++++++++----- core/simulator/src/workload/options.rs | 17 ++++ core/simulator/src/workload/oracle.rs | 53 +++++++++- 5 files changed, 350 insertions(+), 71 deletions(-) diff --git a/core/simulator/src/bin/workload-fuzz.rs b/core/simulator/src/bin/workload-fuzz.rs index 52fe0eeae..4ac434628 100644 --- a/core/simulator/src/bin/workload-fuzz.rs +++ b/core/simulator/src/bin/workload-fuzz.rs @@ -401,24 +401,29 @@ fn main() { ); if quiesce { - if oracle::drive_to_quiesce(&mut sim, &mut workload, 50_000) { - oracle::assert_converged(&sim, &workload); - println!("quiesced and converged (leader-relative + entity oracle)"); - } else { - println!( - "WARN: did not quiesce within budget — expected when crashing to bare quorum; \ - per-tick invariants still held" - ); - } + // A failed drain is a hard failure, not a warning. It used to be one + // because a lost request could not be retried, so a stall was expected + // and unactionable; with the client resending, a request that never gets + // answered inside the budget is either a wedge or a liveness bug, and + // the report says which replicas were live and what they believed. + assert!( + oracle::drive_to_quiesce(&mut sim, &mut workload, 50_000), + "{}", + oracle::quiesce_failure_report(&sim, &workload), + ); + oracle::assert_converged(&sim, &workload); + println!("quiesced and converged (leader-relative + entity oracle)"); } let stats = workload.auditor.stats(); println!( - "coverage: replies_seen={} replies_unknown={} committed_rejections={} samples_none={}", + "coverage: replies_seen={} replies_unknown={} committed_rejections={} \ + samples_none={} resends={}", stats.replies_seen, stats.replies_unknown, stats.committed_rejections, workload.samples_none(), + workload.resends(), ); for action in Action::iter() { let commits = stats.commits(action); diff --git a/core/simulator/src/lib.rs b/core/simulator/src/lib.rs index 2aafb34ee..3b8d621ae 100644 --- a/core/simulator/src/lib.rs +++ b/core/simulator/src/lib.rs @@ -67,6 +67,18 @@ const POLL_BUDGET: u32 = 100_000; /// keeps those draws from perturbing network or workload traces. pub const ENTRY_SHARD_SEED_SALT: u64 = 0x5A1A_F0E5_FACE_0003; +/// Steps a setup handshake (`register_client_with_primary`, `shell_login`) waits +/// before resubmitting its request, and the total it will spend. +/// +/// These helpers used to submit once and give up, which is sound only on a +/// perfect network: under injected packet loss the request or its reply is +/// eventually dropped and the fixture fails before the workload ever starts. +/// Both requests are metadata-plane ops carrying a stable request id, so the +/// client table dedups a resend and replays the cached reply rather than +/// committing twice. +const SETUP_RETRY_STEPS: u32 = 50; +const SETUP_TOTAL_STEPS: u32 = 4_000; + /// One simulated replica: its shards plus the executor bookkeeping needed /// to crash it. One entry per shard in `shards`/`pump_tasks` (a single /// shard until multi-shard lands). @@ -412,7 +424,8 @@ impl Simulator { /// (replica 0), as [`Self::register_client_with_primary`] does. /// /// # Panics - /// If no login reply arrives within 200 steps or it carries no session. + /// If no login reply arrives within [`SETUP_TOTAL_STEPS`], or it carries no + /// session. pub fn shell_login(&mut self, client: &SimClient) { // Register the client's connection metadata on every replica, as // `install_client_fd` does in production. `ensure_transport_connection` @@ -427,24 +440,46 @@ impl Simulator { )); } - let msg = client.login(replica::SHELL_ROOT_USERNAME, replica::SHELL_ROOT_PASSWORD); - self.submit_request(client.client_id(), 0, msg.into_generic()); - let mut session = 0u64; - let mut got_reply = false; - for _ in 0..200 { - if let Some(reply) = self.step().first() { - // The login reply carries the assigned session in `op` - // (`build_reply_with_body` maps the session field to `op`). - session = reply.header().op; - got_reply = true; - break; - } - } - assert!(got_reply, "shell_login: no login reply within 200 steps"); + let msg = client + .login(replica::SHELL_ROOT_USERNAME, replica::SHELL_ROOT_PASSWORD) + .into_generic(); + // The login reply carries the assigned session in `op` + // (`build_reply_with_body` maps the session field to `op`). + let session = self + .await_setup_reply(client.client_id(), 0, &msg, "shell_login") + .map_or(0, |reply| reply.header().op); assert!(session > 0, "shell_login: login reply carried no session"); client.bind_session(session); } + /// Submit `message` to `target` and step until a client reply arrives, + /// resubmitting every [`SETUP_RETRY_STEPS`] steps. + /// + /// Returns the first reply, or `None` once [`SETUP_TOTAL_STEPS`] is spent. + /// The same message is resubmitted verbatim, so the request id is stable and + /// the metadata client table treats a retry as a duplicate. + fn await_setup_reply( + &mut self, + client_id: u128, + target: u8, + message: &Message<GenericHeader>, + label: &str, + ) -> Option<Message<ReplyHeader>> { + for step in 0..SETUP_TOTAL_STEPS { + if step % SETUP_RETRY_STEPS == 0 { + self.submit_request(client_id, target, message.deep_copy()); + } + if let Some(reply) = self.step().into_iter().next() { + return Some(reply); + } + } + panic!( + "{label}: no reply for client {client_id} within {SETUP_TOTAL_STEPS} steps \ + (seed {:#x})", + self.seed, + ); + } + /// Advance the simulation by one tick. Returns client replies delivered. /// /// Every shard runs its real message pump as an executor task, so a @@ -633,17 +668,14 @@ impl Simulator { /// `SimClient`. /// /// # Panics - /// If no reply arrives within 100 steps. + /// If no reply arrives within [`SETUP_TOTAL_STEPS`]. #[allow(clippy::cast_possible_truncation)] pub fn register_client_with_primary(&mut self, client: &SimClient) { - let msg = client.register(); - self.submit_request(client.client_id(), 0, msg.into_generic()); - let mut session = 0u64; - let mut got_reply = false; - for _ in 0..100 { - let replies = self.step(); - if !replies.is_empty() { - let header = replies[0].header(); + let msg = client.register().into_generic(); + let session = self + .await_setup_reply(client.client_id(), 0, &msg, "register_client_with_primary") + .map_or(0, |reply| { + let header = reply.header(); debug_assert_eq!( header.operation, iggy_binary_protocol::Operation::Register, @@ -657,15 +689,8 @@ impl Simulator { client.client_id(), header.client, ); - session = header.commit; - got_reply = true; - break; - } - } - assert!( - got_reply, - "register_client_with_primary: no reply within 100 steps" - ); + header.commit + }); client.bind_session(session); // Partition has no `client_table`: at-least-once, no per-client @@ -2013,6 +2038,70 @@ mod tests { oracle::assert_converged(&sim, &wl); } + /// A lossy network drains, because the client resends. + /// + /// Without [`workload::Workload::due_resends`] this wedges immediately and + /// permanently: a client holds one in-flight slot, nothing times out, so the + /// first dropped request or reply strands that slot for the rest of the run. + /// At 5% loss a 3000-tick run used to drain a handful of replies and then + /// stop, and `drive_to_quiesce` could never finish because the reply it + /// waited on had already been discarded by the network. + /// + /// Asserts the resend path actually ran rather than the seed getting lucky. + #[test] + fn packet_loss_resends_and_drains() { + use crate::workload::{ + self, Workload, + options::{ActionWeights, WorkloadOptions}, + oracle, + }; + server_common::MemoryPool::init_pool(&server_common::MemoryPoolConfigOther { + enabled: false, + size: iggy_common::IggyByteSize::from(0u64), + bucket_capacity: 1, + }); + + let replica_count: u8 = 3; + let client_id: u128 = 1; + let seed = 0x105_5A1A; + let network_opts = packet::PacketSimulatorOptions { + node_count: replica_count, + client_count: 1, + seed, + packet_loss_probability: 0.05, + ..packet::PacketSimulatorOptions::default() + }; + let mut sim = Simulator::new( + usize::from(replica_count), + std::iter::once(client_id), + network_opts, + ); + let client = client::SimClient::new(client_id); + let ns_a = server_common::sharding::IggyNamespace::new(1, 1, 0); + sim.init_partition(ns_a); + sim.register_client_with_primary(&client); + + let mut options = WorkloadOptions::new(seed, replica_count, vec![ns_a]); + options.weights = ActionWeights::partition_only(); + let mut wl = Workload::new(options); + + let clients = [client]; + let replies = workload::run(&mut sim, &mut wl, &clients, 3_000, u64::MAX); + assert!(replies > 0, "lossy workload produced no replies"); + assert!( + wl.resends() > 0, + "no request timed out at 5% packet loss, so the resend path never ran; \ + raise the loss rate or lower request_timeout_ticks" + ); + + assert!( + oracle::drive_to_quiesce(&mut sim, &mut wl, 20_000), + "{}", + oracle::quiesce_failure_report(&sim, &wl), + ); + oracle::assert_converged(&sim, &wl); + } + /// Drive Create-heavy then Delete-heavy workload; assert shadow tracks /// live streams: /// diff --git a/core/simulator/src/workload/mod.rs b/core/simulator/src/workload/mod.rs index 2e641066c..4ac4735da 100644 --- a/core/simulator/src/workload/mod.rs +++ b/core/simulator/src/workload/mod.rs @@ -49,22 +49,49 @@ use rand_xoshiro::Xoshiro256Plus; use rand_xoshiro::rand_core::SeedableRng; use server_common::Message; use shadow::Shadow; -use std::collections::{HashMap, HashSet}; +use std::collections::{BTreeMap, HashSet}; /// Max in-flight requests per client. Must stay under the consensus /// pipeline's queue limits. pub const CLIENT_REQUEST_QUEUE_MAX: usize = 1; +/// An outstanding request, retained so the client can resend it. +/// +/// The encoded message is kept verbatim rather than rebuilt from the sampled +/// `Input`, because rebuilding would draw a fresh request id from the client and +/// a resend must reuse the original: that id is what the metadata plane's client +/// table dedups on, so a renumbered retry commits a second time instead of +/// returning the cached reply. +struct Outstanding { + message: Message<RoutedRequestHeader>, + /// Replica the most recent attempt went to. A resend moves to the next one, + /// so a client whose primary died eventually finds the new one. + target: u8, + /// Tick of the most recent attempt, not of the first. + attempted_tick: u64, + attempts: u32, +} + pub struct Workload { prng: Xoshiro256Plus, pub auditor: ServerAuditor, pub shadow: Shadow, pub options: WorkloadOptions, - /// Number of in-flight requests per client. + /// Outstanding requests keyed exactly as the auditor keys its expectations, + /// so the two are removed together. + /// + /// A `BTreeMap`, not a `HashMap`: [`Self::due_resends`] walks it and the + /// resulting submit order is observable, so hash iteration order would make + /// replay diverge from the seed. /// /// TODO: reap on client disconnect; bounded today by the fixed /// `Simulator::new` set. - in_flight_per_client: HashMap<u128, usize>, + outstanding: BTreeMap<(u128, u64), Outstanding>, + /// Driver tick, advanced by [`Self::tick`]. A driver that never ticks never + /// resends, which is what the hand-written scenario tests rely on. + now: u64, + /// Total resends issued, for the run summary. + resends: u64, /// Debug counter for `sample()` returning `None` (a targeted outcome whose /// shadow precondition is unmet). Flags PRNG-trace drift during development. samples_none: u64, @@ -90,7 +117,9 @@ impl Workload { auditor: ServerAuditor::new(), shadow, options, - in_flight_per_client: HashMap::new(), + outstanding: BTreeMap::new(), + now: 0, + resends: 0, samples_none: 0, strict_outcome_oracle, } @@ -99,18 +128,83 @@ impl Workload { /// True if the client has a free in-flight slot. #[must_use] pub fn client_idle(&self, client_id: u128) -> bool { - self.in_flight_per_client - .get(&client_id) - .copied() - .unwrap_or(0) - < CLIENT_REQUEST_QUEUE_MAX + self.client_in_flight(client_id) < CLIENT_REQUEST_QUEUE_MAX } /// Total in-flight requests across all clients. Read by the /// [`Invariants`]; draws no PRNG. #[must_use] pub(crate) fn total_in_flight(&self) -> usize { - self.in_flight_per_client.values().copied().sum() + self.outstanding.len() + } + + /// Advance the resend clock by one tick. Called once per driver iteration; + /// [`Self::due_resends`] measures against it. + pub fn tick(&mut self) { + self.now += 1; + } + + /// Total resends issued so far. + #[must_use] + pub const fn resends(&self) -> u64 { + self.resends + } + + /// Requests whose reply has not arrived within + /// [`WorkloadOptions::request_timeout_ticks`], each paired with the replica + /// to retry it against. Callers must submit every returned message. + /// + /// This is what a real client's read timeout does, and the harness needs it + /// for two reasons. A dropped request or reply otherwise strands the + /// client's only in-flight slot for the rest of the run, so any packet loss + /// wedges the workload. And a request lost to a crashed primary can only be + /// answered by the next one, which the client reaches by rotating its + /// target. + /// + /// Resending is safe on both planes but not equally cheap: the metadata + /// plane dedups on the retained request id and replays the cached reply, + /// while the partition plane is at-least-once and may commit the op twice. + /// The shadow already models that (`Effect` application is driven by what + /// committed, not by what was targeted). + #[must_use = "returned requests must be submitted or the client stays wedged"] + pub fn due_resends(&mut self) -> Vec<(u8, Message<RoutedRequestHeader>)> { + let timeout = self.options.request_timeout_ticks; + if timeout == 0 { + return Vec::new(); + } + let replica_count = self.options.replica_count.max(1); + let now = self.now; + let mut due = Vec::new(); + for entry in self.outstanding.values_mut() { + if now.saturating_sub(entry.attempted_tick) < timeout { + continue; + } + entry.target = (entry.target + 1) % replica_count; + entry.attempted_tick = now; + entry.attempts += 1; + due.push((entry.target, entry.message.deep_copy())); + } + self.resends += due.len() as u64; + due + } + + /// Outstanding requests as `(client, request, target, attempts)`, in key + /// order. Diagnostic only: names what a run was still waiting on when it + /// failed to drain. + #[must_use] + pub(crate) fn outstanding_summary(&self) -> Vec<(u128, u64, u8, u32)> { + self.outstanding + .iter() + .map(|(&(client, request), entry)| (client, request, entry.target, entry.attempts)) + .collect() + } + + /// In-flight count for one client. Keys are `(client, request)`, so the + /// client's entries are one contiguous range. + fn client_in_flight(&self, client_id: u128) -> usize { + self.outstanding + .range((client_id, 0)..=(client_id, u64::MAX)) + .count() } /// Aggregate in-flight ceiling: one queue's worth per declared client. @@ -173,10 +267,15 @@ impl Workload { request_namespace: header.group, }, ); - *self - .in_flight_per_client - .entry(client.client_id()) - .or_insert(0) += 1; + self.outstanding.insert( + key, + Outstanding { + message: message.deep_copy(), + target, + attempted_tick: self.now, + attempts: 1, + }, + ); Some((target, message)) } @@ -200,7 +299,7 @@ impl Workload { OnReply::Match(entry) => entry, OnReply::NsMismatch => { // Entry consumed; release slot, skip effects (misrouted). - self.decrement_in_flight(header.client); + self.release_outstanding(key); return Vec::new(); } OnReply::Unknown => return Vec::new(), @@ -283,24 +382,27 @@ impl Workload { self.auditor.note_committed(entry.action); } - self.decrement_in_flight(header.client); + self.release_outstanding(key); result.sim_commands } - /// Release one in-flight slot. Panics on underflow so a future - /// double-decrement surfaces instead of being silently clamped. + /// Drop a request's retry entry, freeing the client's slot. Paired with the + /// auditor consuming its expectation for the same key, so the two never + /// disagree about what is outstanding. /// /// # Panics - /// Panics if no entry exists for `client`, or if the counter is 0. - fn decrement_in_flight(&mut self, client: u128) { - let count = self - .in_flight_per_client - .get_mut(&client) - .expect("decrement_in_flight: no entry for client; record_in_flight must precede"); - *count = count - .checked_sub(1) - .expect("in_flight underflow: per-client counter went below 0"); + /// Panics if no entry exists for `key`; the auditor only reports a match or + /// a namespace mismatch for a key it was given, and `build_request` records + /// both sides together, so a miss here means the two drifted. + fn release_outstanding(&mut self, key: (u128, u64)) { + assert!( + self.outstanding.remove(&key).is_some(), + "no outstanding entry for (client={}, request={}); the auditor \ + matched a key the retry buffer never recorded", + key.0, + key.1, + ); } /// Debug counter for `sample()` returning `None`. Surfaces sampling @@ -374,9 +476,13 @@ pub fn run( let mut fault_prng = Xoshiro256Plus::seed_from_u64(workload.options.seed ^ FAULT_SEED_SALT); let mut replies_seen = 0u64; for _ in 0..tick_budget { + workload.tick(); if workload.options.crash_per_tick_ratio > 0.0 { maybe_inject_crash(sim, workload, &mut fault_prng); } + // Resend before sampling: a timed-out request still holds the client's + // slot, so `build_request` would decline it anyway. + resubmit_due(sim, workload); for client in clients { if let Some((target, msg)) = workload.build_request(client) { sim.submit_request(client.client_id(), target, msg.into_generic()); @@ -439,6 +545,17 @@ fn maybe_inject_crash(sim: &mut Simulator, workload: &Workload, prng: &mut Xoshi sim.replica_crash(victim); } +/// Submit every request whose reply is overdue (see [`Workload::due_resends`]). +/// +/// The client id rides the retained message's header, so a resend re-enters the +/// network exactly as the original did, only aimed at the next replica. +pub fn resubmit_due(sim: &mut Simulator, workload: &mut Workload) { + for (target, message) in workload.due_resends() { + let client_id = message.header().client; + sim.submit_request(client_id, target, message.into_generic()); + } +} + /// Apply `SimCommand`s returned by [`Workload::on_reply`]. /// /// Callers must invoke this (or [`run`]) for every batch of returned diff --git a/core/simulator/src/workload/options.rs b/core/simulator/src/workload/options.rs index a4a29447a..d707b2961 100644 --- a/core/simulator/src/workload/options.rs +++ b/core/simulator/src/workload/options.rs @@ -19,6 +19,14 @@ use crate::workload::actions::Action; use server_common::sharding::IggyNamespace; use strum::EnumCount; +/// Default [`WorkloadOptions::request_timeout_ticks`]. +/// +/// Four times the primary's commit-broadcast interval (50 ticks, see +/// `oracle::QUIESCE_SETTLE_TICKS`), so a request whose commit is merely waiting +/// on the next broadcast is never resent, while one lost to a dropped packet or +/// a crashed primary is retried long before the run's budget runs out. +pub const DEFAULT_REQUEST_TIMEOUT_TICKS: u64 = 200; + /// Per-action sampling weights as percentages. Unlisted variants default /// to 0 (never picked). Listed weights must sum to 100. #[derive(Debug, Clone, Copy)] @@ -171,6 +179,14 @@ pub struct WorkloadOptions { /// Floor on live replicas the driver will not crash below, preserving a /// commit quorum. Defaults to `replica_count / 2 + 1`. pub min_survivors: u8, + /// Ticks a request may stay outstanding before the client resends it. + /// + /// Must clear the primary's commit-broadcast interval with room to spare, or + /// a client resends work that was about to be answered and the run spends + /// its budget on duplicates. Must also stay well under the tick budget, or a + /// request lost to a dropped packet never gets retried and the client's slot + /// strands. `0` disables resending. + pub request_timeout_ticks: u64, } impl WorkloadOptions { @@ -191,6 +207,7 @@ impl WorkloadOptions { max_offset: 1_000_000, crash_per_tick_ratio: 0.0, min_survivors: replica_count / 2 + 1, + request_timeout_ticks: DEFAULT_REQUEST_TIMEOUT_TICKS, } } } diff --git a/core/simulator/src/workload/oracle.rs b/core/simulator/src/workload/oracle.rs index b5a2b6ce0..4bdfc50f5 100644 --- a/core/simulator/src/workload/oracle.rs +++ b/core/simulator/src/workload/oracle.rs @@ -38,7 +38,7 @@ use crate::Simulator; use crate::replica::Replica; use crate::workload::shadow::Shadow; -use crate::workload::{Workload, apply_sim_commands}; +use crate::workload::{Workload, apply_sim_commands, resubmit_due}; use consensus::{MetadataHandle, Status}; use metadata::impls::metadata::StreamsFrontend; use std::collections::BTreeSet; @@ -96,6 +96,11 @@ impl CommittedMetadata { pub fn drive_to_quiesce(sim: &mut Simulator, workload: &mut Workload, max_ticks: u64) -> bool { let mut drained = false; for _ in 0..max_ticks { + // The drain keeps resending: a request lost on the way out is never + // answered, so without retries the drain would spend its whole budget + // waiting on a reply that cannot arrive. + workload.tick(); + resubmit_due(sim, workload); for reply in sim.step() { let cmds = workload.on_reply(&reply); apply_sim_commands(sim, &cmds); @@ -117,6 +122,52 @@ pub fn drive_to_quiesce(sim: &mut Simulator, workload: &mut Workload, max_ticks: true } +/// Why the run did not drain, as a multi-line report. +/// +/// A failed drain is either a wedge or a cluster that is merely slow, and the +/// bare boolean [`drive_to_quiesce`] returns cannot tell them apart. Once +/// crashes, restarts and packet loss are all in play, that distinction is the +/// whole diagnosis, so name what is still outstanding and what every live +/// replica thinks the world looks like. +#[must_use] +pub fn quiesce_failure_report(sim: &Simulator, workload: &Workload) -> String { + use std::fmt::Write; + + let mut report = format!( + "did not drain: {} request(s) still outstanding (seed={:#x})\n", + workload.total_in_flight(), + workload.options.seed, + ); + for (client, request, target, attempts) in workload.outstanding_summary() { + let _ = writeln!( + report, + " outstanding client={client} request={request} \ + last_target=replica {target} attempts={attempts}", + ); + } + let _ = writeln!(report, " resends issued: {}", workload.resends()); + for replica_idx in 0..sim.replica_count { + if sim.is_crashed(replica_idx) { + let _ = writeln!(report, " replica {replica_idx}: CRASHED"); + continue; + } + let _ = write!(report, " replica {replica_idx}: live"); + for &ns in &workload.options.namespaces { + let view = sim.consensus_view(usize::from(replica_idx), ns); + let commit = sim + .offsets(usize::from(replica_idx), ns) + .map(|offsets| offsets.commit_offset); + let primary = sim.primary_index(ns); + let _ = write!( + report, + " | ns {ns:?} view={view:?} commit_offset={commit:?} primary={primary:?}", + ); + } + report.push('\n'); + } + report +} + /// Post-drain consensus checks that hold today. /// /// Asserts no live replica is ahead of the leader, and (on a serial run) that
