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 0593626672b0d142c41bc7aceced69bef2f22aa5 Author: Krishna Vishal <[email protected]> AuthorDate: Sat Aug 15 01:33:26 2026 +0530 feat(simulator): drive the workload through the dispatch shell The shell wires the server's real dispatch handlers, and it was reachable only from two hand-written tests. Authorization, session binding and the pre-commit deny replies exist nowhere else -- the raw `on_message` path has no deny site at all -- so none of them had ever seen a workload. Running one immediately found that the workload could not read a denial. `ReplyHeader::status` nonzero means a request refused before commit, and its body is EMPTY by contract, mutually exclusive with the metadata result section. The reply decoder read a result section unconditionally, found no bytes, and reported the reply as corrupt. Denials are now their own outcome: the op never entered the log, so the shadow does not move and there is nothing to classify. Denials are counted per action, not just totalled, because the two causes need telling apart. An op the server refuses for this particular input looks like a mix of commits and denials; an op whose wire shape the dispatch layer cannot decode is denied every single time. The consumer offset ops are the first kind (`ResourceNotFound`, offsets and consumers the workload has not created yet). Both PAT ops are the second: zero commits and `InvalidCommand` on every attempt, so `SimClient` builds a shape the real path rejects, the same class of bug as the consumer kind fixed earlier. Left for the fix pass. The harness also assumed every client-addressed frame is a reply, so an `Eviction` -- a command nothing had exercised before -- failed to decode and surfaced as "wrong command type". Evictions are now classified and recorded. The driver refuses to continue on one, since a lost session makes every outstanding request unanswerable and modelling recovery means re-establishing the session mid-run and discarding the auditor's expectations for that client. Naming it beats presenting as a stall. Shell runs drain and converge 12 of 12 across four op mixes and three seeds. Shell with crashes and restarts hits the eviction limitation instead, which is now a clear message rather than a decode panic. The raw path is unchanged: 12 of 12 with crash and restart injection. --- core/simulator/src/bin/workload-fuzz.rs | 51 +++++++++++---- core/simulator/src/lib.rs | 112 +++++++++++++++++++++++++++++++- core/simulator/src/workload/auditor.rs | 26 ++++++++ core/simulator/src/workload/mod.rs | 44 +++++++++++++ 4 files changed, 220 insertions(+), 13 deletions(-) diff --git a/core/simulator/src/bin/workload-fuzz.rs b/core/simulator/src/bin/workload-fuzz.rs index c41d5be9f..fd918e483 100644 --- a/core/simulator/src/bin/workload-fuzz.rs +++ b/core/simulator/src/bin/workload-fuzz.rs @@ -82,6 +82,12 @@ struct Args { /// Crash the primary too, putting a view change under live traffic. #[arg(long)] crash_primary: bool, + /// Route every client request through the server's real dispatch handlers + /// instead of the raw `on_message` fast path. Clients then log in against the + /// seeded root user and carry a bound session, so the run also covers + /// authorization and session lifecycle, which exist only on this path. + #[arg(long)] + shell: bool, #[arg(long)] no_quiesce: bool, @@ -356,8 +362,8 @@ fn main() { let network_opts = network_options(&args, replicas, clients, seed); println!( "workload-fuzz: seed={seed} ticks={ticks} clients={clients} replicas={replicas} \ - plane={plane:?} faults={:?} crash_prob={crash_prob} quiesce={quiesce}", - args.faults, + plane={plane:?} faults={:?} shell={} crash_prob={crash_prob} quiesce={quiesce}", + args.faults, args.shell, ); println!( "network: loss={} replay={} delay={}..{} partition={:?}/{:?} \ @@ -383,17 +389,38 @@ fn main() { }); let client_ids: Vec<u128> = (1..=u128::from(clients)).collect(); - let mut sim = Simulator::new( - usize::from(replicas), - client_ids.iter().copied(), - network_opts, - ); + let mut sim = if args.shell { + Simulator::with_shards_shell( + usize::from(replicas), + 1, + client_ids.iter().copied(), + network_opts, + ) + } else { + Simulator::new( + usize::from(replicas), + client_ids.iter().copied(), + network_opts, + ) + }; let sim_clients: Vec<SimClient> = client_ids.iter().map(|&id| SimClient::new(id)).collect(); let ns = IggyNamespace::new(1, 1, 0); sim.init_partition(ns); + if args.shell { + // The shell resolves a partition request's namespace against committed + // metadata, so the stream and topic behind it have to exist as well as the + // partition group. + sim.seed_stream_topic_partition(ns); + } for client in &sim_clients { - sim.register_client_with_primary(client); + if args.shell { + // Log in rather than bare-register: the dispatch path admits a request + // only from a bound session, and the login is what mints one. + sim.shell_login(client); + } else { + sim.register_client_with_primary(client); + } } let mut options = WorkloadOptions::new(seed, replicas, vec![ns]); @@ -461,17 +488,19 @@ fn print_coverage(workload: &Workload) { let stats = workload.auditor.stats(); println!( "coverage: replies_seen={} replies_unknown={} committed_rejections={} \ - samples_none={} resends={}", + samples_none={} resends={} denials={}", stats.replies_seen, stats.replies_unknown, stats.committed_rejections, workload.samples_none(), workload.resends(), + stats.denials, ); for action in Action::iter() { let commits = stats.commits(action); - if commits > 0 { - println!(" {action:?}: {commits} commits"); + let (denied, status) = stats.denials_per_action[action as usize]; + if commits > 0 || denied > 0 { + println!(" {action:?}: {commits} commits, {denied} denied (last status {status})"); } } } diff --git a/core/simulator/src/lib.rs b/core/simulator/src/lib.rs index 6223a4f42..6f0ecef54 100644 --- a/core/simulator/src/lib.rs +++ b/core/simulator/src/lib.rs @@ -32,7 +32,7 @@ use deps::SimClock; use deps::SimSuperblock; use deps::{MemStorage, SimJournal}; use executor::{DetExecutor, RunOutcome, TaskId}; -use iggy_binary_protocol::{GenericHeader, ReplyHeader}; +use iggy_binary_protocol::{Command2, GenericHeader, ReplyHeader}; use iggy_common::IggyError; use message_bus::installer::conn_info::{ClientConnMeta, ClientTransportKind}; use metadata::impls::metadata::StreamsFrontend; @@ -176,6 +176,14 @@ pub struct Simulator { entry_rng: Xoshiro256Plus, /// Network seed, kept for livelock diagnostics. seed: u64, + /// Clients the cluster has evicted since the last drain, in delivery order. + /// + /// An eviction ends a session: every request that client has outstanding will + /// never be answered, and it must log in again before it can submit anything. + /// Recorded rather than acted on here, because re-establishing a session needs + /// to step the simulator and drop the workload's expectations for that client, + /// neither of which belongs inside packet delivery. + evicted: Vec<u128>, /// Dispatch-shell mode: when set, inbound client packets are delivered /// through the real `on_client_request` handler (see /// [`shard::IggyShard::deliver_client_request`]) instead of the raw @@ -388,6 +396,7 @@ impl Simulator { client_ids, executor, entry_rng: Xoshiro256Plus::seed_from_u64(seed ^ ENTRY_SHARD_SEED_SALT), + evicted: Vec::new(), seed, shell, } @@ -564,7 +573,17 @@ impl Simulator { } // Crashed or missing: packet silently dropped. } - ProcessId::Client(_) => { + ProcessId::Client(client_id) => { + // Not every client-addressed frame is a reply. `Eviction` + // tells a client its session is gone, which the server sends + // once the client table drops it -- reachable as soon as the + // dispatch shell runs with crashes and restarts. Decoding it + // as a reply fails on the command discriminant, so classify + // first and record the eviction for the driver. + if packet.message.header().command == Command2::Eviction { + self.evicted.push(client_id); + continue; + } let reply: Message<ReplyHeader> = packet .message .deep_copy() @@ -609,6 +628,15 @@ impl Simulator { client_replies } + /// Take the clients evicted since the last call. + /// + /// A driver must consume these: an evicted client's session is gone, so its + /// outstanding requests are unanswerable and its next request is refused + /// until it logs in again. Ignoring them looks exactly like a wedge. + pub fn take_evictions(&mut self) -> Vec<u128> { + std::mem::take(&mut self.evicted) + } + /// Rolling hash of the executor schedule (every poll and timer fire). /// Two runs from the same seed and inputs must agree; determinism /// tests assert on it alongside the reply-trace hash. @@ -3244,6 +3272,86 @@ mod tests { /// pipeline to `PIPELINE_PREPARE_QUEUE_MAX`; a request on `ns_b` /// still commits while `ns_a` is wedged (no quorum without backup /// acks); lifting the block drains `ns_a` completely. + /// The workload drains and converges when every request goes through the + /// server's real dispatch handlers rather than the raw `on_message` path. + /// + /// Worth its own test because the shell path is where authorization, session + /// binding and the pre-commit deny replies live; the raw path has no deny site + /// at all. Running the workload here is the only thing that exercises them, + /// and it is what surfaced that the workload had never modelled a denial: + /// `ReplyHeader::status` nonzero means an empty body, and the reply decoder + /// was reading a result section off it and calling the reply corrupt. + /// + /// Asserts denials were actually observed, so the test cannot pass by taking + /// a path where nothing is ever denied and thus proving nothing about the + /// shell. + #[test] + fn shell_workload_drains_and_converges() { + 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 = 0x5E11_0001; + let network_opts = packet::PacketSimulatorOptions { + node_count: replica_count, + client_count: 1, + seed, + ..packet::PacketSimulatorOptions::default() + }; + let mut sim = Simulator::with_shards_shell( + usize::from(replica_count), + 1, + std::iter::once(client_id), + network_opts, + ); + let ns = IggyNamespace::new(1, 1, 0); + sim.init_partition(ns); + // The dispatch path resolves a partition request's namespace against + // committed metadata, so the stream and topic have to exist too. + sim.seed_stream_topic_partition(ns); + + let client = SimClient::new(client_id); + // Log in rather than bare-register: dispatch admits a request only from a + // bound session. + sim.shell_login(&client); + + let mut options = WorkloadOptions::new(seed, replica_count, vec![ns]); + options.weights = ActionWeights::uniform(); + let mut wl = Workload::new(options); + + let clients = [client]; + let replies = workload::run(&mut sim, &mut wl, &clients, 4_000, u64::MAX); + assert!(replies > 0, "shell workload produced no replies"); + + let stats = wl.auditor.stats(); + assert!( + stats.commits_per_action.iter().sum::<u64>() > 0, + "shell workload committed nothing, so the dispatch path never got past \ + admission" + ); + assert!( + stats.denials > 0, + "no request was denied, so the pre-commit deny path this test exists to \ + cover never ran" + ); + + assert!( + oracle::drive_to_quiesce(&mut sim, &mut wl, 20_000), + "{}", + oracle::quiesce_failure_report(&sim, &wl), + ); + oracle::assert_converged(&sim, &wl); + } + /// The cross-replica equality check actually compares replicas against each /// other, and holds over a metadata workload with crashes and restarts. /// diff --git a/core/simulator/src/workload/auditor.rs b/core/simulator/src/workload/auditor.rs index 411a83af7..903117935 100644 --- a/core/simulator/src/workload/auditor.rs +++ b/core/simulator/src/workload/auditor.rs @@ -58,6 +58,22 @@ pub struct AuditorStats { /// rejection). The shadow does not mutate on these; in a serial run the /// `on_reply` equality oracle asserts the rejection was the targeted outcome. pub committed_rejections: u64, + /// Replies denied before commit, carrying `ReplyHeader::status` and an empty + /// body. Distinct from a committed rejection: the op never entered the log, + /// so the shadow must not move and no result section exists to classify. + /// Only the dispatch shell produces these (authorization runs there); the raw + /// path never denies. + pub denials: u64, + /// Per-action denial counter and the last status seen for it, indexed by + /// `Action as usize`. + /// + /// Per-action rather than a single total because the two causes need telling + /// apart: an op the server legitimately refuses for this input (an offset the + /// partition cannot accept yet) versus an op the dispatch layer cannot decode + /// at all, which means the workload builds a wire shape the real path rejects. + /// The second is a workload bug and shows up as every request for that action + /// being denied with the same status. + pub denials_per_action: [(u64, u32); Action::COUNT], } impl Default for AuditorStats { @@ -67,6 +83,8 @@ impl Default for AuditorStats { replies_unknown: 0, commits_per_action: [0u64; Action::COUNT], committed_rejections: 0, + denials: 0, + denials_per_action: [(0, 0); Action::COUNT], } } } @@ -174,6 +192,14 @@ impl ServerAuditor { /// Record a committed business rejection (nonzero result code). Either /// targeted by outcome-first generation (duplicate name, fabricated missing /// entity) or produced by a race. + /// Record a pre-commit denial (`ReplyHeader::status` nonzero). + pub const fn note_denial(&mut self, action: Action, status: u32) { + self.stats.denials += 1; + let entry = &mut self.stats.denials_per_action[action as usize]; + entry.0 += 1; + entry.1 = status; + } + pub const fn note_committed_rejection(&mut self) { self.stats.committed_rejections += 1; } diff --git a/core/simulator/src/workload/mod.rs b/core/simulator/src/workload/mod.rs index 164367084..564082bfd 100644 --- a/core/simulator/src/workload/mod.rs +++ b/core/simulator/src/workload/mod.rs @@ -315,6 +315,23 @@ impl Workload { OnReply::Unknown => return Vec::new(), }; + // A pre-commit denial short-circuits everything below. `ReplyHeader`'s + // contract makes the two channels mutually exclusive: a reply either + // commits (status 0, result section present) or is denied before commit + // (status set, EMPTY body). Reading a result section off a denial finds + // no bytes, which the metadata branch below would report as a corrupt + // reply. + // + // The op never entered the log, so the shadow must not move and there is + // nothing to classify. Only the dispatch shell produces these, since + // authorization runs there; the raw path has no denial site at all, which + // is why this went unmodelled until the workload ran through the shell. + if header.status != 0 { + self.auditor.note_denial(entry.action, header.status); + self.release_outstanding(key); + return Vec::new(); + } + // Decode the committed result code. Metadata replies carry a // result section (see `ApplyReply::to_reply_body`); partition-plane // replies do not, hence the `is_metadata` gate. @@ -524,6 +541,7 @@ pub fn run_with_faults( apply_sim_commands(sim, &cmds); replies_seen += 1; } + assert_no_evictions(sim); invariants.check(sim, workload); if replies_seen >= replies_target { break; @@ -669,6 +687,32 @@ impl FaultInjector { } } +/// Fail loudly if the cluster evicted a client, which the workload cannot yet +/// survive. +/// +/// An eviction ends the session: the client's outstanding requests become +/// unanswerable and it must log in again before submitting anything. Modelling +/// that means re-establishing the session mid-run and discarding the auditor's +/// expectations for it, which the driver does not do. Until it does, an eviction +/// presents as a client that has silently stopped making progress, so name it +/// here rather than let the run time out with no explanation. +/// +/// Only reachable through the dispatch shell, and in practice only once crashes +/// and restarts are also in play. +/// +/// # Panics +/// If any client was evicted since the last step. +fn assert_no_evictions(sim: &mut Simulator) { + let evicted = sim.take_evictions(); + assert!( + evicted.is_empty(), + "cluster evicted client(s) {evicted:?}: their sessions are gone, so their \ + outstanding requests can never be answered and their next request is \ + refused. The workload does not re-establish a session, so the run cannot \ + continue" + ); +} + /// 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
