Repository navigation
Conversation
|
The latest updates on your projects. Learn more about Vercel for GitHub.
|
…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
force-pushed
the
head-of-line-reservation
branch
from
October 9, 2026 23:14
8bf1ee2 to
2030ea8
Compare
|
|
This branch was successfully deployed
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
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_fitonly slows it). This addshead_of_line_reservation: { after_s }toSimpleSpec: once the action atthe head of the queue has waited
after_ssince it was last queued and apass 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
ApiWorkerSchedulerImplbeside the ledger(
HeadOfLineReservation { operation_id, worker_id, platform_properties }),and
find_available_workerrefuses the held worker to every action but theheld 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_ssince it was last queued (insert_timestampor the state'slast_transition_timestamp, whichever is later, so a requeue starts theclock 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 totalssatisfy 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 onworker_notify_run_actionfor thatoperation (
placed, anywhere) and at the end of a pass whose listing wasread to its end and did not name the operation (
gone; a record the storecannot decode is a load failure and keeps the hold, only
NotFoundcountsas absent). When the held worker is removed, drained or paused
(
remove_worker,set_drain_worker, the backpressure pause, thechannel-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 theheld 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 lastreported free while idle (
Worker::idle_free_kb, set by any idle keepaliveand 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 inexplanations/scheduler-internalsafter the parked-action one, and theregenerated 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: atwelve-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
Queuedafter 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 isheld; 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) andthe waiting small action takes the core.
a_worker_disconnect_moves_the_hold: the held worker disconnects; thehold moves to the survivor for the same operation at once (
moved); thetwelve 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_workeron theheld 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 heldaction 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 ofthe 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(zoneexact): 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(two100 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'srust-version):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_statethe dispatch pathalready 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
goneon the next pass).Risk
None for anyone who does not set the field:
head_of_line_afterisNone,reserve_head_of_lineis never called, the reservation is alwaysNone,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 withthe 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.
WorkerSummarygains
reserved_for_operation, omitted from the JSON when unset, so theadmin listing is byte-identical until a hold exists.
ActionStateResultgains
as_action_info_with_statewith a default, so other implementationsare 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_son afleet 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_sabove 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, nominimumshortfall) that isovertaken 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 thenext 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
unfitonce drained, after which itsidle report caps the estimate.
find_worker_for_action_observedgainsoperation_idandahead_of_holdparameters; the three-argumentfind_worker_for_actionis 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.