Skip to content

feat(scheduler): reserve a worker for the head of the queue after it waits - #2920

Open
rohnnyjoy wants to merge 1 commit into
TraceMachina:mainfrom
rohnnyjoy:head-of-line-reservation
Open

rohnnyjoy wants to merge 1 commit into
TraceMachina:mainfrom
rohnnyjoy:head-of-line-reservation

Conversation

@rohnnyjoy

Copy link
Copy Markdown
Contributor

What and why

On a busy fleet an action larger than the rest never finds room: every core
that frees is taken by the next small action before the pass reaches the big
one again, and it waits as long as the stream lasts (four 8-core actions
queued for an hour on two 16-core workers, No workers matched! ~400/min,
small actions running throughout; best_fit only slows it). This adds
head_of_line_reservation: { after_s } to SimpleSpec: once the action at
the head of the queue has waited after_s since it was last queued and a
pass places something queued after it, the busy worker closest to fitting it
(fewest cores and KiB short, by share of what it advertises, among those
that could take it once drained) is held for it and takes nothing queued
after the action until it is placed, so it drains to the action while the
rest of the fleet keeps backfilling. Slurm's backfill reservation without
the time estimate. Default unset keeps today's behaviour exactly.

Fixes #2919.

The hold lives in ApiWorkerSchedulerImpl beside the ledger
(HeadOfLineReservation { operation_id, worker_id, platform_properties }),
and find_available_worker refuses the held worker to every action but the
held one and the ones queued before it, so neither the main walk nor the
parked replay can dispatch around it. The matching pass makes it at its
end: the first action still parked is the head of the queue, and if it has
waited after_s since it was last queued (insert_timestamp or the state's
last_transition_timestamp, whichever is later, so a requeue starts the
clock again) and the pass placed something queued after it, the pass asks
the worker scheduler to hold. A fleet that places nothing holds nothing, and
an action that is simply next in line is never held for, so a batch of
small actions requeued after a worker loss produces no holds at all. One
predicate, could_hold_for (not draining, not paused, registered totals
satisfy the action's current properties, no idle decline of that size, and
under the live memory veto a free-memory report that, plus what the ledger
had lent to the worker's running actions when it reported
(Worker::reported_lent_kb, so a stale report stays true as they finish),
covers the action), decides both who
can be held and whether a hold still stands. The pick is the smallest
shortfall among the busy workers it accepts and that are not on record as
having refused the action idle since they last spoke; ties go to the fewest
running actions, then least recently used. An idle worker is never held:
one that could take the action would have. When nothing qualifies that is
logged once per action with the count each exclusion took (no_candidate).

One hold at a time, and it follows the queue. An action that sorts before
the held one takes the held worker's room first (FIFO; no event), and takes
the hold for itself only if the held worker could never run it
(superseded). The hold is released on worker_notify_run_action for that
operation (placed, anywhere) and at the end of a pass whose listing was
read to its end and did not name the operation (gone; a record the store
cannot decode is a load failure and keeps the hold, only NotFound counts
as absent). When the held worker is removed, drained or paused
(remove_worker, set_drain_worker, the backpressure pause, the
channel-full pause, a decline) the hold moves to another worker for the
same action with the properties it was made with (moved), or is dropped
(worker_lost, unfit) when none qualifies. Every time the pass meets the
held action and finds no room it re-checks the held worker against the
action's current properties; a worker that no longer passes, or is idle and
still refusing the action, is released (unfit). What a worker last
reported free while idle (Worker::idle_free_kb, set by any idle keepalive
and by that refusal) caps the lent-back projection, so two workers that
refuse the action idle are not drained for it in turn however much their
next running actions borrow; only an idle report with more lifts the cap.
The action info, the operation id and the queued-since time come from the
one read the pass already does (ActionStateResult::as_action_info_with_state,
one borrow of the record in the state manager; the trait default reads
twice), and the dispatch and unsatisfiable paths use that id too, so the
hold costs no extra read on any queue; with the option off no lock is taken
and no id compared against anything but None.

Docs: a section in operate/scaling-workers (the symptom, the mechanism,
the cost, how to set after_s, what to watch), a paragraph in
explanations/scheduler-internals after the parked-action one, and the
regenerated configuration reference (reference/nativelink-config/main.mdx)
so the snippet lint knows the field. The metrics reference is not regenerated:
it is pinned at v1.7.2 and regenerating it rewrites every source line number
on the page.

How was this verified?

Unit tests, in nativelink-scheduler/tests/head_of_line_reservation_test.rs,
with the mock clock, mock workers and the memory state manager. Every hold in
them is made the way the defect makes one: a core frees and a new one-core
action takes it past the waiting big one.

  • without_the_option_a_stream_of_small_actions_starves_a_large_one: a
    twelve-core worker full of one-core actions, an eight-core action queued;
    twenty rounds of a core freeing and a new one-core action taking it; the
    eight-core action is still Queued after every round and nothing is held.
    This is the defect, and it passes with the option off.
  • after_the_wait_a_worker_is_held_and_drains_to_the_large_action
    (after_s: 20, two workers): overtaken at 10 s on both workers, nothing is
    held; overtaken again at 20 s, a worker is held for the big action; a new
    small action is not given that worker's free core and waits; a core
    freeing on the other worker takes it; six more cores freeing on the held
    worker place nothing; the eighth places the big action, the hold is gone
    and the next core that worker frees goes to a small action again.
  • the_hold_is_released_when_the_action_leaves_the_queue (after_s: 0):
    held once overtaken; a freed core is not offered to a waiting small
    action; the client is not heard from for the client timeout, the state
    manager retires the action, the next pass releases the hold (gone) and
    the waiting small action takes the core.
  • a_worker_disconnect_moves_the_hold: the held worker disconnects; the
    hold moves to the survivor for the same operation at once (moved); the
    twelve requeued one-core actions, queued before the big one, take the
    survivor's cores one each as they free, in queue order, with the hold
    unmoved; then the survivor drains to the big action.
  • a_drain_of_the_held_worker_moves_the_hold: set_drain_worker on the
    held worker moves the hold to the other worker at once; the drained
    worker's cores free with nothing landing on them while the survivor drains
    to the big action and places it.
  • a_worker_that_declined_the_action_while_idle_is_not_held_for_it_again
    (memory budgets and the live veto on): the held worker drains, takes the
    action, declines it for load holding nothing else; overtaken again, the
    hold goes to the other, busy worker and the first is never offered the
    action again, also after a keepalive lifts its pause.
  • a_requeue_that_outgrows_the_held_worker_releases_the_hold: the held
    action asks, through the worker scheduler directly (the memory state
    manager has no peer to run an escalation on), for more memory than the
    held worker registered; the hold is released rather than waited on.
  • an_action_ahead_of_the_hold_takes_the_held_workers_room_first
    (after_s: 20): a hold for a priority-0 action; a priority-1 action of
    the same size arrives; the hold stays, and when the held worker has eight
    free cores the priority-1 action takes them, the hold still the other's.
  • an_action_ahead_that_the_held_worker_cannot_run_takes_the_hold (zone
    exact): a priority-1 action pinned to the other worker's zone arrives and
    is overtaken; the hold moves to the other worker, for it (superseded).
  • workers_that_refuse_the_action_idle_are_not_drained_in_turn (two
    100 000 KiB workers reporting 50 000 KiB free while busy, an action
    asking 60 000): a busy worker is held, since its running actions'
    12 000 KiB would cover it once given back; drained and still reporting
    50 000, it is released; the other is held and drained to the same end;
    both refilled with twelve one-core actions (so the lent-back projection
    alone would again say 62 000), four more overtakings, with no new report
    and then with busy reports of the same figure, hold nothing; a worker
    drained and reporting 70 000 idle takes the action with no hold.
  • a_batch_requeued_ahead_is_placed_in_order_without_a_hold (after_s: 0):
    a worker is lost with twelve one-core actions, which requeue ahead of a
    waiting big action; the survivor's cores take them one at a time in order
    with nothing held at any point, then the big action lands.

Without the change the hold tests fail on their first reservation(..)
assertion (the field is never set) and the small action takes the freed
core; the drain, decline, outgrow, ahead and refuse tests fail on the hold
staying put, moving, or draining the second worker in turn.

Ran, on rust 1.97.1 (the workspace's rust-version):

cargo +1.97.1 test -p nativelink-scheduler --test head_of_line_reservation_test
test a_requeue_that_outgrows_the_held_worker_releases_the_hold ... ok
test the_hold_is_released_when_the_action_leaves_the_queue ... ok
test an_action_ahead_of_the_hold_takes_the_held_workers_room_first ... ok
test a_worker_that_declined_the_action_while_idle_is_not_held_for_it_again ... ok
test an_action_ahead_that_the_held_worker_cannot_run_takes_the_hold ... ok
test a_drain_of_the_held_worker_moves_the_hold ... ok
test after_the_wait_a_worker_is_held_and_drains_to_the_large_action ... ok
test without_the_option_a_stream_of_small_actions_starves_a_large_one ... ok
test a_worker_disconnect_moves_the_hold ... ok
test a_batch_requeued_ahead_is_placed_in_order_without_a_hold ... ok
test workers_that_refuse_the_action_idle_are_not_drained_in_turn ... ok
test result: ok. 11 passed; 0 failed; 0 ignored; 0 measured; 0 filtered out; finished in 0.03s

cargo +1.97.1 test -p nativelink-scheduler        # whole crate, 23 test binaries, all ok
# simple_scheduler_test 38 ok, unsatisfiable_action_test 41 ok,
# simple_scheduler_state_manager_test 30 ok, store_awaited_action_db_test 14 ok, ...
cargo +1.97.1 test -p nativelink-util --test metrics_test   # 17 ok
cargo +1.97.1 test -p nativelink-config                     # 11 + 34 + 1 + 1 ok

cargo +1.97.1 check --workspace --all-targets
# Finished `dev` profile [unoptimized + debuginfo] target(s)
cargo +1.97.1 clippy -p nativelink-scheduler -p nativelink-config \
  -p nativelink-util --all-targets
# Finished, no finding on any touched target
rustfmt +nightly-2026-08-26 --edition 2024 --check <each changed file>   # no diff
cd web/apps/docs && node scripts/lint-snippets.mjs
# Checked 98 config snippet(s) across 83 page(s) against 317 known keys; No drift found.
node scripts/gen-config-reference.mjs main   # the new row, the new type, the header sha

Live: the v1.7.5 carry of this commit runs on the two-worker fleet the issue
describes. Not verified: against a Redis-backed queue beyond the existing
store tests (the id reads go through the same as_state the dispatch path
already uses), or with several schedulers sharing one state (each keeps its
own reservation for its own workers; a peer placing the action first shows up
as gone on the next pass).

Risk

None for anyone who does not set the field: head_of_line_after is None,
reserve_head_of_line is never called, the reservation is always None,
the pass takes no worker lock for it, and every new check is a comparison
against None. The pass's one read per action now returns the state with
the action info; the state manager serves both from the same borrow, and the
dispatch and unsatisfiable paths read the id from it instead of borrowing
again, so a store-backed pass does fewer reads than before. WorkerSummary
gains reserved_for_operation, omitted from the JSON when unset, so the
admin listing is byte-identical until a hold exists. ActionStateResult
gains as_action_info_with_state with a default, so other implementations
are unchanged.

With it set, the cost is the one the design chooses: a held worker idles
every core it frees until the action lands, up to the action's size for as
long as the longest action running on that worker when it was held. On a
two-worker fleet that is half the fleet's backfilling for the length of
that action, an hour for an hour-long one, and with a small after_s on a
fleet where large actions are common some worker is held almost all the
time. The field doc and the operator page say so and say to set after_s
above the typical queue wait. Further: (1) a hold needs the action to be
overtaken, so an action that nothing newer ever passes (a full fleet, a
batch requeued ahead of it) is never held for and simply waits its turn; on
such a fleet the option does nothing. (2) A head-of-line action waiting only
on a concurrency cap (max_inflight_tasks, no minimum shortfall) that is
overtaken still holds the fewest-loaded worker, a whole worker for one slot.
(3) A listing that transiently omits the held operation (an eventually
consistent search) releases the hold as gone; it is held for again the
next time it is overtaken. (4) A hold can move: a held worker paused by
backpressure hands the hold to another busy worker, and a hold whose worker
was drained and still refused the action is dropped until the action is
overtaken again, with that worker's idle report capping what a drain is
expected to free from then on; each is one log line and one counter event.
(5) The live-memory projection trusts the ledger's reservations until a
worker has reported idle: a worker whose running actions declared too
little, and which has never reported idle, is held on an estimate that will
not come true, and is released as unfit once drained, after which its
idle report caps the estimate. find_worker_for_action_observed gains
operation_id and ahead_of_hold parameters; the three-argument
find_worker_for_action is unchanged.

AI assistance

An agent drafted the change, the tests and the docs from a written design; I
reviewed every line and ran the verification above.

@vercel

vercel Bot commented Oct 9, 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 9, 2026 11:15pm UTC
nativelink-aidm Ready Ready Preview Oct 9, 2026 11:15pm UTC

Request Review

…tarved

On a busy fleet an action larger than the rest never finds room: every
core that frees is taken by the next small action before the pass
reaches the big one again, and it waits as long as the stream lasts
(four 8-core actions queued for an hour on two 16-core workers, "No
workers matched!" ~400/min, small actions running throughout;
`best_fit` only slows it down).

Add `head_of_line_reservation: { after_s }` to `SimpleSpec`. Once the
action at the head of the queue has waited `after_s` since it was last
queued and a matching pass places something queued after it, the busy
worker closest to fitting it (the fewest cores and KiB short, by share
of what it advertises, among those that could take it once drained) is
held for it and takes nothing queued after the action until the action
is placed, so it drains to the action while the rest of the fleet keeps
backfilling. This is Slurm's backfill reservation without the time
estimate. A full fleet that places nothing holds nothing, and neither
does an action that is simply next in line, so a batch requeued after
a worker loss starts no holds.

One hold at a time, and it follows the queue: an action that sorts
before the held one takes the held worker's room first, in queue
order, and takes the hold itself only if the held worker could never
run it. The hold is part of the ledger, so no placement path dispatches
around it. One predicate decides who can be held and whether a hold
still stands (not draining or paused, registered totals satisfy the
action's current properties, no idle decline of that size, and under
the live memory veto a free-memory report that plus what the ledger
lent to the worker's running actions covers the action); it is checked
every time the pass meets the held action, and a worker that no longer
passes, or is idle and still refusing the action, is released; the
projection counts what the ledger had lent when the worker reported and
is capped at what the worker last reported free while idle, so two
refusing workers are not drained for the action in turn however much
their next running actions borrow. A held worker that is
removed, drained or paused hands the hold to another worker for the
same action, or drops it when none qualifies. The hold also lifts when
the action is placed, there or elsewhere, and when it leaves the queue
unplaced; a record the store cannot decode keeps the hold, only a
missing one counts as gone. The action info, the operation id and the
queued-since time come from the one read the pass already does, and the
dispatch and unsatisfiable paths use that id too, so the hold costs no
extra read; with the option off nothing is locked or read for it.

The cost is idle capacity, by design: a held worker idles every core it
frees until the action lands, up to the action's size for as long as
the longest action running on it when held. Holds are logged at info
with the operation, the worker and the seconds waited;
`scheduler.head_of_line.reservations` counts them by
`scheduler.reservation.event` (reserved, moved, placed, gone,
worker_lost, unfit, superseded, no_candidate); the admin worker listing
shows `reserved_for_operation`. Default unset keeps today's behaviour
exactly.
@rohnnyjoy
rohnnyjoy force-pushed the head-of-line-reservation branch from 8bf1ee2 to 2030ea8 Compare October 9, 2026 23:14
@CLAassistant

Copy link
Copy Markdown

CLA assistant check
Thank you for your submission! We really appreciate it. Like many open source projects, we ask that you sign our Contributor License Agreement before we can accept your contribution.
You have signed the CLA already but the status is still pending? Let us recheck it.

This branch was successfully deployed

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

[Feature]: A large action starves behind a stream of small ones; reserve a worker for the head of the queue

3 participants