Files
app/docs/code/workers.md
T
stroblmeandClaude Opus 5 40f8ad378d Ask a cluster for a machine when nothing here will do
Slurm is not a machine that attaches and stays; it is a queue somebody else
owns. So nothing here submits a node to it. It submits a job whose payload is an
ordinary worker dialling back in, and everything downstream — the protocol, the
artifacts, cancellation, the books — already worked and did not have to learn
what Slurm is.

The alternative, which Covalent takes, is to stage a serialized call and a
runner onto the login node, poll squeue and copy the result back: a second way
of running a node beside the one that exists. The cost of not doing that is one
assumption, that a compute node can open a connection outward. Where that is
false, _payload is the single method a staged variant would replace.

Clusters are configured in provisioners.json beside the alerts, since this is
infrastructure an operator writes rather than anything a flow says. The script
is generated with the system ssh and no new dependency, and prerun owns the
environment — deliberately no pip install, because what is on a cluster is
somebody's decision.

One outstanding request per profile, cancelled if it never attaches and on the
way out. Nothing autoscales.

The run gate needed the same hook: a run held before it starts never reaches the
placer's own wait, so it would have queued forever on a machine nothing had
asked for.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01A6HeySA27EkGANZN95QySW
2026-08-27 09:09:44 +02:00

229 lines
9.6 KiB
Markdown

# Remote workers
The engine runs where the automations are. The GPU is somewhere else, the
Raspberry Pi with the relays is in a shed, and neither of them is on the same
network as the other.
A **worker** is a process that runs the code of nodes marked for it. It dials
*out* to the engine over one authenticated websocket, so nothing on that
machine has to be reachable — and nothing has to expose the engine's state
backend across hosts, which it never should.
## Install and attach
```sh
pip install fluksio-worker
fluksio-worker \
--url wss://api.example.com/api/v1/workers/attach \
--token "$FLUKSIO_WORKER_TOKEN" \
--labels gpu,cuda12 \
--python /opt/torch-venv/bin/python
```
`fluksio-worker` is its own distribution — the agent, the node runner, and
`websockets`. Nothing of the engine, so a GPU box does not install a database
driver in order to run a training step. An engine host already has it, and
`fluksio worker …` is the same program.
| Option | Default | What it does |
|---|---|---|
| `--url` | **required** | `wss://…/api/v1/workers/attach` |
| `--token` | `$FLUKSIO_WORKER_TOKEN` | the credential, minted on the engine |
| `--name` | this host's name | how it shows up in the worker list |
| `--labels` | none | comma-separated; what a node's `device` matches |
| `--python` | this interpreter | the interpreter node code runs on |
| `--parallel` | `1` | how many node calls it will take at once |
| `--artifact-url` | derived from `--url` | where the artifact store is, if not beside the socket |
| `--cpus` | what the job or the machine has | cores to advertise |
| `--gpus` | what the job says, else none | GPUs to advertise; never probed |
| `--ram-mb` | what the job or the machine has | memory to advertise, in MB |
| `--max-idle` | never | stop after this many seconds with nothing running |
`--python` is the important one. It is how this machine keeps its own wheels —
the CUDA build, the vendor SDK, the thing that will not install anywhere else —
without the engine ever installing them or knowing about them.
### What it says it has
A worker reports its inventory when it attaches — cores, GPUs and memory — and
the engine schedules against it: a node asking for two cores and a GPU goes to
a machine that has them free, not merely to one carrying the right label.
Cores and memory are read off the machine, or off the batch job that started
this worker (`SLURM_CPUS_ON_NODE`, `SLURM_MEM_PER_NODE`). **GPUs are never
probed.** Asking a vendor tool would make the one dependency two, so a GPU is
something the job says it was given (`SLURM_GPUS_ON_NODE`, `SLURM_JOB_GPUS`,
or `FLUKSIO_WORKER_GPUS`) or something you say with `--gpus`. A worker that
reports nothing still attaches and is scheduled by its label alone, as every
worker was before any of them reported anything.
The engine tells each call what it may use — thread caps, and the devices it
may see. The worker starts a process per call, so it applies them at the only
moment a numerical library still reads them: before the import.
`--max-idle` is for a worker something else started for one job — a batch
scheduler, say. It exits when nothing has run for that long, so the allocation
goes back rather than idling until its walltime.
## Mint the token
On the engine, as a superuser:
```sh
curl -X POST $FLUKSIO/workers/tokens -H "Authorization: Bearer $TOKEN" \
-H 'Content-Type: application/json' -d '{"name": "gpu-dev"}'
```
Shown once, valid for a year — a worker is a machine somebody sets up and
leaves running. It is signed with the same keypair agent tokens use, so
rotating that key revokes every worker along with them.
??? note "A host where pip is not an option"
The two files work copied into one directory and run with `python agent.py
…`. The engine serves the runner itself at `GET /api/v1/workers/runtime`
it is the same module its own local workers run, deliberately standard
library only.
## Send a node to it
A node declares the label of the machine it needs:
```json
{
"id": "train",
"device": "gpu",
"device_policy": "require",
"timeout": 7200
}
```
| `device_policy` | Behaviour when nothing carrying the label is attached |
|---|---|
| `require` (default) | the run stays `queued` and says what it is waiting for |
| `prefer` | it runs on the engine instead |
`prefer` is what makes a flow work before the GPU box exists. `require` is what
you want once it does.
!!! note "Set from the API"
`device` and `device_policy` are not yet fields in the node panel. Set them
with `PUT /flows/{name}`.
## What follows from this
- **The node's source travels with every call.** Nothing has to be deployed to
the worker, and changing a node's code takes effect on the next execution.
- **A node bound to a device is compiled on that machine.** A node importing
`torch` is correct on the GPU box and a missing module on the engine, so
checking it here would fail something that is fine.
- **`import fluksio` inside a node is the worker's own reporter.** `emit`,
`save_artifact`, `load_artifact` — installed before your code runs, so an
installed `fluksio` package on that box never shadows it.
- **Cancelling a run kills what it is executing**, there or here, and leaves
other runs of the same node alone.
- **If the worker disappears mid-call**, the run fails in seconds with
`worker went away mid-call` rather than waiting out its timeout.
- **A worker sends a heartbeat every ten seconds while it executes**, so a long
node is distinguishable from a dead socket. Ninety seconds of nothing at all —
not even a heartbeat — fails the call as gone. A heartbeat says the *agent* is
alive and nothing about the node, so it never satisfies a node's own timeout:
one set to thirty seconds fires after thirty seconds of the node reporting
nothing, wherever it runs.
## Artifacts across machines
An artifact reference names content by its hash, not a location, so it stays
valid wherever the store is reachable from. A worker that shares the engine's
filesystem writes to it directly; one that does not fetches and uploads over
HTTP, using the artifact endpoint beside the socket it already has. Either way
your node code is the same two calls.
## Seeing what is attached
```sh
curl -s $FLUKSIO/workers -H "Authorization: Bearer $TOKEN" | jq
```
Name, labels, how many calls it will take at once, how many are in flight, when
it attached, when it was last seen, its Python version, a digest of its
environment, and what it says it has: cores, GPUs and memory.
## Upgrading
The engine and the worker speak a version-matched protocol, and a worker
announcing anything else is refused rather than half-understood. Protocol 2 —
the one that carries inventory — is `fluksio-worker` 0.2.0. An older agent is
told so on the socket and stops, rather than retrying against an engine that
will never accept it; `pip install -U fluksio-worker` on that host is the whole
upgrade. Nothing changed in the runner served at `GET /api/v1/workers/runtime`,
so a host that copies its two files copies the same one as before.
## Machines from a batch scheduler
A cluster is not a machine that attaches and stays — it is a queue somebody else
owns. So Fluksio does not submit *nodes* to Slurm. It submits a job whose payload
is an ordinary worker dialling back in, and from there everything works the way
it already does: the same protocol, the same artifacts, the same cancellation.
Write the clusters into `provisioners.json` beside the flows:
```json
[{
"type": "slurm",
"name": "hpc",
"login": "me@login.cluster",
"ssh_key": "/secrets/hpc_ed25519",
"engine_url": "wss://api.example.com/api/v1/workers/attach",
"max_idle_s": 300,
"provision_timeout_s": 900,
"profiles": [{
"name": "gpu-small",
"cpus": 8, "gpus": 1, "ram_mb": 65536,
"labels": ["gpu"],
"sbatch": ["--partition=gpu", "--gres=gpu:1", "--time=04:00:00"],
"prerun": ["module load cuda/12", "source ~/venvs/flux/bin/activate"]
}]
}]
```
A **profile** is what the scheduler is asked for, where a flavor is what a node
asks for. They are separate on purpose, and agree when you set them up to.
When a node needs a machine nothing attached can give, and a profile would fit,
the engine `sbatch`es one over ssh — the system `ssh`, so nothing new is
installed — and the run waits meanwhile, saying so. `prerun` owns the
environment: a `module load`, a venv with `fluksio-worker` already in it. There
is deliberately no `pip install` in the generated script, because what is
installed on a cluster is somebody's decision and not this program's.
One outstanding request per profile, however often it is asked for. A job that
never attaches within `provision_timeout_s` is `scancel`led, as is anything
outstanding when the engine stops. `--max-idle` is what ends the job at the
other end, so an allocation goes back rather than idling to its walltime.
`GET /workers/resources` reports what is outstanding and what last went wrong.
!!! note "It needs a route out"
A compute node must be able to open a connection to the engine. That is
true of most clusters and false of air-gapped ones; there is no staging
path today.
## What a worker is not
It is not a second engine. Subscriptions, schedules, webhooks, the dashboards
and the run queue all stay in one process — that is what keeps a message having
one definition and a cron tick happening once. A worker executes node bodies.
Running two engines against one data directory is not supported. Distribute
work with workers.
## See also
- [Runs: pipelines that finish](../concepts/runs.md#running-a-node-somewhere-else)
- [Writing node code](nodes.md)
- [Getting started: data science](../getting-started/data-science.md)