Skip to content

Architecture

Aether is a distributed control plane for routing structured messages, tracking tasks, and managing connection lifecycles. This document provides a comprehensive overview of the system architecture, core concepts, and design principles.

Aether coordinates external agents, tasks, and engines through a gRPC gateway backed by three categories of infrastructure:

  • Message broker — persistent, replayable message routing
  • KV store — distributed session locks, configuration storage, state management, and checkpointing
  • Relational database — persistent features: task lifecycle, ACL rules, audit logs, orchestration profiles

Full mode uses RabbitMQ Streams (broker), Redis or Valkey (KV store), and PostgreSQL (relational database). AetherLite replaces all three with in-process equivalents: a Badger-backed persistent log (broker), Badger (KV store), and SQLite (relational database). No external services are required in lite mode.

The defining architectural property: the active gRPC connection itself functions simultaneously as the distributed lock for an identity AND as its liveness proof. When the TCP stream closes, the lock releases automatically. This eliminates the need for separate heartbeat APIs and simplifies failure detection.

┌──────────────────────────────────────────────────────────────────┐
│ External Clients │
│ (Agents, Tasks, Users, Orchestrators, WorkflowEngine, Services) │
└────────────────────────┬─────────────────────────────────────────┘
│
┌────────▼──────────┐
│ Gateway Server │
│ (gRPC Stream │
│ Handler) │
└────────┬──────────┘
│
┌────────────────┼────────────────┐
│ │ │
┌──────▼──────┐ ┌─────▼─────┐ ┌──────▼──────┐
│ Router │ │ Session │ │ Admin │
│ │ │ Registry │ │ Server │
│ Topic-to- │ │ │ │ │
│ Stream │ │ (KV-store │ │ (REST API │
│ Mapping │ │ locks) │ │ + UI) │
└──────┬──────┘ └───────────┘ └─────────────┘
│
┌──────┴──────────────────────────────────┐
│ │
┌────▼─────────────────┐ ┌────────▼─────┐
│ ACL & Auth Services │ │ Orchestration │
│ │ │ │
│ - RBAC (Casbin) │ │ - Task dispatch │
│ - Authority grants │ │ - Task claims │
│ - Token validation │ │ - Profile mgmt │
└──────────────────────┘ └──────────────────┘
ComponentPurpose
Gateway ServerHandles bidirectional gRPC streams, authentication, connection lifecycle, message routing, KV/checkpoint operations
RouterMaps topics to broker streams, manages producer pools with idle eviction, implements shared consumer fan-out with offset tracking. Current backend: RabbitMQ Streams (full mode), Badger-backed log (AetherLite). The MessageRouter interface is designed to support additional broker backends in the future.
Session RegistryAtomic conditional lock acquisition at the KV-store layer with 30-second TTL, session metadata storage, stale lock cleanup
KV StoreHierarchical configuration store with eight scopes across two axes: identity scope (global, workspace, user, user-workspace) × sharing (shared vs exclusive). Backed by Redis in full mode; Badger in AetherLite.
Checkpoint StorePersistent state checkpointing for agents and tasks, backed by KV store
Task StoreRelational database-backed task lifecycle management (create, assign, complete, fail, retry, purge)
ACL ServiceRole-based access control using Casbin v3, authority grants, workspace access verification
Audit LoggerBatched, configurable event capture with retention policies; tracks connections, auth, messages, KV, admin, and ACL events
OrchestrationTask dispatch via pluggable TaskDispatcher interface (PostgreSQL NOTIFY + polling in full mode; polling-only in AetherLite), atomic claim-based delivery, orchestrator profile management
Admin ServerREST API and embedded UI for operations; separate ops server on port 9090 for health probes and Prometheus metrics
Auth ServicesmTLS certificate validation (strict / semi-strict / relaxed modes), task token verification, API key validation, OAuth/JWT validation

Aether recognizes eight principal types that determine identity uniqueness and subscription patterns:

TypeIdentity FieldsUniquenessPurpose
Agentworkspace + implementation + specifierExactly one connectionLong-running service or worker
Task (Unique)workspace + implementation + unique_specifierExactly one connectionNamed finite unit of work (e.g., “nightly-backup”)
Task (Non-Unique)workspace + implementation (server assigns ID)Multiple connections allowedWorkers competing on shared broadcast
Useruser_id + window_idOne per windowHuman operator; multiple browser tabs allowed
Workflow Engineoptional workspace scopeOne per shard (currently 1 shard)Subscriber to event::receiver{N} fan-in topics
Metrics Bridge(none)One per shard (currently 1 shard)Subscribes to metric::receiver{N}; receive-only
Orchestratorimplementation + specifierOne per specifierReceives task assignments; spins up compute
Serviceimplementation + specifierOne per specifierCross-workspace HTTP proxy; sv::{impl}::{spec} topic

This is Aether’s defining architectural property:

  1. Implicit Locking: An active gRPC stream connection represents the distributed lock for that identity. No separate lock API exists.

  2. Exclusivity: Two clients cannot hold the same identity simultaneously. A second connection attempt for an occupied identity is rejected with DuplicateIdentityError.

  3. Liveness as Heartbeat: The TCP/gRPC connection state serves as the heartbeat. The gateway refreshes a KV-store lock TTL (30s) every 10 seconds while the connection is alive. If the gateway crashes, the lock auto-expires.

  4. Session Resume: A reconnecting client may provide a resume_session_id to atomically take over an existing lock. This supports brief network partitions without loss of session state.

Messages carry a message_type enum that determines routing behavior. OPAQUE is the default — the SDK convenience helpers (SendToAgent, SendToUser, etc.) use it unless overridden:

TypeDescription
OPAQUE(Default) Sender and receiver own the payload schema. Aether forwards verbatim with no payload decoding, validation, or per-type behavior. Use for binary/structured data.
CHATConversational messages between participants. Use explicitly for true conversational text.
CONTROLSystem control signals (start, stop, configure)
TOOL_CALLTool invocation requests and responses
EVENTBroadcast events routed to the Workflow Engine fan-in
METRICTelemetry data routed to the Metrics Bridge fan-in
Client Gateway KV Store / Broker
| | |
|-- InitConnection ------------>| |
| |-- Authenticate (mTLS/OAuth) -->|
| |-- AcquireLock (atomic CW) --->|
| |<- Lock granted ---------------|
| |-- ACL Check ----------------->|
| |<- Access granted -------------|
| |-- Quota Check + Increment --->|
| |-- Register Session ---------->|
| |-- Subscribe to topic(s) ----->|
|<-- ConnectionAck (sessionID) -| |
|<-- ConfigSnapshot (KV) -------| |
| | |
|<======= message loop ========>|<====== stream I/O ============>|
| | |
|-- (disconnect) -------------->| |
| |-- Unsubscribe --------------->|
| |-- ReleaseLock --------------->|
| |-- UnregisterSession --------->|
  1. Lock Acquisition: Atomic conditional write at the KV-store layer with 30-second TTL. Supports session resume via an atomic take-over operation.
  2. Authentication: Verify credentials (mTLS, task token, API key, OAuth, or none)
  3. ACL Check: Verify workspace access before session becomes discoverable. System principals skip this check.
  4. Quota Check: Atomically check and increment per-workspace connection count at the KV-store layer
  5. Session Registration: Store session metadata with matching TTL
  6. Topic Subscription: Subscribe to appropriate topics based on principal type
  7. Acknowledgment: Send ConnectionAck with session ID and ConfigSnapshot with KV data
  8. Message Loop: Handle incoming/outgoing messages until disconnect
  9. Cleanup: Unsubscribe, release lock, decrement quota, update task state, audit log

Topics follow a hierarchical naming scheme with :: as the separator (the identity separator) that determines routing behavior:

PrefixTargetFormatDescription
agAgentag::{workspace}::{impl}::{spec}Specific long-running agent instance (unicast)
tuUnique Tasktu::{workspace}::{impl}::{unique_spec}Named task instance (unicast)
taAssigned Taskta::{workspace}::{impl}::{task_id}Server-assigned non-unique task instance (unicast)
tbTask Broadcasttb::{workspace}::{impl}Load-balancing broadcast; workers round-robin
usUser (Window)us::{user_id}::{window_id}Specific browser window (unicast)
uwUser (Workspace)uw::{user_id}::{workspace}User scoped to a workspace (multicast)
gaGlobal Agentsga::{workspace}Broadcast to all agents in workspace
guGlobal Usersgu::{workspace}Broadcast to all users in workspace
pgProgresspg::{workspace}Progress updates with server-side recipient filtering
eventWorkflow Engineevent::{workspace} (rewritten to event::receiver{N})Routed to Workflow Engine fan-in shard
metricMetrics Bridgemetric::{workspace} (rewritten to metric::receiver{N})Routed to Metrics Bridge fan-in shard
svServicesv::{impl}::{spec}Cross-workspace HTTP proxy service
brBridgebr::{impl}::{spec}Cross-workspace messaging integration (experimental)

Topic names are validated against allowed prefixes and must be 1–256 characters.

Each principal type automatically subscribes to specific topics on connection:

PrincipalExclusive (offset-tracked)Shared (no offset)
Agentag::{ws}::{impl}::{spec}ga::{ws}, pg::{ws}
Task (Unique)tu::{ws}::{impl}::{spec}—
Task (Non-Unique)ta::{ws}::{impl}::{id}tb::{ws}::{impl}
Userus::{uid}::{wid}gu::{ws}, uw::{uid}::{ws}, pg::{ws}
Workflow Engineevent::receiver{N} (exclusive per shard)—
Metrics Bridgemetric::receiver{N} (exclusive per shard)—
Orchestrator—— (receives tasks via direct gRPC stream)
Servicesv::{impl}::{spec}—

Exclusive subscriptions use broker consumer offset tracking so messages are replayed from the last committed position on reconnection. Shared subscriptions fan out locally from a single broker consumer without persistent offset tracking.

Sharded Fan-In Design (WorkflowEngine and MetricsBridge)

Section titled “Sharded Fan-In Design (WorkflowEngine and MetricsBridge)”

event:: and metric:: topics use a sharded fan-in design: the gateway rewrites sender-side topic addresses (e.g. event::{workspace}) to a specific receiver shard topic (event::receiver{N}) before publishing. Each Workflow Engine or Metrics Bridge instance subscribes exclusively to one shard, enabling horizontal scaling of these consumers without coordination overhead.

Current implementation: The system runs with a single shard (TotalShards() == 1), so all events route to event::receiver0 and all metrics to metric::receiver0. The sharding logic is in place so that adding shards requires no protocol changes — only the shard count needs updating.

Exclusive subscriptions use broker consumer offset tracking so messages are replayed from the last committed position on reconnection.

SenderCan Send ToCannot Send To
AgentAgents, Tasks, Users, Events, MetricsOrchestrators, Progress
TaskAgents, Tasks, Users, Events, MetricsOrchestrators, Progress
UserAgents, Tasks, UsersEvents, Metrics, Progress
Workflow EngineAgents, Tasks, Users, Events, Metrics—
Metrics BridgeNone (receive-only)All
OrchestratorAgent/Task topics onlyEvents, Metrics, Users, Progress
ServiceAgents, Tasks, Users (via proxy)Direct broadcast topics

Cross-workspace sends are default-deny: workspace-scoped principals cannot target topics in other workspaces without an explicit ACL grant. The trigger-aware ACL model allows cross-workspace sends when the sender holds the appropriate capability/cross_workspace_send permission on the target workspace. This is an explicit policy decision, not a transport-layer hard block — same-workspace sends carry an implicit grant without requiring an ACL rule.

All messages are wrapped in a server-stamped MessageEnvelope before publishing to the broker:

message MessageEnvelope {
string source = 1; // Server-verified sender topic (cannot be spoofed)
bytes payload = 2;
MessageType message_type = 3;
int64 timestamp_ms = 4; // Server-assigned timestamp
}

The server injects the source field and timestamps all messages, preventing client spoofing and providing audit trails.

The KV Store provides hierarchical configuration across eight scopes organized on two orthogonal axes:

  • Identity scope: global (tenant-wide), workspace, user, or user-workspace
  • Sharing: shared (all agents see the same key) or exclusive (per-agent namespace)
ScopeSharingKey PatternPurpose
globalSharedkv:globalTenant-wide shared state; all agents in the tenant read/write the same keys
global-exclusiveExclusivekv:agent:{impl}|{spec}:globalTenant-wide but per-agent; agent’s own global state
workspaceSharedkv:ws:{workspace}Workspace-shared configuration; all agents in the workspace share keys
workspace-exclusiveExclusivekv:agent:{impl}|{spec}:ws:{workspace}Per-agent, per-workspace state
user-sharedSharedkv:user:{user_id}Per-user shared across all agents
userExclusivekv:agent:{impl}|{spec}:user:{user_id}Per-agent, per-user state
user-workspace-sharedSharedkv:user:{user_id}:ws:{workspace}Per-user, per-workspace, shared across agents
user-workspaceExclusivekv:agent:{impl}|{spec}:user:{user_id}:ws:{workspace}Per-agent, per-user, per-workspace state

The separator between impl and spec inside the key is | (pipe), not ::, to avoid ambiguity with the identity separator.

When a message targets an offline agent (ag::) or unique task (tu::), orchestration is triggered:

  1. Gateway checks identityIndex (local O(1)) then the KV-store lock (distributed)
  2. If offline: message is published to broker stream (persisted) AND orchestration task created in the database
  3. The TaskDispatcher notifies a gateway instance of the pending task
  4. One gateway claims the task atomically and sends TaskAssignment to a connected orchestrator
  5. Orchestrator spins up compute with a short-lived auth token (24h TTL)
  6. Target connects, validates token, receives persisted messages via offset replay

This pattern provides lazy loading: agents and tasks only run when there is work to do.

In full mode, Aether uses RabbitMQ Streams (not classic queues) for message routing:

  • Persistent, replayable message logs
  • Consumer offset tracking for at-least-once delivery semantics
  • Per-topic producer pools with health checks and idle eviction (5 minutes)
  • Shared consumers with local fan-out to reduce broker connection count

The Router automatically declares streams on first publish or subscribe. The MessageRouter interface abstracts the broker layer, enabling alternative backends (AetherLite uses a Badger-backed persistent log that satisfies the same interface).

In full mode, Redis (or Valkey) serves multiple purposes via a shared client:

  1. Session Registry: Atomic conditional lock acquisition and session metadata (30s TTL)
  2. KV Store: Namespace-scoped configuration with optional TTL
  3. Checkpoint Store: Persistent agent/task state
  4. Task Tokens: Short-lived orchestration auth tokens (24h TTL)
  5. Quota Counters: Per-workspace connection and message rate tracking

Supports single-node, cluster, and auto-detect modes.

In full mode, PostgreSQL is required. It is not optional. Persistent features depend on it:

  • Task lifecycle management (create, assign, complete, fail, retry, purge)
  • Orchestration profiles and queue management
  • ACL rules and authority grants
  • API token storage (HMAC-SHA256 hashed)
  • Audit log (batched writes with retention)
  • Agent registry

Schema is managed by embedded migrations in server/migrations/ that auto-run on startup. In AetherLite mode, PostgreSQL is replaced by an embedded SQLite database using a compatibility driver that rewrites PG-dialect SQL at the driver layer.

Aether is designed for stateless horizontal scaling:

  • Stateless Gateway Instances: All state stored in external KV store and relational database
  • Distributed Locking: Atomic conditional writes at the KV-store layer ensure identity uniqueness across all instances
  • Session Affinity: Load balancer uses ClientIP or cookies for reconnection optimization
  • Multi-Gateway Orchestration: Task claims use database-backed atomic operations; only one gateway delivers each task
  • Graceful Failover: Lock TTLs expire automatically (30s); clients reconnect and the broker preserves message offsets

Deployment manifests are available in server/deployments/:

  • Docker Compose: server/deployments/docker-compose/multi-instance.yaml — available today; multi-instance setup with nginx load balancer
  • Kubernetes: Manifests in server/deployments/k8s/ — available today
  • Helm: Chart in server/deployments/helm/aether-gateway/ — available today

The gateway supports multiple authentication methods configured via auth.modes in the YAML config:

MethodMechanismIdentity Source
mTLS (strict)Certificate CN fully encodes the identity: {type}.{workspace}.{impl}.{spec}Certificate only
mTLS (semi-strict)Certificate CN pins the principal type + workspace + implementation; the specifier is provided or overridden via InitConnection. Enables a single cert to serve multiple horizontally-scaled instances of the same component.Certificate (type, workspace, impl) + InitConnection (specifier)
mTLS (relaxed)Certificate confirms only the principal type; all identity details come from InitConnectionCertificate (type only) + InitConnection
Task TokenShort-lived KV-store-backed token issued by orchestrationToken validates target identity
API KeyLong-lived database-stored token with workspace pattern matchingInitConnection (token authorizes workspace)
OAuth/JWTJWT validated against JWKS endpoint with claims mappingJWT claims
NoneNo authentication (TLS not required, no credentials)InitConnection only

Precedence: mTLS identity > Task token validation > API key / OAuth credential authentication.

The three mTLS modes are configured via auth.mtls.mode in the gateway YAML: "strict", "semi-strict", or "relaxed". An unrecognized value falls back to strict.

Certificate generation and TLS configuration details are covered in the gateway deployment documentation.

Identity Uniqueness and Workspace Isolation

Section titled “Identity Uniqueness and Workspace Isolation”
  • Agents: Globally unique. Two agents cannot connect with the same (workspace, implementation, specifier).
  • Unique Tasks: Globally unique per named instance.
  • Non-unique Tasks: Multiple connections allowed; each gets a server-generated ID.
  • Users: Unique per window (us::{user_id}::{window_id}), allowing multiple browser tabs.

Workspace isolation is enforced via the ACL layer, not as a hard transport-level block:

  • Same-workspace sends carry an implicit grant — no ACL rule is required.
  • Cross-workspace sends are default-deny — a workspace-scoped principal attempting to target a topic in a different workspace will be rejected unless an explicit ACL rule grants the capability/cross_workspace_send permission on the target workspace.
  • ACL rules can permit cross-workspace sends — the trigger-aware ACL model supports delegation chains and authority grants that span workspaces. Operators can explicitly allow specific principals to reach across workspace boundaries.
  • KV Store namespaces isolate configuration by workspace — workspace-scoped KV keys are keyed by workspace slug in the storage layer and are not cross-readable by default.

Aether Gateway exposes:

  • Prometheus metrics at GET http://localhost:9090/metrics with granular connection, message, KV, and orchestration metrics
  • Health probes at GET http://localhost:9090/health/{live,ready,startup}
  • Admin REST API at http://localhost:31880/api/ with connection, KV, and audit endpoints
  • Embedded Admin UI at http://localhost:31880/ for connection browsing and KV inspection

Metrics reference and dashboard setup details are available in the server repository’s monitoring documentation.

AetherLite is an embedded deployment mode that replaces every external service with an in-process alternative. It runs as a single binary (cmd/aetherlite) or as any individual binary with the --lite flag (e.g., ./gateway --lite). No Redis, RabbitMQ, or PostgreSQL installation is required.

Full ModeAetherLite EquivalentNotes
Redis (session locks)Badger KVPrefix sess: in shared Badger database
Redis (KV store)Badger KVPrefix kv: in shared Badger database
Redis (checkpoints)Badger KVPrefix ckpt: in shared Badger database
Redis (task tokens)Badger KVPrefix tok: in shared Badger database
RabbitMQ StreamsBadger-backed persistent logPer-topic offset tracking in Badger (msg:, off: prefixes)
PostgreSQLSQLite (via sqlite_compat driver)Same schema; PG syntax rewritten to SQLite at driver layer
PostgreSQL NOTIFY dispatcherIn-memory polling dispatcherNo NOTIFY connection; polling replaces PG NOTIFY
Per-workspace quotasIn-memory quota managerQuota counters reset on restart

The gateway is built around a set of pluggable interfaces that both full mode and AetherLite satisfy:

InterfaceFull Mode ImplementationLite Mode Implementation
SessionManagerRedis atomic-conditional registryBadgerSessionRegistry
MessageRouterRabbitMQ Streams routerBadgerRouter
KVReadWriterRedis KV storeBadgerKVStore
CheckpointManagerRedis checkpoint storeBadgerCheckpointStore
QuotaCheckerRedis atomic countersMemoryQuotaManager
TaskDispatcherOrchestratorTaskDispatcher (PostgreSQL NOTIFY + polling fallback)MemoryTaskDispatcher (polling-only)

The gateway server code (internal/gateway/server.go) is identical in both modes — only the backend wiring differs.

The TaskDispatcher interface abstracts the dispatch mechanism. In full mode, OrchestratorTaskDispatcher listens for PostgreSQL NOTIFY events and falls back to polling when no NOTIFY connection is available. In AetherLite mode, MemoryTaskDispatcher is a pure polling implementation with no database NOTIFY dependency.

┌───────────────────────────────────────────────┐
│ AetherLite Process │
│ │
│ ┌─────────────┐ ┌──────────┐ ┌──────────┐ │
│ │ Gateway │ │ Workflow │ │MsgBridge │ │
│ │ gRPC :50051 │ │ Server │ │(optional)│ │
│ └──────┬──────┘ └────┬─────┘ └────┬─────┘ │
│ │ │ │ │
│ └──────────────┴──────────────┘ │
│ │ │
│ ┌─────────────┴──────────────┐ │
│ │ │ │
│ ┌──────▼──────┐ ┌────────▼──────┐ │
│ │ Badger DB │ │ SQLite DB │ │
│ │ │ │ │ │
│ │ - Sessions │ │ - Tasks │ │
│ │ - KV store │ │ - ACL rules │ │
│ │ - Checkpts │ │ - Audit log │ │
│ │ - Tokens │ │ - Orchestrat. │ │
│ │ - Messages │ │ - Agent reg. │ │
│ └─────────────┘ └───────────────┘ │
│ │
│ Data directory: ./aether-lite-data/ │
│ badger/ ← Badger files │
│ aether.db ← SQLite database │
└───────────────────────────────────────────────┘

AetherLite is single-node only. It does not support horizontal scaling. For multi-node deployments, use the full stack. See AetherLite for a complete reference.

Aether’s architecture provides:

  1. Simplicity: Connection state is the single source of truth for identity liveness
  2. Scalability: Stateless gateways with external state management (full mode)
  3. Reliability: Persistent message logs, distributed locking, graceful failover
  4. Flexibility: Multiple auth methods, ACL-gated workspace isolation, eight-scope KV store; pluggable backends for lite or full deployment
  5. Observability: Comprehensive metrics, audit logging, admin UI

The system coordinates external agents, tasks, and engines with minimal operational complexity while maintaining strong consistency guarantees for distributed identity and message delivery.