NativeLink
Concepts

Scheduler internals

Action merging, the awaited-action database, how a worker is chosen, and what happens when an action fails.

Who this is for: anyone tuning a scheduler, debugging an action that queued forever, or wondering why two identical builds did not share work. What you'll have at the end: the merge rule, the ordering rule, the matching rule, and the retry rule, each of which has at least one counter-intuitive edge. Time: twenty-five minutes.

The scheduler's job sounds small: take actions from clients, give them to workers. Almost all of its complexity comes from three questions it has to answer that the protocol does not specify: when are two actions the same one, which of the queued actions goes next, and what does "this failed" mean.

No ActionScheduler trait exists

If you go looking for one (and older material about NativeLink will send you looking) you will not find it. The abstraction is split by who is asking, into three traits in operation_state_manager.rs:

ClientStateManager is the client's view: submit an action, subscribe to its progress. WorkerStateManager is the worker's view: report that an operation changed state. MatchingEngineStateManager is the matcher's view: list actions that could be dispatched, assign one to a worker.

Two more sit beside them: WorkerScheduler, which owns worker lifecycle, and KnownPlatformPropertyProvider, which is how the capabilities service learns which platform property keys an instance's scheduler knows about.

The split is what allows the whole state layer to be swapped: SimpleSchedulerStateManager implements all three traits and is generic over an AwaitedActionDb, which is either MemoryAwaitedActionDb or the store-backed StoreAwaitedActionDb used for Redis. It is also why "the scheduler" in configuration can be a stack (the cache_lookup decorator is optional, and none of the shipped example configs use it):

The decorator in front is why a cache hit never appears in queue depth, and why ActionStage::CacheCheck is never emitted by the scheduler crate.

When are two actions the same action?

Two clients asking for the same work should wait on one execution. The key that decides this is ActionUniqueKey, and it has exactly three fields: instance name, digest function, and the action digest.

That key is wrapped in ActionUniqueQualifier, which is either Cacheable or Uncacheable. Only the Cacheable variant merges.

Three things people expect to be in that key and are not:

No salt exists. A comment in the source refers to one; the struct has no such field. If you need two textually identical actions to be distinct, the difference has to be in the action itself.

do_not_cache is never read. The REAPI flag exists on the message and NativeLink parses it, but no code in the repository branches on it. Setting it does not stop a result being written.

The flag that actually matters is skip_cache_lookup, and its effect is asymmetric in a way worth stating plainly: it disables the cache read, and it disables merging, but it does not disable the cache write. An action submitted with skip_cache_lookup runs on its own, ignoring any existing result, and then overwrites that result.

Why that asymmetry is the useful one

It is exactly what you want for "re-run this and fix the cache": force a real execution, and let everyone else benefit from the answer. If it also suppressed the write, the poisoned entry you were trying to replace would survive.

Ordering is one 64-bit integer

AwaitedActionSortKey packs the entire queue ordering into a single u64: the upper 32 bits are the priority (shifted from i32 into unsigned range so byte order matches numeric order), and the lower 32 bits are the insertion timestamp in whole seconds XORed with u32::MAX.

The XOR is what makes an earlier timestamp sort higher, so the rule falls out of the encoding: priority dominates absolutely, and ties break in favour of the action that arrived first, at one-second resolution; actions inserted within the same second are ordered by operation ID, which is effectively arbitrary. It has no aging, no starvation guard, and no fairness across clients. A client that submits everything at priority 1000 will starve a client at priority 0 indefinitely, and nothing in the scheduler will intervene.

Priority is a contract between the people sharing a cluster, enforced socially. If you need it enforced mechanically, that has to happen in front of the scheduler.

Two state enums, and one that never appears

ActionStage is the wire-facing state a client sees. SortedAwaitedActionState is the database index state the scheduler sorts on. They overlap but are not the same, and conflating them while reading the source is a reliable way to get lost.

Notably: ActionStage::CacheCheck is never produced by the scheduler crate. It exists in the protocol, but by the time an action reaches the scheduler the cache lookup has already happened in the decorator in front of it; see where an action's path forks. A client waiting to observe CacheCheck will wait forever.

Matching is first-fit over an LRU

The matching engine walks the queue in sort order, and for each action looks for the first worker that satisfies the action's platform properties. The source still carries an O(n*m) TODO in SimpleScheduler, but the walk is narrowed first: a WorkerCapabilityIndex (an inverted index over Exact and Priority properties) produces the candidate set, and only that set is scanned in LRU order, checking availability and Minimum values per worker.

WorkerCapabilityIndex

For clusters of tens or low hundreds of workers this is not the bottleneck. It is worth knowing about before you plan for thousands.

allocation_strategy is a smaller knob than its name suggests. The two original values choose which end of the LRU the scan starts from, and nothing else. Starting from the least-recently-used end spreads load; starting from the most-recently-used end packs work onto hot workers and leaves the rest idle long enough to be scaled down. least_loaded takes the eligible worker with the fewest running actions, least recently used among equals, so a worker that just joined takes the next action instead of waiting for every older worker to be used once more, and a burst that lands while most of the fleet is still starting does not fill the first worker to its concurrency cap. None of them is a bin-packing algorithm.

The admission ledger only subtracts what an action declares, so an action that declared too little is invisible to it. live_memory_veto names a minimum property, usually memory_kb, that the scheduler checks against the free memory each worker reports with its keepalive (the cgroup limit less current usage, or MemAvailable without a limit). A worker reporting less than the action asks for is skipped for that action; a worker that reports nothing, such as one older than v1.7.3, is never vetoed. Two metrics go with this: worker.keepalive.gap is the seconds between consecutive messages from a worker, which is what worker_timeout_s should be judged against, and worker.disconnections carries a worker.disconnect.reason of disconnected, drained, timeout, kill_unacknowledged, error, removed or shutdown.

Within a priority the queue is ordered by insert time to the nanosecond, on both backends, so a burst submitted together dispatches in arrival order.

The pass keeps that order when a worker frees up while it is running. A pass reads the queue one action at a time, a store round trip each on a store-backed queue (the listing serves the record it already loaded, and only the client keepalive is read), so on a long queue a worker that finishes mid-pass would otherwise hand its room to whatever action the pass happened to be reading, newer than everything it had already passed over. Instead, actions that found no room are parked in the order the queue listed them, highest priority first and then oldest, and when the fleet's capacity changes the pass offers the room to the parked actions in that order before the one in hand, asking the fleet once per property shape. scheduler.matching.parked_dispatched counts the parked actions a worker was actually told to run that way; a non-zero rate is normal on a busy fleet and means the head of the queue is being served.

When a worker reports an action killed for memory, by its own reservation enforcement or by the kernel near the limit, memory_escalation requeues the action with the reservation multiplied by percent (default 200) instead of failing it. With ladder_kb set to the memory values of the fleet's size classes, a kill steps the reservation to the next class up instead of scaling it. The next dispatch carries the raised value, so the worker holds the action to that. The reservation grows to the largest memory any connected worker advertises (or max_kb); the last step reserves that worker whole, its memory and, through cpu_property, its CPU, so the action runs alone. Only a kill there fails the action, with a message saying no worker in the fleet can run it. Escalations have their own budget (max_steps, default 8) and do not spend max_job_retries. The execution histograms carry an execution.outcome attribute so the killed samples can be told from the completed ones.

The historical_resource scheduler that fills reservations reads hints by a key ladder, each step probed only when the action carries that key: (target_id, action_mnemonic), then target_id, then the action's action_digest, then its command_digest (each hash-size or hash/size), then the mnemonic alone. The digest keys reach every client, whether or not it sends Bazel request metadata; the action digest is what origin events carry, so it is what a producer fed by them writes. A hint may name a class instead of numbers, and classes on the same scheduler defines the fleet's size ladder: named points with cpu_count, memory_kb and disk_kb, listed ascending. An action no hint matches takes the class picked by class_by_mnemonic, then class_by_timeout (Bazel test sizes arrive as timeouts), then default_class; without classes the raw cold_start numbers apply. With memory_escalation.ladder_kb set to the classes' memory values, a memory kill steps the reservation to the next class instead of scaling it, and a kill at the top class fails the action. Older hint files still load: every new field is optional.

Platform properties are typed, and Minimum is a resource pool

Every property key a worker advertises has a declared type:

TypeMeaning
ExactThe action's value must equal the worker's
MinimumThe worker's value must be at least the action's, and is consumed
PriorityThe worker must have the key; its value is not compared (the source marks value-based preference as a TODO)
IgnoreAccepted and disregarded

A worker announcing a property key the scheduler does not know is rejected at connect time. An action naming an unknown key is not rejected at submission: the key is only translated during the matching pass, where it produces an Unknown platform property error that is logged, and the action stays queued. That is the harder of the two to diagnose, so if an action sits in the queue with workers idle, check the scheduler log for that message before anything else.

make_prop_value

Minimum is the one that changes how you think about a worker. Its value is decremented when an action is dispatched and restored when the action completes; see Worker. A worker advertising cores: 16 is not saying "I am a 16-core machine, send me anything"; it is a pool of 16 units that actions draw down and return. That is the mechanism behind running several actions concurrently on one worker, and it is why getting the advertised number wrong oversubscribes the machine rather than queueing.

Platform properties covers this from the configuration side.

Retries, and what "failed" means

The attempt counter is incremented by one event, and two other events that used to spend it now have budgets of their own.

UpdateWithError increments the attempt count, unless the error is ResourceExhausted. Backpressure is free. A worker saying "I am too busy" does not spend one of the action's lives, which is what makes it safe for a worker to shed load.

UpdateWithDisconnect counts against max_worker_loss_retries (default 10), not max_job_retries. A worker that vanishes mid-action is usually not the action's doing, so the action goes back to the queue without spending one of its lives; the separate budget is what stops an action that reliably OOM-kills its worker from taking one down forever. Past it the client sees a FAILED_PRECONDITION saying the worker was lost that many times.

UpdateWithEscalation counts against memory_escalation.max_steps (default 8), not max_job_retries either: the scheduler chose to run the action again with more memory, and the fleet's ceiling bounds how far that goes.

The limit is checked as attempts > max_job_retries, strictly, so the default of 3 permits four total executions. That off-by-one is in the name: three retries after one initial attempt.

Two timeouts with very different defaults

max_action_executing_timeout_s defaults to 0, which means disabled. It is not the only guard against a hung action: the worker kills the process when the action's own timeout (or the worker's max_action_timeout_s, 20 minutes by default) elapses, and the scheduler re-queues an executing action whose worker stops sending updates for worker_timeout_s. What this knob adds is a ceiling for the case where the worker is alive and heartbeating but sends no update for a specific action; with it disabled, that case has no ceiling. On a shared cluster it is worth setting.

client_action_timeout_s defaults to 60. If no client is listening for that long, the action is killed with DeadlineExceeded. This is a deliberate anti-orphan measure (a client that pressed Ctrl-C should not leave work running), but it also means a client that disconnects and reconnects slowly can lose work it expected to still be running.

Who writes the action cache entry?

Not the scheduler. The worker writes the action cache entry, and it does so before telling the scheduler that the action finished. The ordering is what guarantees that a client which observes completion and then queries the cache finds the result there.

What gets written depends on the exit code, via should_cache_result in RunningActionsManager. With the defaults, a successful action's result goes to the action cache (SuccessOnly), and a failed action's result goes to the historical results store (FailuresOnly) instead. A non-zero exit code is not an error in the REAPI sense (it is a legitimate, cacheable outcome), but by default NativeLink keeps it out of the action cache so that a failing compile is re-attempted rather than replayed.

Where the exact truth lives

Every scheduler field and default is in the generated scheduler configuration reference. For the operational levers, see tuning.

Common questions

NextWorker execution

What the worker does with an action once the scheduler hands it over, including exactly how much isolation you are and are not getting.

On this page