diff --git a/zakura-network/src/zakura/block_sync/peer_routine.rs b/zakura-network/src/zakura/block_sync/peer_routine.rs index 4b423c790..f295f8902 100644 --- a/zakura-network/src/zakura/block_sync/peer_routine.rs +++ b/zakura-network/src/zakura/block_sync/peer_routine.rs @@ -1158,14 +1158,30 @@ impl PeerRoutine { // continuous backpressure): plausibly transient local write congestion, not // a dead peer. While outbound is full the select loop does not drain inbound // frames (`if outbound_queue_has_capacity`), so a block the peer already sent - // may be waiting behind our write side. Grant one short, BOUNDED grace. This - // is the *only* liveness extension: a peer that stopped reading holds outbound - // full past `request_timeout`, falls through to the disconnect arm, and is - // disconnected at the liveness deadline — it cannot dodge the timer. + // may be waiting behind our write side. Grant one bounded grace period. + // A peer that stops reading remains full and is then disconnected. self.window .extend_liveness_deadline(now, self.config.request_timeout); Ok(()) } + LivenessOutcome::Disconnect if self.sequencer_view.borrow().applying_len > 0 => { + // Bodies are sitting in the LOCAL apply pipeline, so "no accepted + // block progress" cannot be blamed on this peer: nothing can make + // accepted progress while our own commits are not landing. + // This has been observed to exile healthy peers when stuck in a verifier-layer commit stall. + // Defer the disconnect. A genuinely dead peer is caught at the next + // liveness deadline after the pipeline drains (`applying_len == 0`); + // its timed-out requests were already requeued by `expire_due_timeouts`. + self.window + .extend_liveness_deadline(now, self.config.request_timeout); + metrics::counter!("sync.block.liveness.deferred_local_stall").increment(1); + tracing::debug!( + peer = ?self.peer, + applying_len = self.sequencer_view.borrow().applying_len, + "deferring block-sync liveness disconnect: local apply pipeline is stalled" + ); + Ok(()) + } LivenessOutcome::Disconnect => { let error = "block-sync peer made no accepted block progress before liveness deadline"; @@ -1174,7 +1190,7 @@ impl PeerRoutine { now + self.config.effective_no_progress_peer_cooldown(), ); self.trace_protocol_reject_liveness(error); - tracing::debug!( + tracing::info!( peer = ?self.peer, outstanding = self.window.outstanding.len(), "disconnecting Zakura block-sync peer after no accepted block progress" @@ -1218,7 +1234,10 @@ impl PeerRoutine { } if removed { self.publish_outstanding(); - self.window.disarm_liveness_after_progress_if_idle(); + // Other peer satisfied requests below the floor. + // Clear this peer's liveness deadline when no requests remain. + // Unanswered requests expire separately via `expire_due_timeouts`. + self.window.clear_liveness_if_idle(); } } @@ -2413,8 +2432,11 @@ mod tests { use super::super::peer_registry::PeerRegistry; use super::super::request::BlockSizeEstimate; + use super::super::request::{BlockRangeRequest, ExpectedBlock}; use super::super::sequencer_task::{initial_view, SequencerControlInput}; - use super::super::state::{ByteBudget, ThroughputMeter}; + use super::super::state::{ + ByteBudget, LivenessOutcome, OutstandingBlockRange, ReceivedBlockTracker, ThroughputMeter, + }; use super::super::work_queue::WorkQueue; use super::super::{BlockSyncFrontiers, BlockSyncPeerSession, ZakuraBlockSyncConfig}; use super::PeerRoutine; @@ -2680,4 +2702,190 @@ mod tests { request_timeout )); } + + /// Floor-GC of a request satisfied by other peers must clear the liveness + /// deadline: the peer no longer owes those heights. Before the fix, a peer + /// that had delivered earlier and then had a NEWER request GC'd below the + /// floor kept a stale deadline (`request_at > block_at` fails the + /// after-progress disarm) and was disconnected at tip with zero + /// outstanding, then parked in the 180 s no-progress cooldown — serially + /// exiling healthy peers until the block-sync peer set collapsed + /// (reproduced live: 5 → 0 peers after a fleet-wide restart). + #[tokio::test] + async fn floor_gc_clears_stale_liveness_deadline() { + let config = ZakuraBlockSyncConfig::default(); + let budget = ByteBudget::new(1_000_000); + let work = Arc::new(WorkQueue::new(block::Height(0))); + + let cancel = CancellationToken::new(); + let (out_send, _out_recv) = framed_channel(16); + let (_in_send, in_recv) = framed_channel(16); + let peer = ZakuraPeerId::new(vec![8u8; 32]).expect("test peer id is within bounds"); + let session = BlockSyncPeerSession::for_test(peer.clone(), out_send, cancel.clone()); + + let (sequencer_input_tx, _sequencer_input_rx) = mpsc::channel(16); + let (control_tx, _control_rx) = mpsc::unbounded_channel(); + let (actions_tx, _actions_rx) = mpsc::channel(16); + let (routine_to_reactor_tx, _routine_to_reactor_rx) = mpsc::channel(16); + // The download floor sits above the request below, as if other peers + // delivered those heights and the sequencer committed them. + let (_view_tx, view_rx) = watch::channel(initial_view(BlockSyncFrontiers { + finalized_height: block::Height(0), + verified_block_tip: block::Height(100), + verified_block_hash: block::Hash([0; 32]), + })); + + let mut routine = PeerRoutine::new( + peer, + session, + in_recv, + config, + 0, + budget, + work, + Arc::new(PeerRegistry::new()), + Arc::new(Mutex::new(ThroughputMeter::new(Instant::now()))), + sequencer_input_tx, + Arc::new(AtomicU64::new(0)), + control_tx, + actions_tx, + routine_to_reactor_tx, + view_rx, + cancel, + ZakuraTrace::noop(), + ); + + let now = Instant::now(); + let liveness = Duration::from_secs(30); + + // The peer delivered a block for an earlier request... + routine + .window + .arm_liveness(now - Duration::from_secs(120), liveness); + routine + .window + .note_block_progress(now - Duration::from_secs(90), liveness); + + // ...then received a newer request (deadline re-armed) whose heights + // the floor later passed. + routine + .window + .arm_liveness(now - Duration::from_secs(60), liveness); + routine.window.outstanding.push(OutstandingBlockRange { + request: BlockRangeRequest { + start_height: block::Height(99), + count: 2, + anchor_hash: block::Hash([9; 32]), + estimated_bytes: 0, + expected_blocks: vec![ + ExpectedBlock { + height: block::Height(99), + hash: block::Hash([99; 32]), + estimated_bytes: 0, + }, + ExpectedBlock { + height: block::Height(100), + hash: block::Hash([100; 32]), + estimated_bytes: 0, + }, + ], + }, + queued_at: now - Duration::from_secs(60), + deadline: now + Duration::from_secs(60), + delivery_snapshot: routine + .window + .delivery_snapshot(now - Duration::from_secs(60)), + delivered_bytes: 0, + received: ReceivedBlockTracker::default(), + }); + + routine.gc_committed_outstanding(); + + assert!( + routine.window.outstanding.is_empty(), + "the floor passed the whole request, so GC must remove it" + ); + assert_eq!( + routine.window.block_liveness_deadline, None, + "a request satisfied below the floor must not leave a stale liveness deadline" + ); + assert!( + matches!(routine.window.check_liveness(now), LivenessOutcome::Ok), + "the peer owes nothing and must not be disconnected" + ); + } + + /// An expired liveness deadline must not disconnect a peer while the LOCAL + /// apply pipeline holds unfinished bodies: nothing can make "accepted block + /// progress" while our own commits are not landing. During a live + /// verifier-layer commit stall this exiled healthy peers one by one into + /// the no-progress cooldown until the peer set collapsed. Once the + /// pipeline drains, the same expired deadline disconnects as before. + #[tokio::test] + async fn liveness_disconnect_is_deferred_while_local_applies_are_stalled() { + let config = ZakuraBlockSyncConfig::default(); + let budget = ByteBudget::new(1_000_000); + let work = Arc::new(WorkQueue::new(block::Height(0))); + + let cancel = CancellationToken::new(); + let (out_send, _out_recv) = framed_channel(16); + let (_in_send, in_recv) = framed_channel(16); + let peer = ZakuraPeerId::new(vec![9u8; 32]).expect("test peer id is within bounds"); + let session = BlockSyncPeerSession::for_test(peer.clone(), out_send, cancel.clone()); + + let (sequencer_input_tx, _sequencer_input_rx) = mpsc::channel(16); + let (control_tx, _control_rx) = mpsc::unbounded_channel(); + let (actions_tx, _actions_rx) = mpsc::channel(16); + let (routine_to_reactor_tx, _routine_to_reactor_rx) = mpsc::channel(16); + let mut view = initial_view(BlockSyncFrontiers { + finalized_height: block::Height(0), + verified_block_tip: block::Height(100), + verified_block_hash: block::Hash([0; 32]), + }); + view.applying_len = 60; + let (view_tx, view_rx) = watch::channel(view); + + let mut routine = PeerRoutine::new( + peer, + session, + in_recv, + config, + 0, + budget, + work, + Arc::new(PeerRegistry::new()), + Arc::new(Mutex::new(ThroughputMeter::new(Instant::now()))), + sequencer_input_tx, + Arc::new(AtomicU64::new(0)), + control_tx, + actions_tx, + routine_to_reactor_tx, + view_rx, + cancel, + ZakuraTrace::noop(), + ); + + let now = Instant::now(); + // An armed deadline that expired 10s ago. + routine + .window + .arm_liveness(now - Duration::from_secs(60), Duration::from_secs(50)); + + assert!( + routine.check_block_liveness(now).is_ok(), + "an expired deadline must be deferred while local applies are stalled" + ); + + // The pipeline drains: the same expired deadline now disconnects. + view.applying_len = 0; + view_tx + .send(view) + .expect("view receiver is held by the routine"); + // The deferral extended the deadline by request_timeout; move past it. + let later = now + routine.config.request_timeout + Duration::from_secs(1); + assert!( + routine.check_block_liveness(later).is_err(), + "a drained pipeline must restore the liveness disconnect" + ); + } }