This is an automated email from the ASF dual-hosted git repository. numinnex pushed a commit to branch fix_metadata_cancel_safe in repository https://gitbox.apache.org/repos/asf/iggy.git
commit 72bf394ec7ecf0ed086802f43d824cf55844b3cd Author: Grzegorz Koszyk <[email protected]> AuthorDate: Fri Jul 17 15:27:01 2026 +0200 fix(metadata): fix metadata pipeline commit --- core/consensus/src/plane_helpers.rs | 22 +++ core/metadata/src/impls/metadata.rs | 302 ++++++++++++++++++++++++++++++++---- core/server-ng/src/http/state.rs | 21 ++- core/server-ng/src/http/submit.rs | 31 ++-- 4 files changed, 333 insertions(+), 43 deletions(-) diff --git a/core/consensus/src/plane_helpers.rs b/core/consensus/src/plane_helpers.rs index 0b454c078..e385566a4 100644 --- a/core/consensus/src/plane_helpers.rs +++ b/core/consensus/src/plane_helpers.rs @@ -255,6 +255,28 @@ where drained } +/// Header of the pipeline head, iff its op is covered by the commit frontier. +/// +/// Peek-only counterpart of [`drain_committable_prefix`] for commit paths that +/// must survive their driving future being canceled between "committable" and +/// "applied" (see `IggyMetadata::on_ack`): the caller peeks here, performs its +/// awaits with the entry still in the pipeline, then — in a sync region — +/// revalidates that the head is still this exact entry before popping and +/// applying it. A driver dropped at an await strands nothing; a sibling driver +/// that committed the op first fails the caller's revalidation and re-peeks. +pub fn peek_committable_head<B, P>(consensus: &VsrConsensus<B, P>) -> Option<PrepareHeader> +where + B: MessageBus, + P: Pipeline<Entry = PipelineEntry>, +{ + let commit = consensus.commit_max(); + let pipeline = consensus.pipeline().borrow(); + pipeline + .head() + .map(|entry| entry.header) + .filter(|header| header.op <= commit) +} + /// Build reply for a committed prepare. /// /// Every field except `size` comes from `prepare_header`, bytes are identical diff --git a/core/metadata/src/impls/metadata.rs b/core/metadata/src/impls/metadata.rs index 2cddd1c50..dd955e16c 100644 --- a/core/metadata/src/impls/metadata.rs +++ b/core/metadata/src/impls/metadata.rs @@ -27,9 +27,10 @@ use consensus::{ PipelineEntry, Plane, PlaneIdentity, PlaneKind, PreflightOutcome, Project, ReplicaLogContext, RequestLogEvent, Sequencer, SimEventKind, VsrConsensus, ack_preflight, ack_quorum_reached, apply_preflight_consensus_plane, build_eviction_message, build_reply_message, - build_reply_message_with, build_result_rejection_reply, drain_committable_prefix, - emit_sim_event, fence_old_prepare_by_commit, is_caught_up_primary, - panic_if_hash_chain_would_break_in_same_view, pipeline_prepare_common, register_preflight, + build_reply_message_with, build_result_rejection_reply, emit_sim_event, + fence_old_prepare_by_commit, is_caught_up_primary, + panic_if_hash_chain_would_break_in_same_view, peek_committable_head, pipeline_prepare_common, + register_preflight, replicate_preflight, replicate_to_next_in_chain, request_preflight, send_eviction_to_client, send_prepare_ok as send_prepare_ok_common, }; @@ -796,22 +797,27 @@ where "ack quorum received" ); - let drained = drain_committable_prefix(consensus); - let drained_count = drained.len(); - if let (Some(first), Some(last)) = (drained.first(), drained.last()) { - debug!( - target: "iggy.metadata.diag", - plane = "metadata", - replica_id = consensus.replica(), - first_op = first.header.op, - last_op = last.header.op, - drained_count = drained_count, - "draining committed metadata prefix" - ); - } - - for mut entry in drained { - let prepare_header = entry.header; + // Commit loop: peek -> await journal read -> revalidate head -> + // sync {pop, apply, advance_commit_min}. + // + // The entry stays in the pipeline across the journal-read await, + // so a driver of this function that is dropped there — a hyper + // HTTP handler future canceled by peer disconnect, or a parked + // in-process submitter — strands nothing: the next driver + // (sibling submit, shard pump, repair tick) re-peeks the same + // head and commits it. Popping BEFORE the await loses the entry + // forever on cancellation (nothing re-applies a popped entry; + // `repair_primary_self_acks` is re-ack-only), pinning + // `commit_min` below `commit_max` and panicking the next commit + // with "commit_min must advance sequentially". + // + // Concurrent drivers are serialized by the head revalidation: + // only the driver that still finds its peeked header at the + // pipeline head after the await owns that op's commit; everyone + // else re-peeks and moves on to the next committable op. + let mut commits = 0usize; + let mut wire_replies: Vec<(CommitLogEvent, Message<ReplyHeader>)> = Vec::new(); + while let Some(prepare_header) = peek_committable_head(consensus) { // TODO(hubcio): should we replace this with graceful fallback (warn + return)? // When journal compaction is implemented compaction could race // with this lookup if it removes entries below the commit number. @@ -826,6 +832,22 @@ where ) }); + // Revalidate after the await: a sibling driver may have + // committed this op (and more) while we were parked. + let head_is_ours = consensus.pipeline().borrow().head().is_some_and(|head| { + head.header.op == prepare_header.op + && head.header.checksum == prepare_header.checksum + }); + if !head_is_ours { + continue; + } + + let mut entry = consensus + .pipeline() + .borrow_mut() + .pop() + .expect("on_ack: revalidated head exists"); + let pipeline_depth = consensus.pipeline().borrow().len(); let event = CommitLogEvent { replica: ReplicaLogContext::from_consensus(consensus, PlaneKind::Metadata), @@ -841,8 +863,11 @@ where // proof the table is caught up. Table first, counter last: // panic mid-commit leaves the gate closed. // - // Invariant: no .await or panic between client_table.commit_* - // and advance_commit_min. Sync-only. + // Invariant: no .await or panic from the pop above through + // `advance_commit_min` and the subscriber fire below. + // Sync-only — this is what makes pop/apply/advance atomic on + // the single-threaded shard and keeps the head revalidation + // sound. let reply = if prepare_header.operation == Operation::Register { // Register: commit_register creates session, no SM. let reply = build_reply_message(&prepare_header, &bytes::Bytes::new()); @@ -911,10 +936,12 @@ where }; consensus.advance_commit_min(prepare_header.op); emit_sim_event(SimEventKind::OperationCommitted, &event); + commits += 1; // Fire subscriber BEFORE wire send. Slot already updated // (slot-first ordering, see take_reply_sender). Dropped - // receiver: ignored. + // receiver: ignored. Still inside the sync region, so an + // in-process awaiter is woken atomically with its commit. let had_in_process_subscriber = entry.has_reply_sender(); if let Some(sender) = entry.take_reply_sender() { let _ = sender.send(reply.clone()); @@ -926,33 +953,39 @@ where // the same socket. Sending both desyncs the SDK -- it reads // the first frame, fails to decode the typed body, and // leaves the second frame stuck in the socket buffer. - if had_in_process_subscriber { - continue; + if !had_in_process_subscriber { + wire_replies.push((event, reply)); } + } + // Wire replies AFTER the commit loop: this region may await, and + // a driver dropped here loses only reply frames — every commit + // above is applied and its reply cached in the client_table, so + // the SDK recovers it via request replay. + for (event, reply) in wire_replies { let generic_reply = reply.into_generic(); let reply_buffers = freeze_client_reply(generic_reply); emit_sim_event(SimEventKind::ClientReplyEmitted, &event); if let Err(e) = consensus .message_bus() - .send_to_client(prepare_header.client, reply_buffers) + .send_to_client(event.client_id, reply_buffers) .await { error!( - client = prepare_header.client, - op = prepare_header.op, - request_id = prepare_header.request, - operation = ?prepare_header.operation, + client = event.client_id, + op = event.op, + request_id = event.request_id, + operation = ?event.operation, %e, "client reply forward failed, no retransmit path; client will time out", ); } } - // Each commit frees one prepare slot, promote up to - // drained_count buffered requests so the pipeline stays busy. - self.drain_request_queue_into_prepares(drained_count).await; + // Each commit frees one prepare slot, promote up to that many + // buffered requests so the pipeline stays busy. + self.drain_request_queue_into_prepares(commits).await; } } } @@ -2687,4 +2720,209 @@ mod tests { "ServerDefault expiry must be stamped to the configured default at admission" ); } + + /// Bus whose `send_to_client` parks forever while `stall` is set, + /// recording each parked send in `stall_hits`. Models a client whose + /// connection writer stalled, so an `on_ack` driver suspends at a wire + /// send — an await the test can then cancel the driver at. + #[derive(Debug, Default)] + struct StallBus { + stall: std::cell::Cell<bool>, + stall_hits: std::cell::Cell<u32>, + } + + // Cell fields make the futures !Send; fine on the single-threaded shard. + #[allow(clippy::future_not_send)] + impl MessageBus for StallBus { + fn track_background(&self, _handle: JoinHandle<()>) {} + async fn send_to_client( + &self, + _client_id: u128, + _data: Frozen<MESSAGE_ALIGN>, + ) -> Result<(), SendError> { + if self.stall.get() { + self.stall_hits.set(self.stall_hits.get() + 1); + std::future::pending::<()>().await; + } + Ok(()) + } + async fn send_to_replica( + &self, + _replica: u8, + _data: Frozen<MESSAGE_ALIGN>, + ) -> Result<(), SendError> { + Ok(()) + } + fn set_connection_lost_fn(&self, _f: ConnectionLostFn) {} + fn set_replica_forward_fn(&self, _f: ReplicaForwardFn) {} + fn set_client_forward_fn(&self, _f: ClientForwardFn) {} + } + + fn create_stream_request(client: u128, request: u64, name: &str) -> Message<RequestHeader> { + let body = iggy_binary_protocol::requests::streams::CreateStreamRequest { + name: WireName::new(name).unwrap(), + } + .to_bytes(); + let header_size = size_of::<RequestHeader>(); + let total = header_size + body.len(); + let mut message = Message::<RequestHeader>::new(total); + { + let slice = message.as_mut_slice(); + slice[header_size..total].copy_from_slice(&body); + let header = + bytemuck::checked::from_bytes_mut::<RequestHeader>(&mut slice[..header_size]); + *header = RequestHeader { + command: Command2::Request, + operation: Operation::CreateStream, + size: u32::try_from(total).unwrap(), + client, + session: 1, + request, + user_id: 0, + namespace: server_common::sharding::METADATA_CONSENSUS_NAMESPACE, + ..Default::default() + }; + } + message + } + + /// Reproduces the `commit_min must advance sequentially` shard-0 crash. + /// + /// `on_ack` drains (pops) the committable prefix off the pipeline and only + /// then applies it, with awaits in between (journal read, wire send). Any + /// driver of `on_ack` that is dropped at one of those awaits — a hyper + /// HTTP handler future canceled by peer disconnect (`http/state.rs`), or + /// any parked in-process submitter — strands the popped-but-unapplied + /// entries: nothing can re-apply them (`repair_primary_self_acks` is + /// re-ack-only, above `commit_max`), `commit_min` is pinned below + /// `commit_max` (every login rejected `NotCaughtUp`), and the next commit + /// that quorums panics the shard. + /// + /// The test parks a driver mid-commit exactly there, cancels it, and then + /// delivers the next ack. Correct behavior: the stranded op is still in + /// the pipeline and the late ack commits it and everything after it, in + /// order. Broken behavior: panic "expected 2, got 3". + #[compio::test] + async fn dropped_on_ack_driver_must_not_lose_popped_commits() { + use std::future::Future; + + const CLIENT: u128 = 1; + // The session is the Register op's commit number; it must be > 0 and + // sort at-or-under the commits of the three ops below (1..=3). + const SESSION: u64 = 1; + const ACTING_USER: u32 = 7; + + let dir = tempfile::tempdir().unwrap(); + let journal = journal::prepare_journal::PrepareJournal::open( + &dir.path().join("journal.wal"), + 0, + ) + .await + .unwrap(); + let consensus = VsrConsensus::new( + 1, + 0, + 1, + server_common::sharding::METADATA_CONSENSUS_NAMESPACE, + StallBus::default(), + LocalPipeline::new(), + ); + consensus.init(); + let md: IggyMetadata<_, journal::prepare_journal::PrepareJournal, (), TestMux> = + IggyMetadata::new(Some(consensus), Some(journal), None, TestMux::default(), None); + let consensus = md.consensus.as_ref().unwrap(); + + md.client_table.borrow_mut().commit_register( + CLIENT, + ACTING_USER, + register_reply(CLIENT, SESSION), + |_| false, + ); + + // Three prepares through the real primary path: pipeline entry, WAL + // append, self-ack onto the loopback queue. + for (i, name) in ["s1", "s2", "s3"].iter().enumerate() { + let prepare = md + .prepare_request(create_stream_request(CLIENT, i as u64 + 1, name)) + .expect("CreateStream is client-allowed"); + consensus.pipeline_message(PlaneKind::Metadata, &prepare); + md.on_replicate(prepare).await; + } + let mut loopback = Vec::new(); + consensus.drain_loopback_into(&mut loopback); + let mut acks = loopback + .into_iter() + .map(|message| { + message + .try_into_typed::<PrepareOkHeader>() + .expect("loopback holds self PrepareOks") + }) + .collect::<Vec<_>>(); + assert_eq!(acks.len(), 3, "one self-ack per replicated prepare"); + let ack3 = acks.pop().unwrap(); + let ack2 = acks.pop().unwrap(); + let ack1 = acks.pop().unwrap(); + + // Ack op 2 first: quorum for op 2 alone, but the contiguous prefix + // still starts at the un-acked op 1, so nothing commits yet. + md.on_ack(ack2).await; + assert_eq!(consensus.commit_max(), 0); + assert_eq!(consensus.commit_min(), 0); + + // Ack op 1: the quorum walk covers ops 1..=2, so this single driver + // commits both. Poll it by hand until it parks at a stalled wire + // send mid-`on_ack`, then cancel it — the moral equivalent of hyper + // dropping an HTTP handler future on peer disconnect. + consensus.message_bus().stall.set(true); + { + let mut driver = Box::pin(md.on_ack(ack1)); + let waker = std::task::Waker::noop(); + let mut cx = std::task::Context::from_waker(waker); + let mut parked_at_send = false; + for _ in 0..1_000 { + assert!( + driver.as_mut().poll(&mut cx).is_pending(), + "driver must park at the stalled wire send, not complete" + ); + if consensus.message_bus().stall_hits.get() > 0 { + parked_at_send = true; + break; + } + // Let the runtime process the journal-read completion the + // driver is waiting on, then poll again. + compio::time::sleep(std::time::Duration::from_millis(1)).await; + } + assert!( + parked_at_send, + "driver never reached a wire send (commit_min={})", + consensus.commit_min() + ); + // Cancel mid-`on_ack`, with at least op 1 applied and the reply + // send in flight. + drop(driver); + } + consensus.message_bus().stall.set(false); + assert_eq!(consensus.commit_max(), 2, "quorum walk advanced commit_max"); + assert!( + consensus.commit_min() >= 1, + "driver applied op 1 before parking at the wire send" + ); + + // Whatever the canceled driver left behind must still be + // committable: delivering the ack for op 3 has to commit every + // remaining op, in order. The broken commit path lost op 2 with the + // dropped driver (popped, never applied) and panics here with + // "commit_min must advance sequentially: expected 2, got 3". + md.on_ack(ack3).await; + assert_eq!( + consensus.commit_min(), + 3, + "late ack must commit the stranded op 2 and then op 3" + ); + assert_eq!(consensus.commit_max(), 3); + assert!( + is_caught_up_primary(consensus), + "gate must reopen once the prefix is fully applied" + ); + } } diff --git a/core/server-ng/src/http/state.rs b/core/server-ng/src/http/state.rs index 55df0d8e8..777518946 100644 --- a/core/server-ng/src/http/state.rs +++ b/core/server-ng/src/http/state.rs @@ -28,6 +28,7 @@ use axum::http::{HeaderName, HeaderValue}; use axum::response::Response; use configs::server_ng::NgSystemConfig; use consensus::{MetadataHandle, VsrConsensus}; +use futures::channel::oneshot; use iggy_common::{ClusterMetadata, IggyTimestamp}; use message_bus::InstanceToken; use send_wrapper::SendWrapper; @@ -205,8 +206,26 @@ impl HttpInner { } // Shared Register entry point; on shard 0 (always, for HTTP) it runs // `submit_register_in_process` directly on the metadata owner. - let session = submit_register_on_owner(&self.shard, client_id, user_id) + // + // Detached so a client disconnect cannot cancel the Register + // mid-flight: the in-process submit drives shared consensus machinery + // (pipeline push, WAL append, the `on_ack` commit loop), and hyper + // drops this handler future the moment the HTTP peer disconnects. + // A canceled submit used to strand consensus state mid-await; now the + // detached task always drives it to completion and a disconnect only + // drops the receiver half (same discipline as `submit_committed`). + let (result_slot, committed) = oneshot::channel(); + let shard = Rc::clone(&self.shard); + compio::runtime::spawn(async move { + let result = submit_register_on_owner(&shard, client_id, user_id).await; + // A failed send means the handler died mid-await; the Register + // itself has already committed, which is what matters. + let _ = result_slot.send(result); + }) + .detach(); + let session = committed .await + .map_err(|_| AuthError::SessionUnavailable)? .map_err(|error| { warn!(?error, "server-ng HTTP: VSR Register submit failed"); AuthError::SessionUnavailable diff --git a/core/server-ng/src/http/submit.rs b/core/server-ng/src/http/submit.rs index 6e549a599..f343f546f 100644 --- a/core/server-ng/src/http/submit.rs +++ b/core/server-ng/src/http/submit.rs @@ -260,18 +260,29 @@ pub(in crate::http) async fn logout_session(state: &HttpInner, session: &Rc<Http // is terminal for this session, so it needs no gate-issued contiguous id // (mirrors the disconnect path's `u64::MAX`). const LOGOUT_REQUEST_ID: u64 = u64::MAX; - if let Err(error) = submit_logout_on_owner( - &state.shard, - session.client_id, - session.session, - LOGOUT_REQUEST_ID, - ) - .await - { - warn!( + // Detached for the same reason as `submit_committed`: the submit drives + // shared consensus machinery, and axum drops this handler future on + // client disconnect. A cancel mid-await used to strand consensus state; + // the detached task always drives the Logout to completion. + let (result_slot, done) = oneshot::channel(); + let shard = Rc::clone(&state.shard); + let vsr_client_id = session.client_id; + let vsr_session = session.session; + compio::runtime::spawn(async move { + let result = + submit_logout_on_owner(&shard, vsr_client_id, vsr_session, LOGOUT_REQUEST_ID).await; + let _ = result_slot.send(result); + }) + .detach(); + match done.await { + Ok(Ok(_)) => {} + Ok(Err(error)) => warn!( ?error, "server-ng HTTP: VSR Logout submit failed; slot lingers until eviction" - ); + ), + Err(_canceled) => warn!( + "server-ng HTTP: VSR Logout task dropped before replying; slot lingers until eviction" + ), } state.forget_session(session); }
