fix(coordinator): answer a log subscription before replaying its backlog - #3676
Conversation
`LogSubscribe`/`BuildLogSubscribe` replayed every buffered log message (up to 10 000) into the new subscriber's 64-slot channel before sending `found_tx`. The WS task that drains that channel is parked on `found_rx` until then, so after 64 messages each send took the 100 ms timeout: `dora logs --follow` on a busy detached dataflow stalled the whole coordinator event loop for ~10 s, then evicted the subscriber with the rest of the backlog already taken and lost. Register the subscriber and answer `found_tx` first, then replay. The WS task sends the subscribe reply and resumes draining, so the replay flows instead of timing out. Both arms now share `attach_log_subscriber`. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01R8CSx6yQ4JbrB3dqSiU2ME
|
😎 Merged successfully - details. |
Automated review batch — 2026-10-01 (base
|
`cargo deny check advisories` fails every PR's Audit job with `error[yanked]: detected yanked crate` for yoke-derive 0.8.3 (via url -> idna -> icu). `cargo update -p yoke-derive` moves it to 0.8.4; no other lock entry changes. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01R8CSx6yQ4JbrB3dqSiU2ME (cherry picked from commit 2132289)
|
Audit (cargo-audit + cargo-deny) failed with Generated by Claude Code |
|
🤖 This is a fully automated review by Claude Code. No human has checked it. I didn't find any issues. The diagnosis is correct. The WS task in I ran the new test with the fix reverted ( Generated by Claude Code |
Problem
ControlEvent::LogSubscribe(andBuildLogSubscribe) registered the new subscriber, then replayed the whole buffered backlog (up toMAX_BUFFERED_LOG_MESSAGES= 10 000) into it, and only then answeredfound_tx.The other end of that channel is the WS task in
ws_control.rs. It is parked onfound_rx.awaitand does not drainlog_rx(capacity 64) untilfound_txfires. So during the replay:send_log_messagehits its 100 ms timeout;taken, so the rest of it is lost, andfound_txstill answerstrue.How to hit it:
dora start --detacha dataflow that logs a few hundred lines, then rundora logs <df> --follow. You get the first 64 lines and then nothing; meanwhile the whole coordinator was frozen for ~10 s.Fix
Register the subscriber and answer
found_txbefore replaying. The WS task then sends the subscribe reply and returns to itsselect!, where it drainslog_rxwhile the replay flows. Ordering is unchanged: anything queued before the reply is forwarded after it. Both arms now share a smallattach_log_subscriberhelper.Validation
Class: C (
binaries/coordinator/**)handlers::tests::attach_log_subscriber_delivers_backlog_larger_than_channel. A 200-message backlog goes into a 4-slot channel, with a task that, likews_control, waits forfound_rxbefore draining. It asserts all 200 arrive and the subscriber is kept.assertion failed: subscriber must not be evicted.cargo test -p dora-coordinator: ✅ 169 + 23 + 5 passed.cargo clippy -p dora-coordinator --all-targets -D warnings: ✅.cargo fmt --check: ✅.cargo test --all, fault-tolerance E2E, contract tests,cargo check --examples) were run on the combined diff of this review batch. Results are in the batch summary comment.Not addressed
The replay still runs on the coordinator event loop, one awaited send per message. If the same WS session sends another request mid-replay (that arm awaits its reply without draining
log_rx), the stall can still happen. A per-subscriber forwarding task, ortry_sendplus re-buffering, would remove it entirely. That is a larger change, left out to keep this one minimal; the commondora logs --followpath sends nothing else.🤖 Generated with Claude Code
https://claude.ai/code/session_01R8CSx6yQ4JbrB3dqSiU2ME
Generated by Claude Code