Controller Worker Flow
Overview
The worker claims API-admitted PostgreSQL work, rechecks authorization, invokes Compute, and persists results under its live claim. This trace follows Namespace, Agent stop/deletion, and AgentRevision work. The controller reference owns the contract and the deployment guide owns process setup.
Entry Points
- Trigger: Compose or Helm starts
apps/controller/src/worker.mjs; an authenticated API mutation commits Namespace, Agent lifecycle, or AgentRevision work. - Source:
apps/controller/src/worker.mjs:configuration,apps/controller/src/worker.ts:ControllerWorker.start, andpackages/occ/src/state/postgres-state.ts:operations.append. - Assumptions: The initialized database contains the singleton Installation; the API and worker share its application role and selected Drivers. Production supplies trusted Installation YAML. Work records carry the original actor and exact resource ownership.
Flow
graph TD
subgraph Admission["HTTP API and PostgreSQL"]
A["Authenticate and authorize the mutation"] --> B["Commit state, audit, and work together"]
end
subgraph Worker["Independent controller worker"]
C["Validate startup and attach lifecycle Drivers"] --> D["Recover expired claims and claim eligible work"]
B --> D
D --> E["Reload exact resources and current IAM state"]
E --> F{"Authorized and valid?"}
F -->|yes| G["Invoke Compute while renewing the claim lease"]
F -->|no| H["Persist permanent failure under the live claim"]
G --> I{"Observed result"}
end
subgraph Outcome["Claim-protected result and next handoff"]
I -->|ready| J["Publish lifecycle result and complete work"]
I -->|pending| K["Defer without spending failure budget"]
I -->|temporary failure| L["Retry within the attempt budget"]
I -->|invalid or exhausted| H
K --> D
L --> D
H -->|"Initiating caller repeats Agent DELETE"| A
endExecution Trace
1. Initialize the independent worker process
apps/controller/src/worker.mjs:configuration,
apps/controller/src/worker.ts:ControllerWorker.start
The entrypoint validates configuration,
removes stale readiness, and constructs ControllerWorker with an
application-role pool. Development without OCC_CONFIG_PATH preflights Docker;
production requires explicit configuration.
start() loads the bootstrapped Installation, validates native IAM, and attaches
selected Configuration, Sandbox, and IAM hooks to Compute. Shared composition
supplies Kubernetes Compute's optional Sandbox Driver. Selected hooks require
setLifecycleDrivers; unsupported capabilities stop startup. Production runs
Compute preflight before emitting worker.started and entering run().
Metrics scrapes share one read-only connection and
packages/occ/src/state/postgres-metrics.ts:PostgresMetricsSnapshot.collect
for lifecycle and backlog observations without runtime probes. Metrics follow
finalization independently of logging; see the metrics contract.
2. Commit API admission and the durable work record
apps/controller/src/index.ts:perform,
packages/occ/src/index.ts:OpenClawController,
packages/occ/src/state/postgres-state.ts:operations.append
After authentication and authorization, controller operations call
operations.append, which verifies exact ownership and invokes
PostgresWorkQueue.enqueue. State, admission audit and work commit or roll back
together.
Agent deployment keeps the validated deploy authorization request and decision
with the admitted revision until the API appends its audit event. The event uses
that decision's IAM Driver, principal, exact Agent target and evidence even if
Driver selection changes before the append. The API does not reauthorize the deploy decision to construct the event or
relabel it from current selection. Audit failure rolls the revision,
desired runtime state and queued work back with the transaction. An unknown
PostgreSQL commit outcome remains unknown and is not retried.
The queue freezes actor, Namespace owner, lifecycle target, and exact Agent and
immutable AgentRevision for revision work. Agent lifecycle work identifies its
Agent and stopped or deleted target without a revision. Reusing an idempotency
key with a different actor, owner, or target is rejected. The API returns accepted
state before Compute; the worker takes over.
For an already-deleting Agent, OpenClawController.deleteAgent leaves active
work unchanged. The initiator can retry terminal failure; once it lost
permission, another permitted caller takes over as the work's actor, audited as
takeover. After checking delete permission, operations.retryFailedAgentDeletion
resets only the stopped, deleting Agent's terminal work, keeping prior audits.
The worker reauthorizes normally.
deleteNamespace recovers Namespace teardown the same way via
retryFailedNamespaceDeletion, including takeover.
3. Recover expired claims and claim one eligible operation
apps/controller/src/worker.ts:ControllerWorker.run,
packages/occ/src/state/postgres-work-queue.ts:PostgresWorkQueue.claim
Each loop first calls recoverStale(), then claim(). The
PostgreSQL queue selects
eligible queued work with FOR UPDATE SKIP LOCKED, assigns a fresh claim token
and lease deadline, and increments the attempt count. Another live claim for
the same Agent, or the Namespace for Namespace work, prevents concurrent
ownership of that target.
An empty queue waits within a bound. After work or idle, health() queries
work, calls onHealthy, and emits worker.health. Both must succeed for readiness. Serialized failures emit HEALTH_UNAVAILABLE without
consuming retries.
4. Reload ownership and reauthorize before infrastructure effects
apps/controller/src/worker.ts:ControllerWorker.process,
apps/controller/src/worker.ts:ControllerWorker.processRevision,
apps/controller/src/worker.ts:ControllerWorker.authorizeRevision
Namespace work reloads its exact resource and expected provisioning or
deleting status. Work whose target has already changed completes as
SUPERSEDED_TARGET. Otherwise authorize() reloads IAM state and checks the
original actor. Provisioning also checks restrictions on the exact Namespace;
placement into an existing Kubernetes namespace requires Installation
administration permission again.
Revision work reloads its Namespace, Agent, admitted revision, and active revision.
processRevision() rejects mismatched owners, an unready Namespace, an invalid
Agent Principal, a changed Harness descriptor, or a different Compute Driver.
authorizeRevision() rechecks current deploy and, for a ServiceAccount
snapshot, read on that exact ServiceAccount. Before Compute effects, the worker resolves
frozen Backend metadata and rechecks each managed credential's exact Backend, Driver, workspace, and issued
account binding. This read-only path has no Backend client or admin key. The
Backend-managed credential delivery flow owns these checks.
Revocation and denial fail permanently before runtime creation. Older revisions complete as superseded; active revisions enter finalization or maintenance. When Compute requires stopped predecessors, a newer admission supersedes older active maintenance before Compute effects, even after candidate failure. Recovery uses a new revision.
Agent-stop work rechecks current exact-Agent operate. Superseded desired state
completes without shutdown. An absent active pointer does not prove candidates
have no runtime resources, so stop still checks the captured revision history.
Agent-deletion work requires the Agent to remain deleting and stopped, then
rechecks the original actor's exact-Agent delete. It loads every owned revision
and rejects ownership or Compute-Driver mismatches before teardown.
5. Invoke Compute while renewing the live claim
apps/controller/src/worker.ts:ControllerWorker.observe,
apps/controller/src/worker.ts:ControllerWorker.observeRevision,
apps/controller/src/worker.ts:ControllerWorker.withClaimHeartbeat
Namespace dispatch calls ensureNamespace or deleteNamespace. Revision
dispatch optionally binds the exact Agent, then calls prepareRevision with its
immutable snapshot. A wrong owner or invalid observation fails permanently;
a pending observation defers convergence.
apps/controller/src/worker.ts:ControllerWorker.prepareRevision checks Compute's
requiresStoppedPredecessors capability. When selected, it loads earlier
snapshots, including failed candidates, closes their credential sessions, and
calls stopRevision under the claim heartbeat before preparing the candidate; a
release failure prevents preparation. Later passes re-stop after failures and
doubling lease intervals. The per-Agent queue serializes work, and the dispatch
guard keeps maintenance from recreating a predecessor. This path accepts
downtime; recovery needs a higher revision.
The worker validates Compute's startup plugin warning codes and selection keys against the immutable revision. Compute must verify failed selections are disabled before warnings permit deployment; missing or malformed evidence cannot establish success.
Agent-stop dispatch captures and validates the Agent's revisions owned by the current Compute. It stops the active revision first, then the rest, including terminal candidates and predecessors whose retirement failed. Historical revisions pinned to another Compute are excluded; an active revision pinned elsewhere fails closed. Before each shutdown, the worker rechecks the Agent owner and stopped desired state. Later admissions are not added to this cleanup set. Partial failure retries the idempotent shutdowns without clearing the active pointer or deleting retained workspace data.
Before shutdown or retirement, including stopped-revision recovery, the worker
binds the server-owned Namespace and Agent after IAM and exact-resource checks.
Preparation and maintenance recheck desiredRuntimeState; a candidate that
overlaps stop is shut down instead of activated.
Agent-deletion dispatch binds the server-owned Namespace and Agent, retires every owned revision, then invokes the optional Agent credential-deletion capability. This rebuilds Driver-local ownership after restart. A Driver that can provision runtime credentials but cannot delete them fails permanently before binding or retirement. Compute retirement owns workload termination and Sandbox cleanup; the worker invokes neither independently.
withClaimHeartbeat() renews before each effect and every third of the lease
duration, protecting sequences of short effects too. Lease loss, heartbeat failure,
or shutdown aborts Compute and raises WorkClaimLostError. Expired or replaced
claim tokens cannot publish results.
Successful renewals request throttled health updates without delaying effects or renewal. Health failure does not imply lease loss; heartbeat failure aborts Compute.
Compute owns infrastructure and Sandbox dispatch. See the Kubernetes implementation and Docker execution flow.
6. Persist the result and finish revision activation
apps/controller/src/worker.ts:ControllerWorker.finalize,
apps/controller/src/worker.ts:ControllerWorker.finalizeRevision,
apps/controller/src/worker.ts:ControllerWorker.completeActivatedRevision
transactWithQueue() renews the exact claim before atomically publishing
Namespace readiness or completed deletion, lifecycle evidence and queue completion.
Failed provisioning can transition to failed; incomplete deletion cannot publish
successful deletion.
After preparation, activationOrder: beforeCommit activates before the database
pointer changes. Otherwise, an implemented activation stage runs in development
and production after the claim-protected compare-and-set of Agent.activeRevisionId;
the first dedicated revision stays inactive until commit. A changed pointer causes
ACTIVE_REVISION_CHANGED and retry.
After the pointer commit, the worker finishes activation and retires the
predecessor. In a second transaction, completeActivatedRevision() rechecks the
active revision and claim, appends evidence, and completes work. Infrastructure
effects and database state are not atomic. Retried finalization rechecks the
active candidate before effects, so a retained plugin failure cannot become
success after claim loss.
Stop finalization rechecks the live claim, Agent owner, and stopped desired state.
After captured cleanup, it clears activeRevisionId only if it still identifies
the stopped revision. A later deployment supersedes the stop even if it retains
that pointer while preparing. Finalization appends lifecycle-stop evidence and
completes the work item without deleting revision rows or persistent runtime state.
After commit, completeActivatedRevision() and finalizeAgentStop() record
admission-to-completion duration through apps/controller/src/metrics/index.ts:createOccMetrics.
Queue waits and retries count; maintenance and superseded work do not. Process
death before observation can lose a sample.
Deletion finalization uses a restricted database function rather than the
generic queue completion path. In one transaction it validates the live claim,
removes the Agent's revisions, service principal, API keys, and exact IAM
references, records lifecycle-delete success, deletes the Agent, and removes its
work rows. An expired or replaced claim removes nothing; occ_app has no direct
table-level delete privilege for these records.
finalizeAgentDeletion() records committed agent_delete outcomes. Snapshots
count deleting Agents as stopping, permanent cleanup failures as failed,
and remove completed deletions from inventory.
7. Defer, retry, or stop and hand off the next iteration
apps/controller/src/worker.ts:ControllerWorker.finalizeActiveRevision,
packages/occ/src/state/postgres-work-queue.ts:PostgresWorkQueue.defer,
packages/occ/src/state/postgres-work-queue.ts:PostgresWorkQueue.retry
Pending convergence refunds the attempt, requeuing unready revisions after 500 ms
and others with backoff. Dependency failures consume attempts; permanent failure,
exhaustion, deadline, or AUTHENTICATION_FAILED terminates work. See
outcomes and
timing controls.
ControllerWorker.processRepositoryCleanup defers every incomplete pass at the
Driver interval, including closing sessions and failed runtime retirement,
releasing the claim without consuming retries. Obligations survive; lease loss
aborts the pass.
Terminal rows store reason_code and optional result_data: success
{ warnings: [...] }; convergence-deadline failure timeoutMs and optional
runtimeFailure. Compute reads cached startup results from its private status
path, including unready Harnesses without plugins, and verifies runtime incarnation
without repeating the model probe. Missing or invalid evidence leaves cause unspecified.
packages/occ/src/state/controller-work.ts:validateFailureData validates reads
and writes; the PostgreSQL constraint enforces the persisted shape.
Other failure reasons reject data.
PostgresWorkQueue.complete and PostgresWorkQueue.fail publish only under the
live claim; deployment status derives error and warnings from that result.
Completion needs no acknowledgment or post-commit cleanup; maintenance
cannot rewrite deployment warnings.
See deployment status for result semantics.
Indexed workId scopes progress to work; maintenance and unbound history cannot
supply it. getDeploymentStatus reads
findWorkAttempt with the work row in one State snapshot and projects fixed
public explanations. Memory State has no attempt.
Legacy terminal rows derive reason_code from matching activation or terminal
reconcile audit evidence, otherwise LEGACY_OUTCOME_UNKNOWN. Their
result_data remains NULL; pending rows have no terminal outcome.
If Compute declares maintenance, activation schedules exact-revision observations. Incomplete observations, Compute bindings, and dependency retries or expired claims exhausting attempts close the item and schedule another while the revision stays active and running within its credential deadline; outages never retire it. Each claim reauthorizes its actor. Successor keys use strictly later time buckets despite clock skew.
worker.completed reports the target, outcome, and code; polling continues.
Lease loss reports worker.error CLAIM_LOST instead of stale lifecycle state.
On SIGTERM or SIGINT, shutdown removes readiness, aborts in-flight work, waits for the loop, closes PostgreSQL, and emits
worker.stopped. Each PostgresWorkQueue.recoverStale() statement atomically publishes exhausted
work, failure of a still-provisioning Namespace targeted for ready, and audit
evidence for expired claims and exhausted queued work, so final-attempt crashes
cannot strand provisioning.
Debugging and Verification
worker.startednamescomputeDriverIdand optionalsandboxDriverId. Readyworker.healthconfirms a pending-work query; neither proves a model turn.worker.startup-errorreports invalid mode, database, Installation, Driver selection, or preflight before processing. Probes inspect a private readiness marker; the worker has no HTTP endpoint.- For queued operations, compare API and worker database and Installation
configuration, then inspect
worker.completedandworker.error. Check current IAM state forACTOR_REVOKEDorAUTHORIZATION_DENIED;DEPENDENCY_UNAVAILABLEis retryable;CLAIM_LOSTends publication ownership. - Revision and stale-claim tests require PostgreSQL; neither proves real model execution.
- OCC API checks deploy audit attribution and append-failure rollback on the authenticated route after changing IAM Drivers.
- Sandbox startup verifies composition; k3d integration verifies infrastructure.
