Skip to content

Task Management

A task in Aether is a finite unit of work with a typed identity, a claimable lifecycle, durable state, and a built-in audit trail. Tasks are distinct from long-running Agents: an agent runs continuously and holds an identity lock for its lifetime; a task has a clear start, a terminal state, and an explicit claim mechanism that makes it safe to hand off across workers or retry after failure.

AgentTask
LifetimeLong-running, reconnects after disconnectFinite — has a terminal state
IdentityGlobally unique, claimed at connectUnique or non-unique depending on type
Work modelContinuous stream processingClaim → execute → complete/fail
ReconnectResumes from last stream offsetLease tracked; DisconnectReaper fails stale tasks
AuditConnection lifecycle eventsFull state-machine audit trail per record

A unique task has a human-chosen specifier that acts as a globally exclusive name. Only one connection can hold a given unique task identity at a time.

Topic: tu::{workspace}::{implementation}::{unique_specifier}

Use unique tasks when:

  • The work item has a natural key (e.g., nightly-backup, report-2024-Q1, order-88421)
  • You want exactly-once execution semantics enforced at the identity layer
  • An orchestrator should start a dedicated process for a named job

A non-unique task has no caller-chosen specifier. The gateway assigns each connected worker a server-generated ID. Multiple workers with the same implementation can connect simultaneously and compete for work broadcast on the shared pool topic.

Topics:

  • Direct (worker-specific): ta::{workspace}::{implementation}::{server_assigned_id}
  • Broadcast (pool): tb::{workspace}::{implementation}

Use non-unique tasks when:

  • You want a pool of workers competing for queued items
  • Each work item does not need a named identity
  • You are building a fan-out processor or parallel pipeline stage

Aether does not prescribe how work reaches workers. The topic primitives support both push and pull patterns — mix and match as your architecture requires.

Workers self-subscribe to the broadcast topic tb::{ws}::{impl} and compete to process incoming messages. Each worker connected to the pool receives the same broadcast; the application layer handles claiming (e.g., a worker that receives a broadcast does an atomic claim before processing).

Workers W1, W2, W3 all subscribe to tb::default::image-processor
→ A message arrives on tb::default::image-processor
→ All three receive it
→ First to atomically claim the work item proceeds; others discard

This pattern is analogous to a traditional competing-consumers queue but is built on Aether’s stream primitives, giving you replay, offset tracking, and audit for free.

A specific task ID is dispatched directly to a known worker via the ta:: topic. This is the mode used by the orchestration subsystem when it assigns a startup task to a named agent/task identity. Callers that know the exact worker (e.g., an orchestrator or a parent agent tracking child task IDs) use this pattern.

Orchestrator → ta::default::image-processor::a1b2c3d4
→ Only worker a1b2c3d4 receives the message

A worker can claim a unique task identity it selects itself by connecting as a tu:: principal with a specifier it chooses. If the identity lock is available, the worker acquires it and begins processing; if another worker already holds it, the connection is rejected.

Self-assignment is still a push semantics — the worker pre-commits to a specific identity rather than competing for whatever the pool happens to broadcast. It’s how you implement leader-election semantics, singleton jobs, or named reservations without a separate coordination layer.

Tasks tracked by the orchestration subsystem follow this state machine:

pending → assigned → starting → running → completed
↘ failed → (retry w/ backoff) → pending → … → dlq
↘ cancelled
StateDescription
pendingCreated, waiting for an orchestrator to claim it
assignedAn orchestrator has claimed the dispatch record
startingOrchestrator is launching compute; worker not yet connected
runningWorker has connected and validated its auth token
completedWorker disconnected gracefully (clean EOF)
failedWorker disconnected unexpectedly, or explicitly failed
cancelledClient sent TaskOperation_CANCEL
dlqMax retries exhausted; moved to dead-letter queue

The following audit event types are recorded at each transition:

created → assigned → starting → started → completed / failed / cancelled / retry_scheduled / timed_out / moved_to_dlq

The dispatch claim is atomic via a PostgreSQL UPDATE ... WHERE status = 'pending'. Only one gateway wins per task, even when all gateways receive the NOTIFY simultaneously. The claim is idempotent — a second claim attempt returns ErrTaskAlreadyClaimed and the gateway that lost silently ignores the notification.

When an orchestrator fails to deliver a TaskAssignment (or the launched process never connects), the gateway unclaims the dispatch record and increments retry_count. The next retry is delayed by 2^retry_count seconds (2s, 4s, 8s, …). After max_retries attempts, the record is moved to dlq and the task is marked failed.

The DisconnectReaper runs on every gateway instance and periodically scans tasks whose worker has been disconnected longer than the task’s configured grace window. If no active session exists for the task (checked via the in-process session index), the task is failed. This handles cases where a worker crashes after connecting but before completing.

The reaper is multi-gateway safe: FailTask is a state-machine transition; concurrent calls on already-terminal tasks are no-ops.

Auth tokens issued by the gateway for orchestrated startups have a 24-hour TTL (Redis-backed). If an orchestrator does not start compute within that window, the token expires and the launched process cannot authenticate. The task will be failed by the reaper or retry mechanism.

Every state transition is recorded in task_audit_events with the event type, event data, and created_by attribution. The audit trail is append-only and is available via the admin API.

When an orchestrated worker connects, the gateway populates the task_context field of ConfigSnapshot from the TaskAssignment. The map contains:

KeyValue
task_idThe orchestration task ID
workspaceThe workspace the task was launched into
userThe user identity that triggered the launch (if any)
implementationThe agent/task implementation
specifierThe agent/task specifier
profileThe orchestration profile used
lp.*All launch_params entries, prefixed with lp.
Metadata keysAll entries from the TaskAssignment.metadata map

Workers that were not started via orchestration receive an empty task_context.

Workers can checkpoint durable state via the CheckpointRequest / CheckpointResponse proto exchange. Checkpoint data is stored gateway-side and can be read back by the same worker (or a replacement after failure) on reconnect. The checkpointed audit event is recorded on each successful checkpoint.

Task management and orchestration are related but distinct concerns:

  • Task management covers the lifecycle, claim atomicity, durability, and audit of work items.
  • Orchestration is the mechanism for provisioning compute on demand when an offline target needs to be started.

JIT agents are the intersection: orchestration provisions the worker; task management tracks whether it started, ran, and completed correctly.

See Orchestration & JIT Agents for the full provisioning flow.