Skip to content

Orchestration & JIT Agents

Aether supports just-in-time (JIT) agent loading: when a message is sent to an offline target, the gateway automatically triggers compute provisioning through a registered Orchestrator. The sender does not need to know whether the target is running.

Spinning up agents only when work arrives — and letting them idle out when it stops — unlocks several architectural advantages:

  • Serverless and sandboxed compute. Each agent invocation can run in an isolated container or sandbox. The identity connection serves as the boundary; there is no persistent process to manage.
  • Per-invocation version upgrades. Because each JIT launch picks up the current image or binary, rolling out a new agent version is a deploy + idle-out cycle. No coordinated restart required.
  • Cost efficiency. Compute is only consumed when a message is waiting. Agents that are never messaged are never started.
  • Natural back-pressure. If no orchestrator is available, the message is durably stored in the stream and the task stays pending. The sender sees no error.

This pattern is a first-class design goal in Aether, not an afterthought. The orchestration subsystem is the mechanism that makes it work.

The orchestration pattern has three participants:

  • Sender — any connected client that sends a message to a target topic
  • Gateway — detects the target is offline and triggers provisioning
  • Orchestrator — a connected client that receives task assignments and starts compute (containers, subprocesses, VMs, etc.)

The flow is triggered automatically on any SendMessage targeting an ag:: or tu:: topic where the target is not connected.

When the gateway receives a SendMessage targeting an agent topic (ag::) or unique task topic (tu::):

  1. Local check: identityIndex (O(1) in-process lookup) — is the target connected to this gateway instance?
  2. Distributed check: Redis EXISTS on the identity’s lock key — is the target connected to any gateway instance?
  3. If offline: Create an orchestration task record and publish a dispatch notification.

The message is always published to the stream regardless of target availability. This ensures it is waiting when the target eventually connects.

Sender Gateway PostgreSQL Dispatcher Orchestrator
| | | | |
|-- SendMessage ->| | | |
| (target offline) | | |
| |-- Create task ---->| | |
| |<- task_id ---------| | |
| |-- NOTIFY / poll ---------------------->| |
| | | (all gateways receive) |
| |-- Claim task ----->| | |
| |<- claimed ---------| | |
| |-- TaskAssignment ------------------------------------------->|
| | (task_id, target identity, launch params, auth token) |
| | | | |
| | | | (start compute)|
| |<============ Target connects with auth token ============|
| |-- Validate token ->| | |
| |-- Mark running --->| | |
|<- ConnectionAck | | | |
|<- ConfigSnapshot| | | |
|<- Queued msgs via stream offset replay | |
  1. Gateway creates a task record in PostgreSQL with the target identity and the matching orchestrator profile.
  2. The orchestration dispatcher notifies connected gateways. In full mode this uses PostgreSQL NOTIFY (pq.Listener); in AetherLite (embedded/lite mode) it uses an in-memory polling loop (2-second interval). AMQP/RabbitMQ is the message-routing substrate — it is not involved in the orchestration dispatch path.
  3. All connected gateway instances receive the notification.
  4. Each gateway checks whether it has a locally connected orchestrator matching the required profile and workspace.
  5. Only one gateway atomically claims the task (distributed claim via PostgreSQL UPDATE ... WHERE status = 'pending').
  6. The claiming gateway sends a TaskAssignment to its local orchestrator via the gRPC stream. The assignment includes:
    • Task ID and target identity
    • Workspace and launch parameters
    • A short-lived authentication token (Redis-backed, 24-hour TTL)
  7. The orchestrator starts compute (container, subprocess, VM) configured with the token.
  8. The new process connects to the gateway as the target identity using the token as its credential.
  9. The gateway validates the token, acquires the distributed lock, and marks the task as running.
  10. The new agent/task receives the ConfigSnapshot (including task context) and then replays all queued messages from the stream via consumer offset tracking.

Orchestrators register their supported profiles when they connect. An orchestrator profile declares which types of agents it can launch.

// Go SDK: OrchestratorClient with supported profiles
client, err := aether.NewOrchestratorClient(aether.OrchestratorOptions{
ClientOptions: aether.ClientOptions{
ServerAddr: "localhost:50051",
},
Implementation: "kubernetes-orchestrator",
SupportedProfiles: []string{"docker", "kubernetes"},
})

The gateway matches incoming task assignments to available orchestrators by profile. If no matching orchestrator is connected, the task waits in pending state.

Orchestration tasks transition through these states:

pending → assigned → starting → running → completed
↘ failed → (retry with backoff) → pending → … → dlq
↘ cancelled
StateDescription
pendingWaiting for an orchestrator to claim the task
assignedAn orchestrator has claimed the task
startingThe orchestrator has acknowledged and is starting compute
runningThe target agent/task has connected and validated its token
completedThe agent/task disconnected gracefully (clean EOF)
failedThe agent/task disconnected unexpectedly, or was explicitly failed
cancelledA client sent TaskOperation_CANCEL
dlqMax retries exhausted; moved to dead-letter queue

Failed delivery to an orchestrator is retried with exponential backoff (2^n seconds). After exceeding max_retries, the record moves to dlq. Background cleanup jobs purge old task records based on configurable retention policies.

When a task transitions state, the gateway automatically notifies the parent agent that created it. This uses the pg::{workspace} progress infrastructure with server-side recipient filtering — only the spawning agent receives the notification.

TransitionTrigger
→ runningOrchestrated agent connects and validates its task token
→ completedAgent disconnects gracefully (clean EOF)
→ failedAgent disconnects unexpectedly, or task is explicitly failed
→ cancelledClient sends TaskOperation_CANCEL
→ pendingClient sends TaskOperation_RETRY

The ProgressUpdate notification includes:

  • task_id — the task that changed state
  • state — the new state string
  • summary — human-readable description (includes error message for failures)
  • recipient — the parent agent’s topic (for server-side filtering)
client, err := aether.NewOrchestratorClient(aether.OrchestratorOptions{
ClientOptions: aether.ClientOptions{
ServerAddr: "localhost:50051",
},
Implementation: "docker-orchestrator",
SupportedProfiles: []string{"docker"},
})
client.OnTaskAssignment(func(ctx context.Context, task *aether.TaskAssignment) error {
// Launch a Docker container with the provided auth token
token := task.AuthToken
image := task.LaunchParams["image"]
// ... start container with AETHER_TOKEN=token AETHER_GATEWAY=address ...
return nil
})
client.Connect(ctx)
client.Run(ctx)
import { BaseOrchestrator } from "@scitrera/aether-client";
import type { TaskAssignment } from "@scitrera/aether-client";
class DockerOrchestrator extends BaseOrchestrator {
async launchTask(assignment: TaskAssignment): Promise<void> {
const { targetImplementation, authToken, launchParams } = assignment;
// Start a container using authToken as the credential
console.log(`Launching ${targetImplementation} with token`);
}
}
const orch = new DockerOrchestrator({
address: "localhost:50051",
implementation: "docker-orchestrator",
supportedProfiles: ["docker"],
logAssignments: true,
});
await orch.connect();
from scitrera_aether_client import OrchestratorClient
client = OrchestratorClient(
implementation="docker-orchestrator",
supported_profiles=["docker"]
)
def on_task_assignment(assignment):
token = assignment.auth_token
image = assignment.launch_params.get("image")
# ... start container ...
client.on_task_assignment = on_task_assignment
client.connect("localhost:50051")

A complete example of the cold-start flow:

  1. User sends "Process this file" to ag::default::gpu-worker::01.
  2. Gateway checks identityIndex — not locally connected.
  3. Gateway checks Redis — no active lock for ag::default::gpu-worker::01.
  4. Gateway publishes message to the ag::default::gpu-worker::01 stream (persisted).
  5. Gateway creates orchestration task targeting gpu-worker implementation, profile kubernetes.
  6. Dispatcher publishes task notification (NOTIFY in full mode, polling in AetherLite).
  7. All gateways receive the notification. Gateway B has a connected Kubernetes orchestrator matching the kubernetes profile.
  8. Gateway B atomically claims the task in PostgreSQL.
  9. Gateway B sends TaskAssignment (with auth token) to the Kubernetes orchestrator.
  10. Orchestrator starts a pod with AETHER_TOKEN=<token> and AETHER_IMPL=gpu-worker.
  11. Pod connects to the gateway as ag::default::gpu-worker::01 using the token.
  12. Gateway validates the token, acquires the Redis lock, marks the task as running.
  13. Agent receives ConfigSnapshot (with task context) then replays the persisted "Process this file" message from the stream.

The sender saw no error. From its perspective, the message was sent to a valid topic and will be delivered when the target is ready.

  • Orchestration only triggers for offline ag:: and tu:: targets. Non-unique tasks (tb::), users, and broadcast topics do not trigger orchestration.
  • If no matching orchestrator is connected, the task stays pending indefinitely. The sender receives no error — the message is safely stored in the stream.
  • Task tokens have a 24-hour TTL. If an orchestrator does not start compute within that window, the token expires and the launched process cannot authenticate.
  • AetherLite mode uses an in-memory polling dispatcher (2-second interval) instead of PostgreSQL NOTIFY, but the task assignment flow is otherwise identical.