Schedulers and workers
The minimum config for each, the properties contract between them, and the fields that decide what happens when something goes wrong.
Who this is for: anyone adding execution to a config that currently only caches. What you'll know at the end: the minimum scheduler and worker blocks, how the two are wired to each other, and which of the many optional fields actually change behaviour you'll notice. Time: thirty minutes.
Before you start
A config whose stores and servers you understand. See Stores and Servers and services.
The minimum scheduler
A scheduler is one entry in the schedulers array, with a name and a type:
schedulers: [
{
name: "MAIN_SCHEDULER",
simple: {
supported_platform_properties: {
cpu_count: "minimum",
OSFamily: "exact",
"container-image": "priority",
},
},
},
],simple is the scheduler you want. The other four types (grpc,
cache_lookup, property_modifier, and historical_resource) either
forward to another scheduler or wrap one to change its behaviour; none of them
is a starting point.
supported_platform_properties is the whole contract. It declares which
property keys this scheduler understands and how each is matched:
| Value | Meaning |
|---|---|
minimum | Parsed as a number; a worker must advertise at least this much |
exact | Parsed as a string; must match the worker's value exactly |
priority | Not used for filtering; passed to the worker as information |
ignore | Actions may request the key; workers need not advertise it |
A property key an action requests that is not in this map is an error. A key that is in the map but that no worker advertises means the action never matches anything and queues silently. That failure mode gets its own page, Platform properties, because it's the first wall almost everyone hits.
SimpleSpecThe minimum worker
A worker is one entry in the workers array. local is currently the only
type; it means the worker runs actions in this process, on this machine:
workers: [
{
local: {
name: "WORKER_1",
worker_api_endpoint: {
uri: "grpc://${SCHEDULER_ENDPOINT:-127.0.0.1}:50061",
},
cas_fast_slow_store: "WORKER_FAST_SLOW_STORE",
upload_action_result: {
ac_store: "AC_MAIN_STORE",
},
work_directory: "/tmp/nativelink/work",
platform_properties: {
cpu_count: { values: ["16"] },
OSFamily: { values: ["linux"] },
"container-image": { values: [""] },
},
},
},
],Four of those fields are load-bearing.
worker_api_endpoint is the address of the scheduler's worker_api
service, the private port from
Servers and services, not the public
one. The worker dials out; the scheduler never dials in. That direction is why
a worker fleet can sit behind NAT with no inbound rules.
cas_fast_slow_store must name a fast_slow store whose fast half is a
filesystem store. The worker builds each action's input tree by hard-linking
files out of that fast store into the work directory, which requires real
files on a real filesystem. The slow half must eventually resolve to the
same CAS the client and scheduler use, or be a noop, when the fast tier is
the only storage.
work_directory must be on the same filesystem as the fast store's
content_path, for the same hard-linking reason. It is fully managed by the
worker and purged on startup.
platform_properties is the worker's half of the contract. Each key is
either a static list of values or a query_cmd that is run (as a command, not
through a shell) each time the worker connects to the scheduler, with its
output split on newlines:
platform_properties: {
cpu_count: { query_cmd: "nproc" },
OSFamily: { values: ["linux"] },
},Every key the scheduler declares as minimum or exact must appear here, or
this worker will never be matched for any action requesting it.
`upload_action_result` is one level, not two
The correct shape is upload_action_result: { ac_store: "AC_MAIN_STORE" }.
Configs showing upload_action_result: { upload_action_result: { ... } }
are wrong and will be rejected; ac_store is a direct field. Omitting the
block entirely is also valid, and means the worker runs actions but never
caches their results, which looks exactly like a broken cache.
How they find each other
Nothing in the scheduler names a worker. Workers are not declared to the scheduler at all; they connect to it, announce their platform properties, and are added to the pool. Removing a worker means stopping it.
That is what the three pieces of wiring add up to:
- the
executionservice names the scheduler, so clients can submit actions; - the
worker_apiservice names the same scheduler, on a different port, so workers can join; - each worker names the
worker_apiaddress, so it knows where to join.
In a single-process config all three are the same process and the endpoint is
127.0.0.1. In a fleet, the scheduler is its own deployment and that endpoint
is a service address. The config shape does not change.
The fields that change what you'll notice
SimpleSpec has ten optional fields. Most can stay at their defaults
forever. These four are the ones worth setting deliberately:
worker_timeout_s (default 10 s): how long a silent worker stays in the
pool. Any message from the worker counts, not only the heartbeat. Keep it at
least twice the worker's worker_api_endpoint.timeout; a worker evicted
mid-action has every action it held re-queued, each with an attempt spent.
max_action_executing_timeout_s (default 0, disabled): caps how long an
action can sit in Executing with no update, even from a worker that is
otherwise healthy. This is the setting that catches a worker stuck on one
specific action rather than dead. Without it, such an action hangs until a
client gives up.
max_job_retries (default 3): how many times an action that fails with
an internal error is retried before the last error is returned to the client.
It exists to stop one poisonous action from cycling through and destabilising
your whole fleet. A lost worker and a memory escalation do not spend it.
max_worker_loss_retries (default 10): how many times an action whose
worker was lost (disconnected, timed out, evicted, or OOM-killed as a whole)
is queued again before the client gets a FAILED_PRECONDITION. Losing a
worker is not the action's failure, so this budget is separate and larger;
it only stops an action that takes a worker down every time it runs.
retain_completed_for_s (default 60 s): how long a completed action's
result stays available for a late WaitExecution call. Clients that
disconnect and reconnect need this window to be longer than their reconnect
time.
On the worker side, max_inflight_tasks (default 0, meaning unlimited) is
the one to set on any real machine. Unlimited means the worker accepts every
action the scheduler offers and lets the OS sort out the contention.
drain_on_shutdown (off by default) makes a worker announce that it is
draining the moment it gets SIGTERM, so the scheduler stops dispatching to
it while it finishes what it holds. It needs a scheduler that understands the
drain flag; an older one removes the worker on the spot and requeues its
actions.
Three bounds keep one action from taking the worker with it.
max_captured_output_bytes (default 0, unbounded) caps how much of an
action's stdout or stderr the worker holds in memory; past it the output
spills to a file under the action directory and is uploaded from disk, so a
tool that prints gigabytes costs disk, not RAM. precondition_timeout_ms
(default 30000) ends a hanging precondition script and refuses the action as
backpressure. orphan_sweep_interval_s (default 0, off) removes action
directories under work_directory that no running action owns and that a
failed cleanup left behind; the directory is always purged at startup.
When the worker ends an action, for a timeout or a scheduler cancel, the
process group gets SIGTERM first and kill_grace_ms (default 5000) to
write its own cleanup before SIGKILL follows; the memory reservation kill
skips the grace. A killed action keeps the stdout and stderr it produced.
set_tmpdir (default true) gives every action its own TMPDIR under its
action directory, so concurrent actions writing the same temporary file name
no longer collide; a TMPDIR in the action's own environment still wins, and
code that hard-codes /tmp is unaffected. A timeout above max_action_timeout_s
is clamped to it rather than rejected.
When the action's own process has exited, anything it left running in its
process group (a background job, a helper that daemonized) is killed, since
it would otherwise hold the action's output pipes open and the action would
sit there until its timeout; the pipes then get five seconds to drain, and
what was read by then is kept. The worker also makes itself the parent of
every process an action leaves behind and reaps them on a timer, so a build
whose actions fork and do not wait no longer leaves hundreds of zombies on
the worker for the life of the pod. Under use_namespaces the stub does
the reaping inside the action's PID namespace.
Every executed action reports its resource usage to the scheduler with an
outcome: COMPLETED, KILLED_MEMORY (the worker's reservation kill, or a
SIGKILL from the kernel with the last memory sample within 10% of the limit),
KILLED_TIMEOUT or KILLED_EXTERNAL, plus wall time and the reservation the
action ran under. The scheduler logs every kill outcome with its mnemonic.
resource_enforcement (off by default, Linux only) makes the worker hold an
action to the memory reservation the scheduler placed it by. With
memory: soft the worker kills the action's process group once two
consecutive samples exceed memory_kb plus memory_headroom_percent
(default 20), or one sample exceeds twice that, and fails it with
FAILED_PRECONDITION naming the reservation and the observed peak. Samples
are 250 ms apart, so an action that allocates faster than the pod's whole
limit in half a second can still reach the cgroup limit first; the
reservations admitted onto a worker must also sum to less than its limit
with the headroom included. Without it the kill comes from the pod's cgroup limit,
and on cgroup v2 that kill takes the worker and every action on it, not just
the one that grew. The sample is the group's proportional set size, so a
parent's pages shared with its forked children count once. Actions run
through a persistent worker are not covered: that process outlives the action
and serves others, so only the pod's limit bounds it. With enforcement on, the
worker reads an action's properties from the scheduler's dispatch rather than
the client's request, so a reservation the scheduler placed the action by is
the one enforced and exported to the environment.
disk: guard adds a check before any input is fetched: an action whose
disk_kb reservation exceeds the free space under the work directory, less
what the actions already admitted reserved, is refused with RESOURCE_EXHAUSTED, which the scheduler requeues without
counting an attempt and holds off this worker until its next keepalive,
instead of running the fetch into a full disk.
disk: soft is the guard plus a kill: every ten seconds the worker adds up
the files the action has written under its own directory (those modified
since its command started; inputs were materialized before that and are
not counted), and an action past its disk_kb reservation by
disk_headroom_percent is killed and fails with FAILED_PRECONDITION, its
outcome KILLED_DISK. The walk runs on the blocking pool as its own task,
so memory sampling carries on meanwhile. With either disk mode the action's
usage report carries peak_disk_kb, what it wrote, from one walk of the
action directory at completion, which is the only walk guard pays. The
figure is reported and exported for now; the scheduler's sizing consumer
comes with the size ladder.
"resource_enforcement": {
"memory": "soft",
"memory_property_name": "memory_kb",
"memory_headroom_percent": 20,
"disk": "soft",
"disk_property_name": "disk_kb",
"disk_headroom_percent": 20
}capacity makes the worker advertise CPU and memory from its own cgroup
limits instead of from platform_properties: it reads cpu.max and
memory.max at its cgroup v2 root, keeps back overhead_cpu_millicores
(thousandths of a core) and overhead_memory_kb for itself, divides the
memory by the enforcement headroom (memory_headroom_percent, which follows
resource_enforcement unless set), and sets the two properties to the
result at registration. The limit is the number the kernel enforces, so the
number the scheduler packs against is derived from it rather than typed next
to it. Where the cgroup cannot be read, hand-typed cpu_count and
memory_kb stand with a warning; with neither typed the worker refuses to
start, since it could take no action that asks for either. cpu_unit is the scale of the CPU
property: cores by default, rounded down, the scale cpu_count: { query_cmd: "nproc" } and an action's cpu_count=1 use; a fleet that
advertises and requests thousandths of a core sets millicores. The
overhead is always in thousandths of a core.
"capacity": {
"cpu_unit": "cores",
"overhead_cpu_millicores": 1000,
"overhead_memory_kb": 4194304
}Scheduler state is in memory by default
experimental_backend defaults to memory, so restarting the scheduler loses
the queue. A Redis backend exists for shared state across scheduler
replicas. Both are covered in
Production configuration; for a single
scheduler, the default is correct.
Putting the pieces together
Adding execution to a cache-only config takes four edits, all of which you've now seen:
Add a
schedulersarray with onesimpleentry, naming the platform properties you intend to match on.Add
executionto the public server's services, naming the CAS store and the scheduler.Add a second server on a private port with
worker_apinaming the same scheduler, and nothing else client-facing on that port.Add a
workersarray with onelocalentry pointing at that private port, plus thefast_slowstore it needs.
You did it right if
- The worker logs
Worker registered with schedulerwith its assigned worker id. - A build submitted with matching platform properties reports actions executing remotely rather than queuing.
- Every key the scheduler declares
minimumorexactappears in the worker'splatform_properties.
Your first full config does exactly this, from an empty file, one block at a time.
FAQ
Empty file to a running cache-and-execution cluster, one block at a time, with a check after each.
SidewaysPlatform propertiesThe matching contract in depth, and how to diagnose a queue that never drains.