From b6bd48b14b0f4cb9f8b2a833653288f91d7be135 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 2 Oct 2026 03:30:58 +0000 Subject: [PATCH 1/2] fix(coordinator): run the pruned-log fallback param replay off the event loop When a reconnecting daemon's state-log ack had been pruned, the daemon-status handler awaited `replay_persisted_params_for_daemon` inline. That replay sends every persisted param one round-trip at a time, each bounded only by `TCP_READ_TIMEOUT` (30s). A slow or half-dead daemon could therefore stall the whole coordinator for N x 30s per affected dataflow: no CLI requests, no heartbeats, and no other daemons' events. Spawn the replay as the other two replay call sites already do. It reports back through a new internal `ParamFallbackReplayFinished` event, which advances the daemon's ack to the `state_log_sequence` captured *before* the replay started (never moving it backwards), and only if every param was replayed and the daemon is still on the connection the replay was sent on. A per-daemon in-flight mark, keyed by connection, keeps a later status report on the same connection from starting a duplicate replay; a reconnect starts a fresh one. Fixes #3684 Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_016ZCEkLBowwvtAMhopshZwt --- binaries/coordinator/src/daemon_events.rs | 37 +++- binaries/coordinator/src/events.rs | 17 ++ binaries/coordinator/src/handlers.rs | 2 + binaries/coordinator/src/lib.rs | 26 ++- binaries/coordinator/src/params.rs | 149 +++++++++---- binaries/coordinator/src/state.rs | 5 + binaries/coordinator/src/tests.rs | 245 +++++++++++++++++++++- 7 files changed, 421 insertions(+), 60 deletions(-) diff --git a/binaries/coordinator/src/daemon_events.rs b/binaries/coordinator/src/daemon_events.rs index 64fcf430ad..fe9f1f5386 100644 --- a/binaries/coordinator/src/daemon_events.rs +++ b/binaries/coordinator/src/daemon_events.rs @@ -5,10 +5,10 @@ use crate::{ Coordinator, DaemonRequest, apply_disconnect_actions, build_result_timeout, check_build_timeouts, check_spawn_timeouts, cleanup_disconnected_daemons_from_running_builds, cleanup_disconnected_daemons_from_running_dataflows, expire_stopped_nodes, - handle_pruned_state_catchup_fallback, notify_daemons_about_disconnected_peers, + finish_pruned_state_catchup_fallback, notify_daemons_about_disconnected_peers, reestablish_running_dataflow, replay_all_nodes_ready, restore_topic_debug_streams_for_daemon, - send_heartbeat_with_timeout, state, status_report_should_stop_orphan, - stop_orphaned_dataflow_on_daemon, + send_heartbeat_with_timeout, start_pruned_state_catchup_fallback, state, + status_report_should_stop_orphan, stop_orphaned_dataflow_on_daemon, }; use dora_coordinator_store::DataflowStatus as StoreDataflowStatus; use dora_message::{ @@ -750,22 +750,47 @@ impl Coordinator { "state catch-up: log pruned for dataflow {uuid}, \ falling back to full param replay for daemon {daemon_id}" ); - handle_pruned_state_catchup_fallback( + start_pruned_state_catchup_fallback( *uuid, df, &daemon_id, self.store.clone(), &mut self.daemon_connections, self.clock.clone(), + self.internal_events.clone(), now, - ) - .await; + ); } } } Ok(()) } + pub(crate) fn handle_param_fallback_replay_finished( + &mut self, + dataflow_id: DataflowId, + daemon_id: DaemonId, + connection_id: Uuid, + ack_sequence: u64, + succeeded: bool, + ) { + let Some(df) = self.running_dataflows.get_mut(&dataflow_id) else { + return; + }; + let current_connection_id = self + .daemon_connections + .get_mut(&daemon_id) + .map(|connection| connection.connection_id); + finish_pruned_state_catchup_fallback( + df, + &daemon_id, + connection_id, + current_connection_id, + ack_sequence, + succeeded, + ); + } + pub(crate) async fn handle_state_catch_up_ack( &mut self, daemon_id: DaemonId, diff --git a/binaries/coordinator/src/events.rs b/binaries/coordinator/src/events.rs index 75bedfb40f..6774262bd1 100644 --- a/binaries/coordinator/src/events.rs +++ b/binaries/coordinator/src/events.rs @@ -92,6 +92,22 @@ pub enum Event { node_id: NodeId, clean_stop: bool, }, + /// A full param replay started by the pruned-state-log fallback has + /// finished. It runs in a spawned task, off the event loop (#3684), and + /// reports back here so the loop can advance the daemon's ack. + ParamFallbackReplayFinished { + dataflow_id: DataflowId, + daemon_id: DaemonId, + /// The connection the replay was sent on: a reply for an older + /// connection must not mark a reconnected daemon as caught up. + connection_id: Uuid, + /// `state_log_sequence` when the replay started. The ack advances to + /// this, not to the sequence at completion, which may include entries + /// appended while the replay ran. + ack_sequence: u64, + /// Whether every param was replayed. + succeeded: bool, + }, } impl Event { @@ -122,6 +138,7 @@ impl Event { Event::DaemonStatusReport { .. } => "DaemonStatusReport", Event::DaemonStateCatchUpAck { .. } => "DaemonStateCatchUpAck", Event::DaemonNodeStopped { .. } => "DaemonNodeStopped", + Event::ParamFallbackReplayFinished { .. } => "ParamFallbackReplayFinished", } } } diff --git a/binaries/coordinator/src/handlers.rs b/binaries/coordinator/src/handlers.rs index 459ee48f27..9cf3921e5e 100644 --- a/binaries/coordinator/src/handlers.rs +++ b/binaries/coordinator/src/handlers.rs @@ -767,6 +767,7 @@ pub(crate) async fn start_dataflow( store_generation: 0, last_recovery_attempt: BTreeMap::new(), last_replay_attempt: BTreeMap::new(), + fallback_replay_in_flight: BTreeMap::new(), uv, state_log_sequence: 0, state_log: Vec::new(), @@ -1002,6 +1003,7 @@ mod tests { store_generation: 0, last_recovery_attempt: BTreeMap::new(), last_replay_attempt: BTreeMap::new(), + fallback_replay_in_flight: BTreeMap::new(), uv: false, state_log_sequence: 0, state_log: Vec::new(), diff --git a/binaries/coordinator/src/lib.rs b/binaries/coordinator/src/lib.rs index 1f4ecac4b1..dfe1cc29f6 100644 --- a/binaries/coordinator/src/lib.rs +++ b/binaries/coordinator/src/lib.rs @@ -145,8 +145,8 @@ pub(crate) use daemon_liveness::{ }; pub(crate) use params::{ ensure_delete_param_forward_applied, ensure_set_param_forward_applied, - handle_pruned_state_catchup_fallback, replay_persisted_params_for_daemon, - schedule_param_replay_for_ready_dataflow, + finish_pruned_state_catchup_fallback, replay_persisted_params_for_daemon, + schedule_param_replay_for_ready_dataflow, start_pruned_state_catchup_fallback, }; pub(crate) use ready_barrier::{ broadcast_all_nodes_ready, nodes_on_daemon, replay_all_nodes_ready, @@ -362,6 +362,8 @@ pub(crate) struct Coordinator { pub(crate) otel_metrics: otel_metrics::SharedMetrics, /// Aborts the event stream on `dora down` / Ctrl-C. pub(crate) abort_handle: futures::stream::AbortHandle, + /// Feeds events back into the event loop, for tasks it spawns. + pub(crate) internal_events: tokio::sync::mpsc::Sender, } impl Coordinator { @@ -393,9 +395,13 @@ async fn start_inner( tokio_stream::wrappers::IntervalStream::new(tokio::time::interval(Duration::from_secs(3))) .map(|_| Event::DaemonHeartbeatInterval); + // results of work the event loop spawned off itself + let (internal_events_tx, internal_events_rx) = tokio::sync::mpsc::channel::(64); + let internal_events = ReceiverStream::new(internal_events_rx); + // events that should be aborted on `dora down` let (abortable_events, abort_handle) = - futures::stream::abortable((events, daemon_heartbeat_interval).merge()); + futures::stream::abortable((events, daemon_heartbeat_interval, internal_events).merge()); let mut events = abortable_events; @@ -467,6 +473,7 @@ async fn start_inner( #[cfg(feature = "metrics")] otel_metrics, abort_handle, + internal_events: internal_events_tx, }; while let Some(event) = events.next().await { @@ -574,6 +581,19 @@ async fn start_inner( .handle_daemon_node_stopped(daemon_id, dataflow_id, node_id, clean_stop) .await? } + Event::ParamFallbackReplayFinished { + dataflow_id, + daemon_id, + connection_id, + ack_sequence, + succeeded, + } => coordinator.handle_param_fallback_replay_finished( + dataflow_id, + daemon_id, + connection_id, + ack_sequence, + succeeded, + ), } // warn if event handling took too long -> the main loop should never be blocked for too long diff --git a/binaries/coordinator/src/params.rs b/binaries/coordinator/src/params.rs index 3306f9cfa5..e092fa7cc3 100644 --- a/binaries/coordinator/src/params.rs +++ b/binaries/coordinator/src/params.rs @@ -2,13 +2,15 @@ //! replaying persisted parameters to a daemon that (re)joins a dataflow. use crate::state::{DaemonConnections, RunningDataflow}; -use crate::{FALLBACK_REPLAY_BACKOFF, nodes_on_daemon}; +use crate::{Event, FALLBACK_REPLAY_BACKOFF, nodes_on_daemon}; use dora_core::uhlc::HLC; use dora_message::{DataflowId, common::DaemonId, daemon_to_coordinator::DaemonCoordinatorReply}; use eyre::eyre; use serde::Serialize; use serde_json::value::RawValue; use std::{sync::Arc, time::Instant}; +use tokio::sync::mpsc; +use uuid::Uuid; pub(crate) struct ParamReplayItem { pub(crate) node_id: dora_core::config::NodeId, @@ -22,15 +24,40 @@ pub(crate) struct ParamReplaySummary { pub(crate) failed: usize, } -pub(crate) async fn handle_pruned_state_catchup_fallback( +/// Start a full param replay for a daemon whose state-log ack was pruned. +/// +/// The replay sends every persisted param one round-trip at a time, each +/// bounded by `TCP_READ_TIMEOUT`, so it runs in a spawned task rather than on +/// the event loop (#3684). It reports back through `events` with +/// [`Event::ParamFallbackReplayFinished`], handled by +/// [`finish_pruned_state_catchup_fallback`]. +#[allow(clippy::too_many_arguments)] +pub(crate) fn start_pruned_state_catchup_fallback( dataflow_id: DataflowId, dataflow: &mut RunningDataflow, daemon_id: &DaemonId, store: Arc, daemon_connections: &mut DaemonConnections, clock: Arc, + events: mpsc::Sender, now: Instant, ) { + let Some(connection) = daemon_connections.get_mut(daemon_id).cloned() else { + tracing::warn!( + "failed to run fallback replay for dataflow {dataflow_id}: \ + daemon {daemon_id} is not connected" + ); + return; + }; + let connection_id = connection.connection_id; + + if dataflow.fallback_replay_in_flight.get(daemon_id) == Some(&connection_id) { + tracing::debug!( + "skipping fallback replay for dataflow {dataflow_id} on daemon {daemon_id}: \ + a replay is already in flight" + ); + return; + } if let Some(last_replay_attempt) = dataflow.last_replay_attempt.get(daemon_id) && now.duration_since(*last_replay_attempt) < FALLBACK_REPLAY_BACKOFF { @@ -41,56 +68,88 @@ pub(crate) async fn handle_pruned_state_catchup_fallback( return; } - let Some(connection) = daemon_connections.get_mut(daemon_id).cloned() else { - tracing::warn!( - "failed to run fallback replay for dataflow {dataflow_id}: \ - daemon {daemon_id} is not connected" - ); - return; - }; - - let node_ids_on_daemon: Vec<_> = dataflow - .node_to_daemon - .iter() - .filter(|(_, did)| *did == daemon_id) - .map(|(node_id, _)| node_id.clone()) - .collect(); + let node_ids_on_daemon = nodes_on_daemon(dataflow, daemon_id); dataflow.last_replay_attempt.insert(daemon_id.clone(), now); + dataflow + .fallback_replay_in_flight + .insert(daemon_id.clone(), connection_id); + // Captured before the replay reads the store: everything up to here is + // in the store, so the replay covers it. Entries appended while it runs + // were forwarded to the daemon by their own handlers. + let ack_sequence = dataflow.state_log_sequence; + let daemon_id = daemon_id.clone(); - let replay_summary = replay_persisted_params_for_daemon( - dataflow_id, - daemon_id.clone(), - node_ids_on_daemon, - store, - connection, - clock, - ) - .await; + tokio::spawn(async move { + let replay_summary = replay_persisted_params_for_daemon( + dataflow_id, + daemon_id.clone(), + node_ids_on_daemon, + store, + connection, + clock, + ) + .await; + if replay_summary.failed != 0 { + tracing::warn!( + "fallback replay incomplete for dataflow {dataflow_id} on daemon \ + {daemon_id}: attempted={}, failed={}", + replay_summary.attempted, + replay_summary.failed, + ); + } + let finished = Event::ParamFallbackReplayFinished { + dataflow_id, + daemon_id, + connection_id, + ack_sequence, + succeeded: replay_summary.failed == 0, + }; + // Fails only when the coordinator is shutting down. + let _ = events.send(finished).await; + }); +} - let last_ack = dataflow +/// Record the outcome of a replay started by +/// [`start_pruned_state_catchup_fallback`]. +/// +/// `current_connection_id` is the daemon's connection now. A replay sent on +/// an older connection says nothing about the daemon behind the new one, so +/// it neither advances the ack nor clears that connection's in-flight mark. +pub(crate) fn finish_pruned_state_catchup_fallback( + dataflow: &mut RunningDataflow, + daemon_id: &DaemonId, + connection_id: Uuid, + current_connection_id: Option, + ack_sequence: u64, + succeeded: bool, +) { + if dataflow.fallback_replay_in_flight.get(daemon_id) == Some(&connection_id) { + dataflow.fallback_replay_in_flight.remove(daemon_id); + } + if current_connection_id != Some(connection_id) { + tracing::debug!( + "ignoring fallback replay result for daemon {daemon_id}: \ + it was sent on a connection that has since been replaced" + ); + return; + } + if !succeeded { + // Advancing the ack on partial failure could silently diverge + // runtime and store state; the next status report retries. + return; + } + // Replay is authoritative for pruned history. Individual SetParam + // events don't trigger StateCatchUpAck, so set the ack here to avoid + // repeated fallback replays on every status-report cycle. Never move it + // backwards past an ack that arrived meanwhile. + let current = dataflow .daemon_ack_sequence .get(daemon_id) .copied() .unwrap_or(0); - if replay_summary.failed == 0 { - // Mark daemon as caught up only when full replay succeeds: - // replay is authoritative for pruned history, and advancing - // ack on partial failure can silently diverge runtime/store state. - // Individual SetParam events don't trigger StateCatchUpAck, so - // we set ack here for successful full replay to avoid repeated - // fallback replays on every status-report cycle. - dataflow - .daemon_ack_sequence - .insert(daemon_id.clone(), dataflow.state_log_sequence); - } else { - tracing::warn!( - "fallback replay incomplete for dataflow {dataflow_id} on daemon \ - {daemon_id}: attempted={}, failed={}; leaving ack at {}", - replay_summary.attempted, - replay_summary.failed, - last_ack - ); - } + dataflow + .daemon_ack_sequence + .insert(daemon_id.clone(), current.max(ack_sequence)); } /// Load the persisted params of every node on a daemon, flattened into diff --git a/binaries/coordinator/src/state.rs b/binaries/coordinator/src/state.rs index 638d9e75e3..2296429dc4 100644 --- a/binaries/coordinator/src/state.rs +++ b/binaries/coordinator/src/state.rs @@ -390,6 +390,10 @@ pub(crate) struct RunningDataflow { /// Per-daemon timestamp of last full fallback param replay attempt /// (for state catch-up backoff on pruned logs). pub(crate) last_replay_attempt: BTreeMap, + /// Fallback param replays in flight, by daemon, with the connection ID + /// each was sent on. A status report from the same connection does not + /// start a second one; one from a new connection does. + pub(crate) fallback_replay_in_flight: BTreeMap, /// Whether UV was used for this dataflow (needed for restart). pub(crate) uv: bool, @@ -597,6 +601,7 @@ impl RunningDataflow { store_generation: record.generation, last_recovery_attempt: BTreeMap::new(), last_replay_attempt: BTreeMap::new(), + fallback_replay_in_flight: BTreeMap::new(), uv: record.uv, state_log_sequence: 0, state_log: Vec::new(), diff --git a/binaries/coordinator/src/tests.rs b/binaries/coordinator/src/tests.rs index 5ea930d967..ee8812c556 100644 --- a/binaries/coordinator/src/tests.rs +++ b/binaries/coordinator/src/tests.rs @@ -314,6 +314,7 @@ fn test_running_dataflow( store_generation: 0, last_recovery_attempt: BTreeMap::new(), last_replay_attempt: BTreeMap::new(), + fallback_replay_in_flight: BTreeMap::new(), uv: false, state_log_sequence: 0, state_log: Vec::new(), @@ -1771,16 +1772,17 @@ async fn fallback_replay_keeps_ack_unchanged_when_daemon_is_disconnected() { dataflow.daemon_ack_sequence.insert(daemon_id.clone(), 3); let mut daemon_connections = DaemonConnections::default(); - handle_pruned_state_catchup_fallback( + let (events_tx, _events_rx) = tokio::sync::mpsc::channel(8); + start_pruned_state_catchup_fallback( dataflow_id, &mut dataflow, &daemon_id, store, &mut daemon_connections, Arc::new(HLC::default()), + events_tx, Instant::now(), - ) - .await; + ); assert_eq!(dataflow.daemon_ack_sequence.get(&daemon_id), Some(&3)); assert!(!dataflow.last_replay_attempt.contains_key(&daemon_id)); @@ -1801,21 +1803,252 @@ async fn fallback_replay_respects_backoff_window() { .insert(daemon_id.clone(), Instant::now()); let mut daemon_connections = DaemonConnections::default(); - handle_pruned_state_catchup_fallback( + let (events_tx, _events_rx) = tokio::sync::mpsc::channel(8); + start_pruned_state_catchup_fallback( dataflow_id, &mut dataflow, &daemon_id, store, &mut daemon_connections, Arc::new(HLC::default()), + events_tx, Instant::now(), - ) - .await; + ); // No replay attempted while backoff is active, so ack remains unchanged. assert_eq!(dataflow.daemon_ack_sequence.get(&daemon_id), Some(&2)); } +#[tokio::test] +async fn fallback_replay_does_not_wait_for_an_unresponsive_daemon() { + // #3684: the fallback runs from the daemon-status handler on the event + // loop. A daemon that never answers must not stall that loop for + // `TCP_READ_TIMEOUT` per persisted param. + let store: Arc = Arc::new(InMemoryStore::new()); + let dataflow_id = DataflowId::from(Uuid::new_v4()); + let daemon_id = DaemonId::new(Some("m1".to_string())); + let node_id: dora_core::config::NodeId = "camera".to_string().into(); + let value_bytes = serde_json::to_vec(&serde_json::json!(1)).unwrap(); + store + .put_node_param(&dataflow_id, &node_id, "threshold", &value_bytes) + .unwrap(); + + let mut dataflow = test_running_dataflow(dataflow_id, daemon_id.clone(), node_id); + dataflow.state_log_sequence = 10; + dataflow.daemon_ack_sequence.insert(daemon_id.clone(), 3); + + // Accepts commands but never replies. + let (tx, mut stalled_rx) = tokio::sync::mpsc::channel::(8); + let connection = crate::state::DaemonConnection::new( + tx, + Arc::new(tokio::sync::Mutex::new(HashMap::new())), + BTreeMap::new(), + ); + let connection_id = connection.connection_id; + let mut daemon_connections = DaemonConnections::default(); + daemon_connections.add(daemon_id.clone(), connection); + + let (events_tx, mut events_rx) = tokio::sync::mpsc::channel(8); + start_pruned_state_catchup_fallback( + dataflow_id, + &mut dataflow, + &daemon_id, + store.clone(), + &mut daemon_connections, + Arc::new(HLC::default()), + events_tx.clone(), + Instant::now(), + ); + + // Returned without an answer, with the replay running in the background. + tokio::time::timeout(Duration::from_secs(10), stalled_rx.recv()) + .await + .expect("replay must still be sent") + .expect("connection open"); + assert_eq!(dataflow.daemon_ack_sequence.get(&daemon_id), Some(&3)); + assert_eq!( + dataflow.fallback_replay_in_flight.get(&daemon_id), + Some(&connection_id) + ); + assert!(events_rx.try_recv().is_err(), "replay has not finished"); + + // A second status report on the same connection, past the backoff, + // does not start a duplicate replay while the first is in flight. + start_pruned_state_catchup_fallback( + dataflow_id, + &mut dataflow, + &daemon_id, + store, + &mut daemon_connections, + Arc::new(HLC::default()), + events_tx, + Instant::now() + FALLBACK_REPLAY_BACKOFF * 2, + ); + assert!( + tokio::time::timeout(TokioDuration::from_millis(50), stalled_rx.recv()) + .await + .is_err(), + "no duplicate replay while one is in flight" + ); +} + +#[tokio::test] +async fn fallback_replay_reports_back_and_acks_the_sequence_it_started_at() { + // #3684: once the replay is off the event loop, entries may be appended + // while it runs. The ack must advance to the sequence captured when the + // replay started, not the one current at completion. + #[derive(serde::Deserialize)] + struct OutboundRaw { + id: String, + } + + let store: Arc = Arc::new(InMemoryStore::new()); + let dataflow_id = DataflowId::from(Uuid::new_v4()); + let daemon_id = DaemonId::new(Some("m1".to_string())); + let node_id: dora_core::config::NodeId = "camera".to_string().into(); + let value_bytes = serde_json::to_vec(&serde_json::json!(1)).unwrap(); + store + .put_node_param(&dataflow_id, &node_id, "threshold", &value_bytes) + .unwrap(); + + let mut dataflow = test_running_dataflow(dataflow_id, daemon_id.clone(), node_id); + dataflow.state_log_sequence = 10; + dataflow.daemon_ack_sequence.insert(daemon_id.clone(), 3); + + let (tx, mut rx) = tokio::sync::mpsc::channel::(8); + let pending_replies = Arc::new(tokio::sync::Mutex::new(HashMap::new())); + let connection = + crate::state::DaemonConnection::new(tx, pending_replies.clone(), BTreeMap::new()); + let connection_id = connection.connection_id; + let mut daemon_connections = DaemonConnections::default(); + daemon_connections.add(daemon_id.clone(), connection); + + let (events_tx, mut events_rx) = tokio::sync::mpsc::channel(8); + start_pruned_state_catchup_fallback( + dataflow_id, + &mut dataflow, + &daemon_id, + store, + &mut daemon_connections, + Arc::new(HLC::default()), + events_tx, + Instant::now(), + ); + // A `dora param set` appended to the log while the replay is running. + dataflow.state_log_sequence = 12; + + let outbound = rx.recv().await.expect("replay command"); + let request_id: Uuid = serde_json::from_str::(&outbound) + .unwrap() + .id + .parse() + .unwrap(); + let reply = serde_json::to_string(&DaemonCoordinatorReply::SetParamResult(Ok(()))).unwrap(); + let _ = pending_replies + .lock() + .await + .remove(&request_id) + .expect("pending reply sender") + .send(reply); + + let event = tokio::time::timeout(Duration::from_secs(10), events_rx.recv()) + .await + .expect("replay must report back") + .expect("channel open"); + let Event::ParamFallbackReplayFinished { + dataflow_id: finished_df, + daemon_id: finished_daemon, + connection_id: finished_connection, + ack_sequence, + succeeded, + } = event + else { + panic!("unexpected event {event:?}"); + }; + assert_eq!(finished_df, dataflow_id); + assert_eq!(finished_daemon, daemon_id); + assert_eq!(finished_connection, connection_id); + assert_eq!(ack_sequence, 10); + assert!(succeeded); + + finish_pruned_state_catchup_fallback( + &mut dataflow, + &daemon_id, + finished_connection, + Some(connection_id), + ack_sequence, + succeeded, + ); + assert_eq!(dataflow.daemon_ack_sequence.get(&daemon_id), Some(&10)); + assert!(!dataflow.fallback_replay_in_flight.contains_key(&daemon_id)); +} + +#[test] +fn fallback_replay_result_from_a_replaced_connection_is_ignored() { + let dataflow_id = DataflowId::from(Uuid::new_v4()); + let daemon_id = DaemonId::new(Some("m1".to_string())); + let node_id: dora_core::config::NodeId = "camera".to_string().into(); + let mut dataflow = test_running_dataflow(dataflow_id, daemon_id.clone(), node_id); + dataflow.state_log_sequence = 10; + dataflow.daemon_ack_sequence.insert(daemon_id.clone(), 3); + let old_connection = Uuid::new_v4(); + let new_connection = Uuid::new_v4(); + // The reconnected daemon already has its own replay in flight. + dataflow + .fallback_replay_in_flight + .insert(daemon_id.clone(), new_connection); + + finish_pruned_state_catchup_fallback( + &mut dataflow, + &daemon_id, + old_connection, + Some(new_connection), + 10, + true, + ); + assert_eq!(dataflow.daemon_ack_sequence.get(&daemon_id), Some(&3)); + assert_eq!( + dataflow.fallback_replay_in_flight.get(&daemon_id), + Some(&new_connection) + ); +} + +#[test] +fn failed_fallback_replay_leaves_the_ack_and_never_moves_it_back() { + let dataflow_id = DataflowId::from(Uuid::new_v4()); + let daemon_id = DaemonId::new(Some("m1".to_string())); + let node_id: dora_core::config::NodeId = "camera".to_string().into(); + let mut dataflow = test_running_dataflow(dataflow_id, daemon_id.clone(), node_id); + let connection = Uuid::new_v4(); + dataflow.daemon_ack_sequence.insert(daemon_id.clone(), 3); + dataflow + .fallback_replay_in_flight + .insert(daemon_id.clone(), connection); + + finish_pruned_state_catchup_fallback( + &mut dataflow, + &daemon_id, + connection, + Some(connection), + 10, + false, + ); + assert_eq!(dataflow.daemon_ack_sequence.get(&daemon_id), Some(&3)); + assert!(!dataflow.fallback_replay_in_flight.contains_key(&daemon_id)); + + // A StateCatchUpAck past the replay's start sequence arrived meanwhile. + dataflow.daemon_ack_sequence.insert(daemon_id.clone(), 15); + finish_pruned_state_catchup_fallback( + &mut dataflow, + &daemon_id, + connection, + Some(connection), + 10, + true, + ); + assert_eq!(dataflow.daemon_ack_sequence.get(&daemon_id), Some(&15)); +} + #[tokio::test] async fn ready_boundary_schedules_replay_for_daemon_nodes() { #[derive(serde::Deserialize)] From 2763e8881deeaffa8e8e5a95a316af0fc63affb4 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 2 Oct 2026 03:33:29 +0000 Subject: [PATCH 2/2] chore(deps): bump yanked yoke-derive 0.8.3 to 0.8.4 Port of #3682: cargo-deny fails the Audit job on every PR because yoke-derive 0.8.3 was yanked. No-op once #3682 lands on main. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_016ZCEkLBowwvtAMhopshZwt --- Cargo.lock | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index a71c0f182c..f61c88e63c 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -9770,9 +9770,9 @@ dependencies = [ [[package]] name = "yoke-derive" -version = "0.8.3" +version = "0.8.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "33811428bee40dbceb6d545e95754741d17a6aef9a4849f0fd62e2ba4f412a78" +checksum = "ec8ebde2db3681e8c9980cc27822030e68752690ddfa9473e739aeb4dbde6d71" dependencies = [ "proc-macro2", "quote",