Skip to content

fix(coordinator): run the pruned-log fallback param replay off the event loop - #3686

Open
phil-opp wants to merge 2 commits into
mainfrom
claude/elegant-fermat-0cnwo7-fallback-replay-offloop
Open

phil-opp wants to merge 2 commits into
mainfrom
claude/elegant-fermat-0cnwo7-fallback-replay-offloop

Conversation

@phil-opp

@phil-opp phil-opp commented Oct 2, 2026 •

Copy link
Copy Markdown
Collaborator

Fixes #3684

Problem

Suppose a reconnecting daemon's state-log ack has been pruned. The daemon-status handler then called handle_pruned_state_catchup_fallback(...).await inline, on the coordinator's sequential event loop. That awaits replay_persisted_params_for_daemon, which sends every persisted param one send_and_receive at a time, each bounded only by TCP_READ_TIMEOUT (30s). A slow or half-dead daemon could stall the whole coordinator for N × 30s per affected dataflow. During that time it handled no CLI requests, no heartbeats, and no other daemons' events. The other two replay call sites already tokio::spawn the replay.

Fix

  • start_pruned_state_catchup_fallback (replaces handle_pruned_state_catchup_fallback) is now synchronous. It spawns the replay and returns immediately.
  • The spawned task reports back through a new internal Event::ParamFallbackReplayFinished. It travels on a coordinator-internal channel merged into the event stream (Coordinator::internal_events).
  • finish_pruned_state_catchup_fallback handles that event:
    • The ack advances to the state_log_sequence captured before the replay started, as the issue requires. Entries appended while the replay runs are not skipped. The ack never moves backwards past a StateCatchUpAck that arrived meanwhile.
    • The ack only advances if every param was replayed. That is unchanged.
    • The ack only advances if the daemon is still on the connection the replay was sent on. A result from a replaced connection says nothing about the daemon behind the new one.
  • RunningDataflow::fallback_replay_in_flight maps a daemon to the connection ID its replay is in flight on. A later status report on the same connection does not start a duplicate. A reconnect does start a fresh replay, so a replay stuck on a dead connection cannot block the new one.

Note: status reports are sent once per daemon connection. That is why the result comes back as an event rather than being polled at the next status report.

Related: #3683 / #3685 also touches replay_persisted_params_for_daemon. The two PRs are independent but will need a trivial merge (one extra argument at the call site) depending on merge order.

Also carries the yoke-derive 0.8.3 → 0.8.4 Cargo.lock bump from #3682, because the yanked crate fails the Audit job on every PR. It becomes a no-op once #3682 lands.

Validation

Class: C (coordinator)

  • cargo test -p dora-coordinator: ✅ 172 + 23 + 5 passed
  • New tests in binaries/coordinator/src/tests.rs:
    • fallback_replay_does_not_wait_for_an_unresponsive_daemon: the daemon never replies. The call returns, the replay is in flight, the ack is unchanged, and a second status report does not duplicate the replay.
    • fallback_replay_reports_back_and_acks_the_sequence_it_started_at: state_log_sequence goes from 10 to 12 while the replay runs, and the ack lands at 10.
    • fallback_replay_result_from_a_replaced_connection_is_ignored
    • failed_fallback_replay_leaves_the_ack_and_never_moves_it_back
  • Existing fallback_replay_keeps_ack_unchanged_when_daemon_is_disconnected and fallback_replay_respects_backoff_window were adapted to the new signature.
  • cargo clippy -p dora-coordinator --all-targets -- -D warnings: ✅, cargo fmt --check: ✅
  • make qa-fast: everything passes locally except audit/typos (not installed here) and breaking-changes (shallow clone, no tags). CI covers all three. Event gains a variant. dora-coordinator is not in the semver-checked crate list.
  • Fault-tolerance E2E (cargo test --test fault-tolerance-e2e -- --test-threads=1): ✅ 16 passed on 2763e88
  • Adversarial review: skipped. I did a focused manual review of the ack/connection race handling instead (see above).

🤖 Generated with Claude Code

https://claude.ai/code/session_016ZCEkLBowwvtAMhopshZwt

…ent 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 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_016ZCEkLBowwvtAMhopshZwt
@trunk-io

trunk-io Bot commented Oct 2, 2026

Copy link
Copy Markdown
Contributor

Merging to main in this repository is managed by Trunk.

  • To merge this pull request, check the box to the left or comment /trunk merge below.

After your PR is submitted to the merge queue, this comment will be automatically updated with its status. If the PR fails, failure details will also be posted here

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 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_016ZCEkLBowwvtAMhopshZwt

phil-opp commented Oct 2, 2026

Copy link
Copy Markdown
Collaborator Author

🤖 This is a fully automated review by Claude Code. No human has checked it.

I found one issue.

Running the fallback replay off the event loop opens the #3683 race on this path.
On main, the coordinator awaits handle_pruned_state_catchup_fallback inline on the event loop, so it can't interleave with handle_set_param/handle_delete_param. This PR moves it into a tokio::spawn (start_pruned_state_catchup_fallback in binaries/coordinator/src/params.rs). The spawned replay_persisted_params_for_daemon reads the store once (collect_param_replay_items) and then sends the items one round-trip at a time.

Example: a daemon reconnects with a pruned ack. The fallback replay starts with gain = 1 in its snapshot. While it is still waiting on earlier items, the user runs dora param set node gain 2. The event loop persists 2 and forwards it, and the daemon applies it. The replay then sends its stale gain = 1, and the daemon applies that too. A delete in the same window brings the key back. Advancing the ack to the sequence captured at the start doesn't repair this. The daemon counts as caught up, so nothing re-sends the newer value.

This path is race-free on main, so this would be a new regression. #3685 fixes the same race for the other two spawned replays. It would help to land this PR after #3685, or combine the two, and pass the dataflow's param_write_lock into the spawned replay. The two PRs currently conflict in params.rs, state.rs, handlers.rs and tests.rs. Note that #3685's lock is awaited on the event loop. With both PRs merged, a dora param set during a fallback replay to an unresponsive daemon would still stall the loop for up to 30 s. That is less than N × 30 s, but not zero. I raised this on #3685.

I checked the tests by reading them. I didn't run them with the fix reverted, because they call functions that don't exist on main. fallback_replay_does_not_wait_for_an_unresponsive_daemon covers the stall, and the ack-sequence and replaced-connection tests cover the new bookkeeping. None of them covers a set during the replay.


Generated by Claude Code

phil-opp commented Oct 3, 2026

Copy link
Copy Markdown
Collaborator Author

Automated review by Claude (fully automated, not a human review) — reviewed at 2763e88.

The race with dora param set/delete that I raised earlier (#3683, fixed for the other replays by #3685) is still open on this branch. One new point, about whether the fallback needs its own replay at all:

The pruned-log fallback replays the same params that the same status report has already sent to this daemon.
In handle_daemon_status_report (daemon_events.rs), a dataflow that the daemon reports as running goes through reestablish_running_dataflow. For an existing entry, that function returns df.ready_barrier_released. When it returns true, replay_all_nodes_ready spawns replay_persisted_params_for_daemon for exactly the nodes on this daemon (ready_barrier.rs:145). A few lines later, the state catch-up loop reaches the None (pruned) branch for the same daemon and starts a second full replay of the same store contents. When the barrier has not been released yet, schedule_param_replay_for_ready_dataflow replays to every daemon of the dataflow once it is. So every param the fallback sends reaches the daemon through another path anyway. The only thing the fallback adds is the daemon_ack_sequence update.

This means #3684 could likely be fixed with much less code. Drop the second replay in the pruned branch and only advance the ack, either from the outcome of the ready-barrier replay or on the next successful StateCatchUpAck. That avoids the new Event variant, the internal channel, and the fallback_replay_in_flight bookkeeping. It also halves the number of SetParam round trips to a reconnecting daemon. Once #3685 lands, it also stops two replays to the same daemon from queueing on param_write_lock against each other and against dora param set. If the separate replay is kept on purpose, a comment explaining why it is still needed next to the ready-barrier replay would help.

I found this by reading the code (daemon_events.rs status-report handler, daemon_liveness.rs:237-250, ready_barrier.rs:114-157). I did not run a test for it.

How the two PRs interact. Whichever lands second needs the param_write_lock argument at the new spawned call site in start_pruned_state_catchup_fallback. Without it, the fallback replay is exposed to the #3683 race again. The ack-at-start semantics here work with #3685's per-item re-read: a value newer than ack_sequence that gets re-sent is harmless.


Generated by Claude Code

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

coordinator: pruned-log fallback param replay runs inline in the event loop, so it can stall the coordinator for N × 30s

2 participants