Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

37 changes: 31 additions & 6 deletions binaries/coordinator/src/daemon_events.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::{
Expand Down Expand Up @@ -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,
Expand Down
17 changes: 17 additions & 0 deletions binaries/coordinator/src/events.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -122,6 +138,7 @@ impl Event {
Event::DaemonStatusReport { .. } => "DaemonStatusReport",
Event::DaemonStateCatchUpAck { .. } => "DaemonStateCatchUpAck",
Event::DaemonNodeStopped { .. } => "DaemonNodeStopped",
Event::ParamFallbackReplayFinished { .. } => "ParamFallbackReplayFinished",
}
}
}
Expand Down
2 changes: 2 additions & 0 deletions binaries/coordinator/src/handlers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
Expand Down Expand Up @@ -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(),
Expand Down
26 changes: 23 additions & 3 deletions binaries/coordinator/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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<Event>,
}

impl Coordinator {
Expand Down Expand Up @@ -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::<Event>(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;

Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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
Expand Down
149 changes: 104 additions & 45 deletions binaries/coordinator/src/params.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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<dyn dora_coordinator_store::CoordinatorStore>,
daemon_connections: &mut DaemonConnections,
clock: Arc<HLC>,
events: mpsc::Sender<Event>,
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
{
Expand All @@ -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<Uuid>,
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
Expand Down
5 changes: 5 additions & 0 deletions binaries/coordinator/src/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<DaemonId, Instant>,
/// 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<DaemonId, Uuid>,

/// Whether UV was used for this dataflow (needed for restart).
pub(crate) uv: bool,
Expand Down Expand Up @@ -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(),
Expand Down
Loading
Loading