Repository navigation
Worker: a stalled scheduler connection reconnects instead of leaving the worker evicted and idle - #2856
Conversation
…alives and acknowledgements that reconnects, and HTTP/2 keepalive on the scheduler channel
|
The latest updates on your projects. Learn more about Vercel for GitHub.
|
corcillo
left a comment
There was a problem hiding this comment.
Seven comments; two blocking. The loop restructure (never await a send inline, since this loop is the stream's only reader) and the bounded shutdown GoingAway are right. The deadlines on every send are at the wrong layer: see the two blocking comments.
nativelink-config/src/cas_server.rs:880(not in the diff): Problem:EndpointConfig.timeoutnow governs the send deadlines, the HTTP/2 PING interval, the PING timeout and the shutdown bound, and its doc still says "timeout a request should take";DEFAULT_ENDPOINT_TIMEOUT_Sin local_worker.rs says that doc must track it. Fix: document each behavior it governs, or split the knobs as in comment 4.
… no deadlines on stream messages
|
The lone red (asan) looks like it isn't a sanitizer or test failure, the tests pass. It's a compile The run-loop restructuring deepens the async combinator nesting in Nice fix. |
MarcusSorealheis
left a comment
There was a problem hiding this comment.
Reviewed at 1d79d28 with cluster stability in mind: many workers on one shared scheduler, rolling restarts, and backend latency spikes.
The latest push is a big improvement over the first version. Dropping the per-message deadlines removes the risk of a slow-but-healthy scheduler making every worker kill its in-flight actions at once, and the 30 s / 20 s HTTP/2 keepalive defaults avoid correlated disconnects across the fleet. Reading the ConnectWorker stream on its own task, and never awaiting a send inline in the worker loop, removes the circular wait described in the bed report.
Verified locally: local_worker_test 24/24 and worker_api_server_test 14/14 pass, and CI is green.
Two things inline are worth addressing for cluster deployments, plus a nit. Separately, the filesystem_store try_update change looks unrelated to this fix; fine to keep if the toolchain needs it, but worth a line in the description.
| // action once it is acknowledged. A decline above | ||
| // must not spend it, or every dispatch it turns | ||
| // away costs a pod. | ||
| let spent_mark = self.config.single_use.then(|| self.accepted_action.clone()); |
There was a problem hiding this comment.
Single-use isolation. This moves the "spent" mark from the dispatch to after the ExecuteAccepted send completes. The only local guard against a second action is accepted_action (line 535): registration advertises one slot to the scheduler, but admission() checks max_inflight_tasks, which isn't forced to 1 for single-use workers.
So from the moment action #1 is dispatched until its acknowledgement has gone through the one-slot channel, a second StartExecute passes the guard and runs in the same container. That window can now be long, because the acknowledgement queues behind anything else this worker is sending. With a scheduler that doesn't speak the acknowledgement, the flag is set only when this future is first polled, and select! can read the stream before that.
The scheduler shouldn't send that second action, but this guard is the defence in depth behind the single-use promise ("never accepts a second action"). Suggest setting accepted_action synchronously here, before the loop can read the next message, which the old inline await guaranteed. If the acknowledgement then fails, the loop ends and the spent worker exits instead of reconnecting: that costs one pod, but keeps the isolation guarantee.
| // purpose: a worker's messages are bounded by its in-flight | ||
| // actions and its keepalive interval, and backpressure here is | ||
| // exactly the stall. | ||
| let (tx, mut rx) = mpsc::unbounded_channel(); |
There was a problem hiding this comment.
Liveness still waits behind processing. Moving the reads off the processor fixes the flow-control stall, but the worker's liveness is still refreshed only when a message is processed (the keepalive handling and touch_liveness run in the loop below). A keepalive that arrives on time therefore waits behind this worker's queued results before it counts.
During a backend latency spike (a Redis failover, a slow update_action write), this queue grows the same way on every worker connected to the scheduler. Once it's longer than worker_timeout_s (10 s by default), the scheduler evicts workers whose keepalives it has already received, and requeues all their actions onto the same slow scheduler.
Suggest recording "last heard" in this reader task on receipt (an atomic timestamp the timeout sweep consults), or handling KeepAliveRequest on a fast path here, so liveness reflects the connection rather than the backend. This could be a follow-up, but it's the remaining way a healthy worker gets evicted under load.
| @@ -319,13 +362,47 @@ impl<'a, T: WorkerApiClientTrait + 'static, U: RunningActionsManager> LocalWorke | |||
|
|
|||
| /// Tells the scheduler the action is refused; only meaningful when it | |||
There was a problem hiding this comment.
Nit: these two lines describe decline, but they now sit above going_away_deadline, so its rustdoc starts with them and decline has none. Moving them down to fn decline fixes both.
…ner.lock (TL1)
Registry is the second liveness source remove_timedout_workers consults ("alive
if EITHER source"), but refresh_worker stamped it only AFTER the inner.lock block
— which holds the lock across an await (unacked-dispatch sweep). A punctual
keepalive whose inner.lock is contended, or whose task is starved on a saturated
runtime, refreshed NEITHER liveness source before the eviction sweep read them,
falsely evicting a live worker. Transport/runtime-starvation liveness class,
seeded by Aman TraceMachina#2856 (h2-read starvation); TL1 is the sibling where a read
keepalive must still LAND through the lock.
Fix: stamp update_worker_heartbeat on receipt (lightweight RwLock insert, no
await-under-lock/store/sweep); make it monotone so an out-of-order keepalive
cannot regress last_seen. Regression heartbeat_is_monotone_never_regresses is
fails-without by construction; decoupling half exercised by the soak harness.
NOT YET built/run on metal.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
…ner.lock (TL1)
Registry is the second liveness source remove_timedout_workers consults ("alive
if EITHER source"), but refresh_worker stamped it only AFTER the inner.lock block
— which holds the lock across an await (unacked-dispatch sweep). A punctual
keepalive whose inner.lock is contended, or whose task is starved on a saturated
runtime, refreshed NEITHER liveness source before the eviction sweep read them,
falsely evicting a live worker. Transport/runtime-starvation liveness class,
seeded by Aman TraceMachina#2856 (h2-read starvation); TL1 is the sibling where a read
keepalive must still LAND through the lock.
Fix: stamp update_worker_heartbeat on receipt (lightweight RwLock insert, no
await-under-lock/store/sweep); make it monotone so an out-of-order keepalive
cannot regress last_seen. Regression heartbeat_is_monotone_never_regresses is
fails-without by construction; decoupling half exercised by the soak harness.
NOT YET built/run on metal.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
…ner.lock (TL1)
Registry is the second liveness source remove_timedout_workers consults ("alive
if EITHER source"), but refresh_worker stamped it only AFTER the inner.lock block
— which holds the lock across an await (unacked-dispatch sweep). A punctual
keepalive whose inner.lock is contended, or whose task is starved on a saturated
runtime, refreshed NEITHER liveness source before the eviction sweep read them,
falsely evicting a live worker. Transport/runtime-starvation liveness class,
seeded by Aman TraceMachina#2856 (h2-read starvation); TL1 is the sibling where a read
keepalive must still LAND through the lock.
Fix: stamp update_worker_heartbeat on receipt (lightweight RwLock insert, no
await-under-lock/store/sweep); make it monotone so an out-of-order keepalive
cannot regress last_seen. Regression heartbeat_is_monotone_never_regresses is
fails-without by construction; decoupling half exercised by the soak harness.
NOT YET built/run on metal.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
…ner.lock (TL1)
Registry is the second liveness source remove_timedout_workers consults ("alive
if EITHER source"), but refresh_worker stamped it only AFTER the inner.lock block
— which holds the lock across an await (unacked-dispatch sweep). A punctual
keepalive whose inner.lock is contended, or whose task is starved on a saturated
runtime, refreshed NEITHER liveness source before the eviction sweep read them,
falsely evicting a live worker. Transport/runtime-starvation liveness class,
seeded by Aman TraceMachina#2856 (h2-read starvation); TL1 is the sibling where a read
keepalive must still LAND through the lock.
Fix: stamp update_worker_heartbeat on receipt (lightweight RwLock insert, no
await-under-lock/store/sweep); make it monotone so an out-of-order keepalive
cannot regress last_seen. Regression heartbeat_is_monotone_never_regresses is
fails-without by construction; decoupling half exercised by the soak harness.
NOT YET built/run on metal.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
…ner.lock (TL1)
Registry is the second liveness source remove_timedout_workers consults ("alive
if EITHER source"), but refresh_worker stamped it only AFTER the inner.lock block
— which holds the lock across an await (unacked-dispatch sweep). A punctual
keepalive whose inner.lock is contended, or whose task is starved on a saturated
runtime, refreshed NEITHER liveness source before the eviction sweep read them,
falsely evicting a live worker. Transport/runtime-starvation liveness class,
seeded by Aman TraceMachina#2856 (h2-read starvation); TL1 is the sibling where a read
keepalive must still LAND through the lock.
Fix: stamp update_worker_heartbeat on receipt (lightweight RwLock insert, no
await-under-lock/store/sweep); make it monotone so an out-of-order keepalive
cannot regress last_seen. Regression heartbeat_is_monotone_never_regresses is
fails-without by construction; decoupling half exercised by the soak harness.
NOT YET built/run on metal.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
What and why
On the test bed a worker a few minutes old was evicted for a missed keepalive right after its first burst of completions, stayed idle, and never reconnected. Everything a worker sends is a message on its one
ConnectWorkerstream, and the scheduler read that stream inline, running each message to completion before the next read, so a burst filled the flow-control window behind a slow message and the keepalives waited there until the eviction. The scheduler now reads the stream in a task that only hands messages on and processes them in order in another, and the worker's run loop, the stream's only reader, no longer awaits any send inline, so it sees the scheduler close the stream and reconnects.An action starts once its acknowledgement has gone out, so a failed one leaves the action unstarted and a single-use worker unspent. There are no deadlines on stream messages, the scheduler's liveness timeout is the eviction clock; only the shutdown
GoingAwayis bounded, and a miss still drains in-flight actions. The scheduler connection gets the TCP and HTTP/2 keepalive knobs a store endpoint has.How was this verified?
worker_api_server_test: a result whose store write is held no longer stalls the stream reads, and fails on the old reader.local_worker_test: a hung acknowledgement leaves the loop reading and the action unstarted, a failed one leaves a single-use worker unspent, a hung keepalive waits for the stream close. The bed rerun is inRUN-rc-20260930.md, finding 2.Risk
A worker's messages queue on the scheduler without backpressure, bounded by its in-flight actions and keepalive interval. Nothing changes for a stream that never stalls.
AI assistance
An agent drafted the change and I reviewed every line.