Cluster Mode
September 2, 2026 · View on GitHub
Audience: SREs running Bernstein on more than one host.
Overview
Cluster mode lets one Bernstein server coordinate work across many remote worker nodes. Use it when one host can no longer absorb the agent workload, when you need fault tolerance across machines, or when teams want a single multi-tenant control plane in front of dedicated worker pools (GPU boxes, region-pinned hosts).
Do not reach for cluster mode for single-host workloads. The single-process orchestrator is simpler, faster to debug, and avoids the network failure modes covered below. The architecture is intentionally a coordinator-with-workers topology, not peer-to-peer; if you only have one machine the extra moving parts are pure overhead.
Architecture
There is exactly one central server (bernstein run --remote - the run
command is also exposed under the hidden alias conduct; both dispatch the
same function) and zero-or-more worker nodes (bernstein worker).
Workers never talk to each other. They register with the central server, send
heartbeats, claim tasks, and report completion. All scheduling decisions
happen on the central server.
Quick start:
# 1) start a central server reachable from workers
bernstein run --remote --goal "Build feature X"
# 2) start one or more workers
bernstein worker --server http://central-host:8052
+--------------------+
| Central server |
| (FastAPI+task DB) |
| |
| NodeRegistry |
| TaskStealPolicy |
| ClusterAuthen- |
| ticator |
+--------+-----------+
|
+--------------+----------------+
| | |
+------v-----+ +-----v------+ +------v-----+
| bernstein | | bernstein | | bernstein |
| worker A | | worker B | | worker C |
| (GPU box) | | (us-east) | | (laptop) |
+------------+ +------------+ +------------+
Code references:
src/bernstein/core/protocols/cluster/cluster.py:37-NodeRegistry, the in-memory (optionally disk-persisted) registry of nodes.src/bernstein/core/protocols/cluster/cluster.py:340-TaskStealPolicymatches over- and under-loaded nodes.src/bernstein/core/protocols/cluster/cluster.py:396-NodeHeartbeatClientlibrary used bybernstein workerto call back.src/bernstein/core/routes/task_cluster.py:25- every cluster HTTP endpoint listed below.src/bernstein/core/protocols/cluster/cluster_auth.py:49-ClusterAuthenticator(JWT issuance, verification, revocation).src/bernstein/cli/commands/worker_cmd.py:50-WorkerLoop, the worker-side run loop.src/bernstein/core/fleet/- fleet aggregator that can roll up the same cluster status across multiple Bernstein projects (seecore/fleet/aggregator.py).
Worker setup
A worker process is a thin loop that registers, heartbeats, claims tasks for a
configured set of roles, and spawns local CLI agents to execute them
(worker_cmd.py:355). Start one with:
bernstein worker --server http://central:8052 --token "$BERNSTEIN_AUTH_TOKEN"
Common flags (worker_cmd.py:371-430):
| Flag | Default | Purpose |
|---|---|---|
--server URL | required | Central server URL. Also reads BERNSTEIN_SERVER_URL. |
--token TOKEN | from env | Bearer token for cluster auth. Also reads BERNSTEIN_AUTH_TOKEN. |
--name NAME | hostname | Worker name shown in /cluster/nodes. |
--slots N | 6 | Maximum concurrent agents on this worker (= node capacity). |
--roles a,b,c | backend,qa,security,frontend | Roles this worker accepts. |
--label k=v | (repeat) | Free-form labels for affinity routing. |
--adapter NAME | auto-detect | Which CLI agent to invoke (claude, codex, ...). |
--model NAME | adapter default | Default model for tasks with no explicit model. Non-Claude adapters fall back to their own discovered default. |
--poll-interval SECONDS | 10 | Task poll cadence. |
--poll-interval-ms MS | none | Override poll cadence in milliseconds. |
--heartbeat-interval-ms | 15000 | Heartbeat cadence. |
Capacity is set with --slots. Internally the worker tracks available_slots = max_agents - len(active_tasks) and reports it on every heartbeat
(worker_cmd.py:104-108). Pick a value that reflects how many parallel
adapter processes the host can comfortably run. There is no auto-detection -
oversubscribing a worker just queues longer.
Worker environment requirements:
- A logged-in CLI agent (
claude,codex,gemini,qwen,aider). Auto-detection usesagent_discovery.discover_agents_cached()(worker_cmd.py:38-46). - Python 3.12+ runtime with the same Bernstein version as the central server.
- Network reachability to the central server's HTTP port.
- A git checkout of the target repository at the worker's workspace
(
--workdir, defaultcwd;/workspacein the published cluster image). Every claimed task runs in a git worktree created under this path, so the workspace must be a git work tree with at least one commit. Mount or clone the target repo there before starting the worker. A worker whose workspace is not a usable git repository refuses to start (see Workspace preflight below) rather than registering and claiming tasks it cannot run.
Docker/Kubernetes: the published image does not ship a git checkout at
/workspace. Bind-mount the repo (-v /path/to/repo:/workspace) or run an init step that clones it, so/workspaceis a git checkout with a commit beforebernstein workerstarts.
Workspace preflight
Before it registers or accepts any claim, the worker verifies its workspace is
a usable git repository (WorkerLoop._workspace_setup_error). If the workspace
is missing, is not a git work tree, or is a git repo with no commits, the
worker prints an actionable setup error naming the path and exits non-zero
without registering. This prevents the failure mode where a worker registers,
claims a task, and then cannot spawn the agent (git worktree add -> fatal: not a git repository), leaving the task stranded in claimed with no live
agent (#3018).
Readiness signal
A worker is ready when it prints Registered as node <id> — the central
server has accepted its registration and it is eligible for claims. Until
that line appears the worker is not yet part of the cluster; it retries
registration on its poll interval, printing the connection error each time,
rather than exiting. Automation should wait on that line rather than on a
fixed delay, and tests/integration/test_first_run_long_running_surfaces.py
holds the surface to it.
JWT node authentication
Cluster auth is JWT-based and is the recommended deployment mode in production. The flow is:
- Server config - pass a
ClusterAuthConfig(secret=..., require_auth=True)intoClusterAuthenticatorand mount it onapp.state.cluster_authenticator(cluster_auth.py:32-46,task_cluster.py:54). Whenrequire_auth=False, the verifier returns a synthetic anonymous payload (cluster_auth.py:113-121) - useful for tests only. - Token issuance - the server issues a node token via
ClusterAuthenticator.issue_node_token(node_id)(cluster_auth.py:70-93). By default the token carriesnode:registerandnode:heartbeatscopes; admin tokens additionally requirenode:admin. - Worker presents token - the worker sends
Authorization: Bearer <token>on every cluster HTTP call (worker_cmd.py:98-102). - Per-request verification -
task_cluster._verify_cluster_auth()resolves the required scope per route (SCOPE_NODE_REGISTER,SCOPE_NODE_HEARTBEAT,SCOPE_NODE_ADMIN) and rejects with HTTP 401 on any failure (task_cluster.py:44-60). - Heartbeat-bound identity - the heartbeat verifier checks that the
token's
user_idmatches the path'snode_id(cluster_auth.py:262-265); a stolen heartbeat token cannot be replayed against a different node. - Revocation -
revoke_token(token)andrevoke_node(node_id)mark a token unusable for subsequent verifications (cluster_auth.py:174-191). Token revocation is in-memory; for persistence across restarts use short-lived tokens (default 24h,ClusterAuthConfig.token_expiry_hours).
Default scopes, defined in cluster_auth.py:23-25:
| Scope | Required by |
|---|---|
node:register | POST /cluster/nodes, ClusterService.RegisterNode |
node:heartbeat | POST /cluster/nodes/{id}/heartbeat, POST /cluster/claims/gossip, ClusterService.Heartbeat, ClusterService.StreamHeartbeats |
node:admin | cordon, uncordon, drain, DELETE /cluster/nodes, POST /cluster/steal, and their ClusterService counterparts |
POST /cluster/steal reassigns other nodes' claimed work from
caller-reported queue depths, so it is scoped with the node-registry
mutations rather than with gossip: gossip verifies each receipt's own
Ed25519 signature and chain link inside the handler, so its bearer scope
only has to establish fleet membership. A default node token carries
register + heartbeat, so triggering a rebalance uses the cluster secret or
an admin-scoped token — the same credential the drain and cordon primitives
already need.
The worker side of the flow lives in worker_cmd.py:134-164 (registration)
and :165-190 (heartbeat). On HTTP 404 from the heartbeat the worker
re-registers automatically (eviction recovery).
The gRPC cluster surface
ClusterService (proto/bernstein/v1/cluster.proto) mirrors the routes above
for node-to-node traffic. It is not started by anything in src/ today; the
notes here apply once an operator wires BernsteinGrpcServer.start() into a
deployment.
- Same scopes, gRPC status codes. Pass the
ClusterAuthenticatorasstart(..., cluster_authenticator=...)and every mutating RPC verifies it. A missing or invalid credential is refusedUNAUTHENTICATED; a valid credential without the required scope is refusedPERMISSION_DENIED. Reads (ListNodes,GetClusterStatus) are unauthenticated, matchingGET /cluster/nodesandGET /cluster/status. - Where the credential goes. In the
authorizationcall metadata asBearer <token>-- the gRPC counterpart of the HTTP header. SetGrpcClientConfig.auth_tokento the cluster secret or an operator-minted node JWT. RegisterNodeResponse.auth_token. Populated with a node JWT minted against the registered node id, carryingnode:registerandnode:heartbeat(notnode:admin).ClusterClientadopts it for every subsequent call, so heartbeats travel under the node's own token rather than the join credential. With no authenticator wired the field stays empty -- nothing would verify a token minted there.- Re-registration is idempotent on node identity.
RegisterNoderesolves an existing entry by(name, url)and updates it. A worker that restarts therefore refreshes its row instead of adding one, andGET /cluster/statuskeeps counting its capacity once. A blanknameorurlnever matches, so anonymous registrations stay distinct. - Insecure port. Without
tls_cert_path/tls_key_paththe server binds an insecure port. Credential enforcement is unaffected -- it lives in the servicer, not the transport -- but node tokens then cross the wire in cleartext, and the server logs a warning saying so. Terminate TLS in front of the port or configure the cert pair.
Operational primitives: drain / cordon / uncordon / steal
Every primitive is one HTTP POST. The server-side handlers all live in
task_cluster.py and ultimately call into NodeRegistry
(cluster.py:155-185).
| Operation | Endpoint | What it does |
|---|---|---|
| Cordon a node | POST /cluster/nodes/{id}/cordon | Sets status to CORDONED. Node still heartbeats but is excluded from scheduling (cluster.py:155-163). |
| Uncordon a node | POST /cluster/nodes/{id}/uncordon | Restores ONLINE status from CORDONED or DRAINING (cluster.py:165-174). |
| Drain a node | POST /cluster/nodes/{id}/drain | Sets status to DRAINING. Equivalent to cordon + signal - agents finish their current tasks but no new work is assigned (cluster.py:176-184). |
| Trigger task stealing | POST /cluster/steal | Server runs TaskStealPolicy.find_steal_pairs() and resets eligible tasks back to open so quieter nodes can claim them (task_cluster.py:197-246, cluster.py:340-393). |
| Unregister (graceful exit) | DELETE /cluster/nodes/{id} | Removes the node from the registry. Workers call this on SIGINT / SIGTERM (worker_cmd.py:310-322). |
Task stealing thresholds (cluster.py:351-358):
overload_threshold=5- a node with more than 5 queued tasks is a donor.idle_threshold=2- a node with at least 2 free slots is a receiver.max_steal_per_tick=3- at most 3 tasks move per call.
The body of POST /cluster/steal is {"queue_depths": {"node-a": 7, "node-b": 0, ...}}; the response lists (donor_node_id, receiver_node_id, task_ids) actions and a total count (task_cluster.py:213-246).
Operationally, drain is the right primitive for a planned restart: cordon →
drain → wait for active_agents = 0 → unregister. Cordon is the right
primitive for an unhealthy host you want to investigate without rebooting.
Failure modes
Bernstein's cluster recovery is conservative - there is no global lock service, no leader election, no consensus protocol. Behaviours below are exact descriptions of what the code does today.
A worker disappears (network partition, crash, host reboot). The central
server keeps the node entry. Each tick the orchestrator (or any caller of
NodeRegistry.mark_stale(), cluster.py:197-210) checks
last_heartbeat; if it exceeds node_timeout_s the node is flipped to
OFFLINE. The node's claimed tasks are not automatically reassigned -
they remain in claimed status. Operators recover them with
POST /cluster/steal (donor=offline-node) or by deleting and re-claiming
the tasks. When the worker returns, its first heartbeat re-registers it
(cluster.py:185-195).
Partial network partition (worker can heartbeat but not claim tasks).
The worker's _claim_task() tolerates HTTP errors silently
(worker_cmd.py:202-204); the heartbeat path is independent. The node
appears healthy in /cluster/status but does no useful work. Detection:
watch for nodes with available_slots == max_agents despite
active_agents == 0 and a non-zero pending task count.
Server restart with active workers. NodeRegistry optionally persists
to a JSON file (persist_path in cluster.py:50-81). On startup all
loaded nodes are marked OFFLINE until they heartbeat
(cluster.py:72-76); fresh tokens are still valid (JWT secret survives the
restart). Workers detect eviction via HTTP 404 on heartbeat and re-register
via _register_with_retry() (worker_cmd.py:260-270).
Worker holds a stale token after revocation. The next mutating call
returns 401. The worker does not automatically renew - issue a fresh
token externally and pass it via BERNSTEIN_AUTH_TOKEN. Heartbeats are
separately scoped, so a heartbeat-only token cannot be promoted to
node:admin operations.
Worker's CLI adapter is missing or logged out. Auto-detection falls
back to "claude" (worker_cmd.py:46-47). The _spawn_agent() call will
raise inside AgentSpawner.spawn_for_task() and the worker logs a warning
without crashing. Because the spawn failed after the task was claimed, the
worker releases the claim back to the pool via POST /tasks/{id}/release
(WorkerLoop._release_task), so the task returns to open and another node
can pick it up instead of being stranded in claimed (#3018).
Observability for cluster health
Topology at a glance: bernstein cluster status / nodes
Two CLI subcommands render the registry without hand-rolling curl:
$ bernstein cluster status
Cluster topology=star nodes=2/3 online
capacity: 4 active / 8 free / 18 total slots
Cluster Nodes
Node ID Name Status Adapter Heartbeat Claimed Slots
node-alpha alpha online codex 4s 2 4/6
node-bravo bravo online claude 11s 2 4/6
node-charlie charlie offline gemini never 0 6/6
$ bernstein cluster nodes # node table only
$ bernstein cluster nodes --json-output
- Heartbeat is the age since the node's last heartbeat (
neverfor a persisted-but-not-yet-rejoined node). - Claimed is the node's self-reported
active_agents- the tasks it has claimed and is actively running. - Adapter comes from the node's advertised
adapterlabel (workers set it from their configured adapter;-when unset). - Both accept
--server-url(defaultBERNSTEIN_SERVER_URL) and--json-output. They read the endpoints below, so a stopped server prints a clear "cannot connect" hint instead of a stack trace.
Endpoints relevant to cluster operations:
GET /cluster/nodes[?status=online|cordoned|draining|offline]- full node list (task_cluster.py:163-177).GET /cluster/status- aggregate summary (task_cluster.py:180-194): topology, total/online/offline node counts, total capacity, available slots, active agents, full node list.- Per-node fields surfaced:
last_heartbeat,registered_at,capacity.{max_agents,available_slots,active_agents,gpu_available,supported_models},labels,cell_ids(cluster.py:274-292).
Pair these with the standard observability surface
(/metrics, /grafana/dashboard, /slo) - see
Observability overview. The fleet aggregator
in core/fleet/aggregator.py is the right primitive when you have multiple
Bernstein projects you want rolled into a single dashboard; it scrapes each
project's /cluster/status and /status endpoints.
For node JWT issuance and revocation, audit events flow through the
standard audit log (core/security/audit.py); see
Security and identity for the integrity
guarantees and how to export them.
Governing workloads that are not ours
bernstein cluster and the Helm chart govern the orchestrator's own workload.
Agent workloads that run next to it on the same cluster are governed by
declaration instead: a workload carries the label bernstein.io/govern and is
inventoried, routed, and diffed without its manifest being edited by Bernstein.
kubectl get deploy,sts -A -o json > workloads.json
bernstein cluster govern-inventory --manifests workloads.json --json > inventory.json
The inventory holds one record per workload, not one per governed workload:
| Label | State | Meaning |
|---|---|---|
bernstein.io/govern: enabled | governed | Telemetry is routed to the ingest boundary. |
bernstein.io/govern: disabled | opted_out | Declined governance, still listed. |
| (no label) | ungoverned | Nobody enrolled it. Listed, so it can be found. |
A label value that is neither spelling is refused rather than resolved to a posture nobody declared, so a typo fails loudly instead of reading as "not governed".
Pass a previous inventory to see what changed:
bernstein cluster govern-inventory --manifests workloads.json --previous inventory.json
A workload that drops the label reports opted_out; a workload that is gone from
the cluster reports withdrawn. The two are distinct, and neither is a row that
quietly stops appearing.
The inventory document is in the shape bernstein governance plan consumes, so
the same file feeds the posture diff without a second format. inventory_hash
is a pure function of the listing's contents, not its order: two operators
running against the same cluster get the same hash.
Governed workloads' OTLP spans reach the ingest boundary under the source label
k8s:<namespace>/<kind>/<name>, so the signed receipt names the workload the
spans came from. Nothing here sits in a workload's data path, and there is no
admission webhook: a workload is never blocked from starting.
Code pointers
| Concern | File |
|---|---|
| Cluster routes | src/bernstein/core/routes/task_cluster.py |
| NodeRegistry, persistence | src/bernstein/core/protocols/cluster/cluster.py |
| JWT cluster auth | src/bernstein/core/protocols/cluster/cluster_auth.py |
| Task-stealing policy | src/bernstein/core/protocols/cluster/cluster.py:340-393 |
| Worker CLI | src/bernstein/cli/commands/worker_cmd.py |
| Worker run loop | src/bernstein/cli/commands/worker_cmd.py:50-368 |
| Heartbeat client (library) | src/bernstein/core/protocols/cluster/cluster.py:396-... |
| Cluster autoscaler (optional) | src/bernstein/core/protocols/cluster/cluster_autoscaler.py |
| Fleet aggregator | src/bernstein/core/fleet/aggregator.py |
| Workload governance inventory | src/bernstein/core/govern/cluster_inventory.py |
| Models / data classes | src/bernstein/core/models.py - NodeInfo, NodeCapacity, NodeStatus, ClusterConfig |