Skip to content

Worker: a stalled scheduler connection reconnects instead of leaving the worker evicted and idle - #2856

Merged
amankrx merged 10 commits into
TraceMachina:mainfrom
amankrx:fix/worker-scheduler-stall
Oct 1, 2026
Merged

amankrx merged 10 commits into
TraceMachina:mainfrom
amankrx:fix/worker-scheduler-stall

Conversation

@amankrx

@amankrx amankrx commented Oct 1, 2026 •

Copy link
Copy Markdown
Member

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 ConnectWorker stream, 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 GoingAway is 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 in RUN-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.

…alives and acknowledgements that reconnects, and HTTP/2 keepalive on the scheduler channel
@vercel

vercel Bot commented Oct 1, 2026 •

Copy link
Copy Markdown

The latest updates on your projects. Learn more about Vercel for GitHub.

Project Deployment Actions Updated
nativelink Ready Ready Preview Oct 1, 2026 2:47pm UTC
nativelink-aidm Ready Ready Preview Oct 1, 2026 2:47pm UTC

Request Review

@corcillo corcillo left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

  1. nativelink-config/src/cas_server.rs:880 (not in the diff): Problem: EndpointConfig.timeout now 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_S in local_worker.rs says that doc must track it. Fix: document each behavior it governs, or split the knobs as in comment 4.

Comment thread nativelink-worker/src/local_worker.rs Outdated
Comment thread nativelink-worker/src/local_worker.rs Outdated
Comment thread nativelink-worker/src/local_worker.rs Outdated
Comment thread nativelink-worker/src/local_worker.rs Outdated
Comment thread nativelink-worker/src/local_worker.rs Outdated
Comment thread nativelink-worker/src/local_worker.rs Outdated
@b7r6

b7r6 commented Oct 1, 2026

Copy link
Copy Markdown
Contributor

The lone red (asan) looks like it isn't a sanitizer or test failure, the tests pass. It's a compile
error building //:nativelink:

overflow evaluating the requirement Pin<Box<...>>: CoerceUnsized<_>

The run-loop restructuring deepens the async combinator nesting in local_worker::run (the
Metrics::wrap(… AndThen/Map/TryFlatten …) chain) past the coercion/type-length limit, and the
asan config tips over where Dev/Bazel don't. .boxed() on the inner async chain — breaking the
concrete type before the Box<dyn Future> unsize — usually clears it without touching the logic,
and it's cheaper than bumping type_length_limit/recursion_limit. Hit this exact shape recently,
figured I'd save you a round-trip.

Nice fix.

@MarcusSorealheis MarcusSorealheis left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Lgtm

@amankrx
amankrx merged commit 389cf73 into TraceMachina:main Oct 1, 2026
45 of 46 checks passed

@MarcusSorealheis MarcusSorealheis left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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());

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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();

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

b7r6 added a commit to b7r6/nativelink that referenced this pull request Oct 4, 2026
…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>
b7r6 added a commit to b7r6/nativelink that referenced this pull request Oct 7, 2026
…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>
b7r6 added a commit to b7r6/nativelink that referenced this pull request Oct 7, 2026
…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>
b7r6 added a commit to b7r6/nativelink that referenced this pull request Oct 7, 2026
…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>
b7r6 added a commit to b7r6/nativelink that referenced this pull request Oct 8, 2026
…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>

This branch was successfully deployed

2 active deployments
Preview – nativelink — 1d79d280 Deployed Oct 1, 2026 by vercel[bot]
Preview – nativelink-aidm — 1d79d280 Deployed Oct 1, 2026 by vercel[bot]
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.

4 participants