NativeLink

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:

ValueMeaning
minimumParsed as a number; a worker must advertise at least this much
exactParsed as a string; must match the worker's value exactly
priorityNot used for filtering; passed to the worker as information
ignoreActions 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.

SimpleSpec

The 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.

LocalWorkerConfig

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 execution service names the scheduler, so clients can submit actions;
  • the worker_api service names the same scheduler, on a different port, so workers can join;
  • each worker names the worker_api address, 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:

  1. Add a schedulers array with one simple entry, naming the platform properties you intend to match on.

  2. Add execution to the public server's services, naming the CAS store and the scheduler.

  3. Add a second server on a private port with worker_api naming the same scheduler, and nothing else client-facing on that port.

  4. Add a workers array with one local entry pointing at that private port, plus the fast_slow store it needs.

You did it right if

  • The worker logs Worker registered with scheduler with its assigned worker id.
  • A build submitted with matching platform properties reports actions executing remotely rather than queuing.
  • Every key the scheduler declares minimum or exact appears in the worker's platform_properties.

Your first full config does exactly this, from an empty file, one block at a time.

FAQ

NextYour first full config

Empty file to a running cache-and-execution cluster, one block at a time, with a check after each.

SidewaysPlatform properties

The matching contract in depth, and how to diagnose a queue that never drains.

On this page