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 deb2e974e2fb0b901653172f38f929437b2e57ff Author: Krishna Vishal <[email protected]> AuthorDate: Sat Aug 15 07:33:30 2026 +0530 feat(simulator): cover register and logout forwarding, and measure command reach Register and logout forwarding landed with no deterministic coverage and could not have had any: forwarding fires only when a client dials a BACKUP, and the harness only ever dialed the primary. The backup verifies the credentials itself, sends the consensus proposal on as `ForwardRegister`, parks the login until `ForwardRegisterResult` returns, then answers on the connection it owns; logout works the same way. All four frames are now exercised by a login and logout issued against a backup. The test asserts the frames were DELIVERED rather than inferring them from a working session, because dialing a backup would also appear to work if the server had silently redirected the client, which is the design forwarding replaced. Pointing the same test at the primary fails on that assertion, so it cannot pass vacuously. Proving that needed a way to observe the wire, so the packet simulator now counts deliveries per command. That answers a question the harness could not previously answer about itself: which of the protocol it actually reaches. A run with crashes and restarts delivers ten commands and notably includes RequestStartView, RequestPrepares, RepairPrepare and RepairDone, so restart injection does drive rejoin and log repair rather than merely being configured to. Adding `--crash-primary` adds StartViewChange and DoViewChange. Nineteen commands are still never delivered, among them RangeEvicted and the four state-transfer frames, which is the next task's target. Logging out needed a `SimClient::logout`, which the client had never had because nothing tore a session down. --- core/simulator/src/bin/workload-fuzz.rs | 28 +++++++- core/simulator/src/client.rs | 12 ++++ core/simulator/src/lib.rs | 116 +++++++++++++++++++++++++++++++- core/simulator/src/network.rs | 18 ++++- core/simulator/src/packet.rs | 71 +++++++++++++++++++ 5 files changed, 241 insertions(+), 4 deletions(-) diff --git a/core/simulator/src/bin/workload-fuzz.rs b/core/simulator/src/bin/workload-fuzz.rs index fd918e483..940120387 100644 --- a/core/simulator/src/bin/workload-fuzz.rs +++ b/core/simulator/src/bin/workload-fuzz.rs @@ -47,7 +47,7 @@ use server_common::sharding::IggyNamespace; use server_common::{MemoryPool, MemoryPoolConfigOther}; use simulator::Simulator; use simulator::client::SimClient; -use simulator::packet::{PacketSimulatorOptions, PartitionMode, PartitionSymmetry}; +use simulator::packet::{COMMAND_LABELS, PacketSimulatorOptions, PartitionMode, PartitionSymmetry}; use simulator::workload::actions::Action; use simulator::workload::options::{ActionWeights, WorkloadOptions}; use simulator::workload::{FaultInjector, Workload, oracle, run_with_faults}; @@ -480,9 +480,35 @@ fn main() { print_coverage(&workload); } + print_command_coverage(&sim); println!("workload-fuzz: OK (seed={seed})"); } +/// Which protocol commands the run actually delivered, and which it never +/// reached. +/// +/// The harness wires far more of the command space than any one scenario drives, +/// and "is this path covered?" was previously answered by grepping the source. +/// Counted at delivery, so a command listed here really arrived somewhere. +fn print_command_coverage(sim: &Simulator) { + let counts = sim.network.command_counts(); + let mut seen: Vec<String> = Vec::new(); + let mut unseen: Vec<&str> = Vec::new(); + for (discriminant, &count) in counts.iter().enumerate() { + let label = COMMAND_LABELS[discriminant]; + if label == "Reserved" { + continue; + } + if count > 0 { + seen.push(format!("{label}={count}")); + } else { + unseen.push(label); + } + } + println!("commands delivered: {}", seen.join(" ")); + println!("commands never delivered: {}", unseen.join(" ")); +} + /// Reply, rejection and resend counters plus per-action commits. fn print_coverage(workload: &Workload) { let stats = workload.auditor.stats(); diff --git a/core/simulator/src/client.rs b/core/simulator/src/client.rs index 02ea37323..88e0d9d08 100644 --- a/core/simulator/src/client.rs +++ b/core/simulator/src/client.rs @@ -231,6 +231,18 @@ impl SimClient { .expect("login request must be valid") } + /// Tear down this client's bound session. + /// + /// Replicates through the metadata plane like any other session op, so it + /// carries the bound session and a metadata request id and needs no body. A + /// logout issued to a BACKUP is what produces `ForwardLogout`: the backup owns + /// the connection but not the log, so it asks the primary to commit the + /// teardown and answers the client itself once `ForwardLogoutResult` returns. + #[must_use] + pub fn logout(&self) -> Message<RoutedRequestHeader> { + self.build_request(Operation::Logout, &[]) + } + /// # Panics /// Panics if the stream name is not a valid wire name. pub fn create_stream(&self, name: &str) -> Message<RoutedRequestHeader> { diff --git a/core/simulator/src/lib.rs b/core/simulator/src/lib.rs index 6f0ecef54..9ad034d00 100644 --- a/core/simulator/src/lib.rs +++ b/core/simulator/src/lib.rs @@ -462,6 +462,23 @@ impl Simulator { /// If no login reply arrives within [`SETUP_TOTAL_STEPS`], or it carries no /// session. pub fn shell_login(&mut self, client: &SimClient) { + self.shell_login_via(client, 0); + } + + /// [`Self::shell_login`] against a chosen replica. + /// + /// Dialing a BACKUP takes a different path than dialing the primary, and it + /// is the only way to reach register forwarding: the backup verifies the + /// credentials itself and sends only the consensus proposal onward as + /// `ForwardRegister`, parking the login until the matching + /// `ForwardRegisterResult` comes back, then answering on the connection it + /// owns. A client that always dials the primary never produces one of those + /// four frames. + /// + /// # Panics + /// If no login reply arrives within [`SETUP_TOTAL_STEPS`], or it carries no + /// session. + pub fn shell_login_via(&mut self, client: &SimClient, target: u8) { // Register the client's connection metadata on every replica, as // `install_client_fd` does in production. `ensure_transport_connection` // reads it to admit the connection into the SessionManager, which the @@ -481,7 +498,7 @@ impl Simulator { // 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") + .await_setup_reply(client.client_id(), target, &msg, "shell_login") .map_or(0, |reply| reply.header().op); assert!(session > 0, "shell_login: login reply carried no session"); client.bind_session(session); @@ -3272,6 +3289,103 @@ 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. + /// A client that dials a BACKUP still gets a working session, and the login + /// travels as a forwarded consensus proposal rather than a redirect. + /// + /// Register forwarding exists so a client need not find the primary itself: + /// the backup verifies the credentials locally, sends only the proposal on as + /// `ForwardRegister`, parks the login until the matching + /// `ForwardRegisterResult` returns, then answers on the connection it owns. + /// The whole subsystem landed with no deterministic coverage, and could not + /// have any while the harness only ever dialed the primary. + /// + /// Asserts the four frames were delivered rather than inferring them from a + /// working session: dialing a backup would also "work" if the server had + /// silently redirected the client instead, which is the design this replaced. + #[test] + fn login_via_backup_forwards_the_register_to_the_primary() { + use iggy_binary_protocol::Command2; + + 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 network_opts = packet::PacketSimulatorOptions { + node_count: replica_count, + client_count: 1, + seed: 0xF0_2D_0001, + ..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); + sim.seed_stream_topic_partition(ns); + + // Replica 0 leads both planes at view 0, so replica 1 is a backup and the + // login has to be forwarded. + let client = SimClient::new(client_id); + sim.shell_login_via(&client, 1); + + assert!( + sim.network.delivered_any(Command2::ForwardRegister), + "no ForwardRegister crossed the wire: the backup answered the login \ + itself, so this covers nothing" + ); + assert!( + sim.network.delivered_any(Command2::ForwardRegisterResult), + "the forwarded register was never answered, so the login below \ + succeeded by some other route" + ); + + // Log out on the SAME backup, which both covers the other half of the + // forwarding subsystem and proves the session was real: only a bound + // session can be torn down, and the teardown replicates through the + // primary exactly as the register did. + // + // A logout, not a data request, because the session belongs to the + // connection: a real client holds one connection to one replica, and a + // partition write is refused by a backup for routing reasons + // (`TransientNotAccepted`) which says nothing about the session. + let logout = client.logout(); + let request = logout.header().request; + sim.submit_request(client_id, 1, logout.into_generic()); + let mut answered = false; + for _ in 0..400 { + if let Some(reply) = sim + .step() + .into_iter() + .find(|reply| reply.header().request == request) + { + assert_eq!( + reply.header().status, + 0, + "logout on the forwarded session was refused (status {})", + reply.header().status, + ); + answered = true; + break; + } + } + assert!(answered, "no reply to the logout issued on the backup"); + assert!( + sim.network.delivered_any(Command2::ForwardLogout), + "the backup committed the logout without asking the primary" + ); + assert!( + sim.network.delivered_any(Command2::ForwardLogoutResult), + "the forwarded logout was never answered" + ); + } + /// The workload drains and converges when every request goes through the /// server's real dispatch handlers rather than the raw `on_message` path. /// diff --git a/core/simulator/src/network.rs b/core/simulator/src/network.rs index 6a4aa795e..abe7eeb66 100644 --- a/core/simulator/src/network.rs +++ b/core/simulator/src/network.rs @@ -22,9 +22,10 @@ //! process-to-bus routing, and node enable/disable logic. use crate::packet::{ - ALLOW_ALL, BLOCK_ALL, LinkFilter, Packet, PacketSimulator, PacketSimulatorOptions, ProcessId, + ALLOW_ALL, BLOCK_ALL, COMMAND_COUNT_MAX, LinkFilter, Packet, PacketSimulator, + PacketSimulatorOptions, ProcessId, }; -use iggy_binary_protocol::GenericHeader; +use iggy_binary_protocol::{Command2, GenericHeader}; use server_common::Message; /// Network layer for the cluster simulation. @@ -164,4 +165,17 @@ impl Network { pub fn packets_in_flight(&self) -> usize { self.simulator.packets_in_flight() } + + /// Packets delivered so far, per [`Command2`](iggy_binary_protocol::Command2) + /// discriminant. Names come from [`COMMAND_LABELS`]. + #[must_use] + pub const fn command_counts(&self) -> &[u64; COMMAND_COUNT_MAX] { + self.simulator.command_counts() + } + + /// Whether any packet of this command reached its destination. + #[must_use] + pub const fn delivered_any(&self, command: Command2) -> bool { + self.simulator.delivered_any(command) + } } diff --git a/core/simulator/src/packet.rs b/core/simulator/src/packet.rs index 093615aa9..c01016c68 100644 --- a/core/simulator/src/packet.rs +++ b/core/simulator/src/packet.rs @@ -241,8 +241,62 @@ pub struct PacketSimulator { auto_partition_nodes: Vec<usize>, /// Reusable buffer for delivered packets. delivered: Vec<Packet>, + /// Packets actually delivered, per [`Command2`] discriminant. + /// + /// Counted at delivery, past every drop path, so a command shows up only if a + /// process really received one. This is how a run answers which parts of the + /// protocol it exercised: the harness reaches far more of the command space + /// than any single scenario drives, and a test asserting a command was + /// observed is the difference between covering a path and merely compiling it. + command_counts: [u64; COMMAND_COUNT_MAX], } +/// One past the highest [`Command2`] discriminant, sizing [`COMMAND_LABELS`] and +/// the delivery counters. Raising it is part of adding a command. +pub const COMMAND_COUNT_MAX: usize = 30; + +/// Names for each [`Command2`] discriminant, so a coverage report reads as +/// protocol rather than as integers. Indexed by discriminant; the trailing +/// assert keeps it aligned with the enum. +pub const COMMAND_LABELS: [&str; COMMAND_COUNT_MAX] = [ + "Reserved", + "Ping", + "Pong", + "PingClient", + "PongClient", + "Request", + "Prepare", + "PrepareOk", + "Reply", + "Commit", + "StartViewChange", + "DoViewChange", + "StartView", + "Eviction", + "ReplicaHello", + "ReplicaChallenge", + "ReplicaFinish", + "RequestStartView", + "RequestPrepares", + "RepairPrepare", + "RepairDone", + "RangeEvicted", + "RequestStateTransfer", + "StateTransferTarget", + "RequestStateChunk", + "StateChunk", + "ForwardRegister", + "ForwardRegisterResult", + "ForwardLogout", + "ForwardLogoutResult", +]; + +const _: () = { + // Adding a command without extending the table above would silently report it + // under the wrong name, or index past the end. + assert!(Command2::ForwardLogoutResult as usize == COMMAND_COUNT_MAX - 1); +}; + impl PacketSimulator { /// Create a new packet simulator. /// @@ -322,6 +376,7 @@ impl PacketSimulator { auto_partition_stability: initial_stability, auto_partition_nodes: (0..node_count).collect(), delivered: Vec::new(), + command_counts: [0; COMMAND_COUNT_MAX], } } @@ -524,6 +579,7 @@ impl PacketSimulator { delivered, max_processes, next_index, + command_counts, .. } = self; @@ -580,6 +636,7 @@ impl PacketSimulator { tracing::trace!("packet replayed"); } + command_counts[command as usize] += 1; delivered.push(packet); } } @@ -588,6 +645,20 @@ impl PacketSimulator { std::mem::take(&mut self.delivered) } + /// Packets delivered so far, per [`Command2`] discriminant. See + /// [`Self::command_counts`]'s field docs for why this counts at delivery. + #[must_use] + pub const fn command_counts(&self) -> &[u64; COMMAND_COUNT_MAX] { + &self.command_counts + } + + /// Whether any packet of this command has been delivered. The question a + /// coverage assertion actually asks. + #[must_use] + pub const fn delivered_any(&self, command: Command2) -> bool { + self.command_counts[command as usize] > 0 + } + /// Return a previously taken buffer for reuse. pub fn recycle_buffer(&mut self, mut buf: Vec<Packet>) { buf.clear();
