Python SDK
Python client library for the Aether distributed control plane. Aether is a system for routing structured messages, tracking tasks, and managing connection lifecycles for agents, tasks, users, and other principals.
Features
Section titled “Features”- Sync and Async Support: Both synchronous (threading-based) and asynchronous (asyncio-based) client implementations
- Multiple Client Types: Agent, Task, User, Orchestrator, WorkflowEngine, and MetricsBridge clients
- Key-Value Store: Hierarchical configuration store with multiple scopes (global, workspace, user, user-workspace)
- Task Management: Create and manage tasks with different assignment modes
- Checkpointing: Persist and restore agent/task state across restarts
- Auto-Reconnection: Configurable exponential backoff with automatic reconnection
- TLS/mTLS Support: Secure connections with optional mutual TLS authentication
- Typed + Catch-all Handlers: Register
on_chat_message,on_control_message, etc. alongsideon_message— both fire for matching messages
Installation
Section titled “Installation”pip install scitrera-aether-clientFor development:
pip install scitrera-aether-client[dev]Quick Start
Section titled “Quick Start”Synchronous Client
Section titled “Synchronous Client”from scitrera_aether_client import AgentClient, CHAT
# Create an agent clientclient = AgentClient( workspace="default", implementation="my-agent", specifier="agent-01")
# Set up message callbackdef on_message(msg): print(f"Received from {msg.source_topic}: {msg.payload.decode()}")
client.on_message = on_message
# Connect to the gatewayclient.connect("localhost:50051")
# Send a message to another agentclient.send_message_to_agent( workspace="default", implementation="other-agent", specifier="01", payload=b"Hello!")
# Keep running until interruptedtry: while True: import time time.sleep(1)except KeyboardInterrupt: passfinally: client.close()Asynchronous Client
Section titled “Asynchronous Client”import asynciofrom scitrera_aether_client import AsyncAgentClient
async def main(): client = AsyncAgentClient( workspace="default", implementation="my-async-agent", specifier="agent-01" )
async def on_message(msg): print(f"Received: {msg.payload.decode()}")
client.on_message = on_message
await client.connect("localhost:50051")
await client.send_message_to_agent( workspace="default", implementation="my-async-agent", specifier="agent-01", payload=b"Hello from async!" )
# Wait until disconnected await client.wait_until_disconnected()
asyncio.run(main())Using Async Context Manager
Section titled “Using Async Context Manager”async with AsyncAgentClient("default", "my-agent", "01") as client: await client.connect("localhost:50051") await client.send_message_to_agent("default", "other", "01", b"Hello!") await asyncio.sleep(1)# Connection automatically closedClient Types
Section titled “Client Types”AgentClient / AsyncAgentClient
Section titled “AgentClient / AsyncAgentClient”For long-running agent processes that need unique identities.
from scitrera_aether_client import AgentClient
client = AgentClient( workspace="default", implementation="python-worker", specifier="worker-01")TaskClient / AsyncTaskClient
Section titled “TaskClient / AsyncTaskClient”For task execution. Supports both unique (named) and non-unique (pooled) tasks.
from scitrera_aether_client import TaskClient
# Unique task (named)unique_task = TaskClient( workspace="default", implementation="data-processor", unique_specifier="job-123")
# Non-unique task (pooled, server assigns ID)pooled_task = TaskClient( workspace="default", implementation="worker")UserClient / AsyncUserClient
Section titled “UserClient / AsyncUserClient”For user session connections (e.g., from browser windows).
from scitrera_aether_client import UserClient
client = UserClient( user_id="user-123", window_id="window-abc")OrchestratorClient / AsyncOrchestratorClient
Section titled “OrchestratorClient / AsyncOrchestratorClient”For managing agent/task lifecycle and compute resources.
from scitrera_aether_client import OrchestratorClient
client = OrchestratorClient( implementation="kubernetes-orchestrator", supported_profiles=["docker", "kubernetes"])WorkflowEngineClient / AsyncWorkflowEngineClient
Section titled “WorkflowEngineClient / AsyncWorkflowEngineClient”For processing events and triggering downstream actions.
from scitrera_aether_client import WorkflowEngineClient
client = WorkflowEngineClient()MetricsBridgeClient / AsyncMetricsBridgeClient
Section titled “MetricsBridgeClient / AsyncMetricsBridgeClient”For collecting telemetry data from agents and tasks.
from scitrera_aether_client import MetricsBridgeClient
client = MetricsBridgeClient()Callbacks
Section titled “Callbacks”All clients support the following callbacks:
| Callback | Description | Signature |
|---|---|---|
on_message | Every incoming message (catch-all) | (msg: IncomingMessage) -> None |
on_chat_message | CHAT-typed messages | (msg: IncomingMessage) -> None |
on_control_message | CONTROL-typed messages | (msg: IncomingMessage) -> None |
on_tool_call | TOOL_CALL-typed messages | (msg: IncomingMessage) -> None |
on_event | EVENT-typed messages | (msg: IncomingMessage) -> None |
on_metric | METRIC-typed messages | (msg: IncomingMessage) -> None |
on_config | Configuration snapshot received | (config: ConfigSnapshot) -> None |
on_signal | Signal received | (signal: Signal) -> None |
on_error | Error occurred | (error: ErrorResponse) -> None |
on_kv_response | Async KV operation response | (kv: KVResponse) -> None |
on_task_assignment | Task assigned (Orchestrators) | (assignment: TaskAssignment) -> None |
on_checkpoint_response | Async checkpoint response | (response: CheckpointResponse) -> None |
on_connect | Connection established | () -> None |
on_disconnect | Connection lost | (reason: str) -> None |
Typed handlers (on_chat_message, etc.) and the catch-all on_message are independent. If both are registered, the typed handler fires first, then on_message fires as well. This matches the behavior of the Go and TypeScript SDKs.
For async clients, callbacks can be either sync or async functions:
# Sync callbackdef on_message(msg): print(msg.payload.decode())
# Async callbackasync def on_message(msg): await process_message(msg)Messaging
Section titled “Messaging”Message Types
Section titled “Message Types”from scitrera_aether_client import OPAQUE, CHAT, CONTROL, TOOL_CALL, EVENT, METRIC
# OPAQUE - Default type; forwarded verbatim (callers own the payload schema)client.send_message_to_agent(..., message_type=OPAQUE)
# CHAT - Conversational text messagesclient.send_message_to_agent(..., message_type=CHAT)
# CONTROL - Control/command messagesclient.send_message_to_agent(..., message_type=CONTROL)
# EVENT - Events for workflow engineclient.send_event(payload)
# METRIC - Structured telemetry for metrics bridge — use new_metric() builderfrom scitrera_aether_client import new_metricmetric = new_metric().add("latency", "gauge", 42.0).build()client.send_metric(metric)Sending Messages
Section titled “Sending Messages”# To a specific agentclient.send_message_to_agent( workspace="default", implementation="worker", specifier="01", payload=b"Hello!")
# To a taskclient.send_message_to_task( workspace="default", implementation="processor", payload=b"Process this", unique_specifier="task-123" # Optional for unique tasks)
# To a user sessionclient.send_message_to_user_session( user_id="user-123", window_id="window-abc", payload=b"Notification")
# Broadcast to all agents in workspaceclient.send_broadcast_to_agents( workspace="default", payload=b"Broadcast message")
# Send event (agents/tasks only)client.send_event(b'{"event": "completed"}')
# Send metric (agents/tasks only) — must pass a Metric protobuf, not raw bytesfrom scitrera_aether_client import new_metricmetric = new_metric().add("latency", "gauge", 42.0).build()client.send_metric(metric)KV Operations
Section titled “KV Operations”The KV store supports multiple scopes:
| Scope | Description | Required Parameters |
|---|---|---|
global | Global configuration | None |
workspace | Workspace-specific | workspace |
user | User-specific | user_id |
user-workspace | User + workspace scoped | user_id, workspace |
Synchronous Client KV
Section titled “Synchronous Client KV”# Store a value (fire-and-forget)client.kv_put( key="config/setting", value=b"value", scope="global")
# Store with workspace scopeclient.kv_put( key="team/setting", value=b"team-value", scope="workspace", workspace="default")
# Get a value (response via callback)client.kv_get(key="config/setting", scope="global")
# List keysclient.kv_list(key_prefix="config/", scope="global")
# Delete a keyclient.kv_delete(key="config/old", scope="global")Async Client KV
Section titled “Async Client KV”Async clients support both fire-and-forget (_nowait) and awaitable operations:
# Fire-and-forgetawait client.kv_put_nowait( key="setting", value=b"value", scope="global")
# Await responseresponse = await client.kv_get( key="setting", scope="global", timeout=5.0)if response: print(f"Value: {response.value}")
# Put and await confirmationresponse = await client.kv_put( key="setting", value=b"new-value", scope="global", timeout=5.0)Task Creation
Section titled “Task Creation”Create tasks with different assignment modes:
from scitrera_aether_client import SELF_ASSIGN, TARGETED, POOL
# Self-assigned task (creator handles it)client.create_task( task_type="data-processing", workspace="default", assignment_mode=SELF_ASSIGN, metadata={"priority": "high"})
# Targeted task (assign to specific agent)client.create_task( task_type="specialized-work", workspace="default", assignment_mode=TARGETED, target_agent_id="ag.default.worker.specialist-01", launch_param_overrides={"memory": "4G"})
# Pool task (load-balanced to available workers)client.create_task( task_type="batch-job", workspace="default", assignment_mode=POOL)Checkpointing
Section titled “Checkpointing”Save and restore agent/task state:
Synchronous
Section titled “Synchronous”# Save checkpointclient.checkpoint_save(data=b"state data", key="my-checkpoint")
# Save and wait for confirmationresponse = client.checkpoint_save_sync( data=b"state data", key="my-checkpoint", timeout=5.0)
# Load checkpointresponse = client.checkpoint_load_sync(key="my-checkpoint", timeout=5.0)if response and response.data: print(f"Restored state: {response.data}")
# List checkpointsresponse = client.checkpoint_list_sync(timeout=5.0)if response: print(f"Checkpoints: {response.keys}")
# Delete checkpointclient.checkpoint_delete_sync(key="my-checkpoint", timeout=5.0)Asynchronous
Section titled “Asynchronous”# Save and waitresponse = await client.checkpoint_save( data=b"state data", key="my-checkpoint", timeout=5.0)
# Loadresponse = await client.checkpoint_load(key="my-checkpoint", timeout=5.0)
# Fire-and-forget operationsawait client.checkpoint_save_nowait(data=b"state", key="quick-save")await client.checkpoint_delete_nowait(key="old-checkpoint")Connection Configuration
Section titled “Connection Configuration”All clients support configurable connection behavior:
client = AgentClient( workspace="default", implementation="worker", specifier="01", # Retry configuration max_retries=5, # Max connection attempts (0 = infinite) initial_backoff=1.0, # Initial retry delay in seconds max_backoff=30.0, # Maximum retry delay auto_reconnect=True # Auto-reconnect on connection loss)Reconnection Behavior
Section titled “Reconnection Behavior”- On connection loss, clients automatically attempt to reconnect (if
auto_reconnect=True) - Exponential backoff with jitter prevents thundering herd
- Non-recoverable errors (authentication failures, etc.) stop reconnection attempts
- Session IDs are preserved for session resumption when possible
Workspace Switching
Section titled “Workspace Switching”Agents and tasks can switch workspaces:
client.switch_workspace("new-workspace")TLS Configuration
Section titled “TLS Configuration”# Simple TLS (server authentication using system CA)client = AgentClient( workspace="default", implementation="worker", specifier="w1", tls_enabled=True,)
# Custom CA certificateclient = AgentClient( workspace="default", implementation="worker", specifier="w1", tls_enabled=True, tls_root_cert_path="/path/to/ca.pem",)
# mTLS (mutual authentication)client = AgentClient( workspace="default", implementation="worker", specifier="w1", tls_enabled=True, tls_root_cert_path="/path/to/ca.pem", tls_client_cert_path="/path/to/client.pem", tls_client_key_path="/path/to/client-key.pem",)
# Pass certificate bytes directly instead of pathsclient = AgentClient( workspace="default", implementation="worker", specifier="w1", tls_enabled=True, tls_root_cert=ca_cert_bytes, tls_client_cert=client_cert_bytes, tls_client_key=client_key_bytes,)Constants
Section titled “Constants”from scitrera_aether_client import ( # Message types OPAQUE, # Default — forwarded verbatim by Aether CHAT, # Conversational text messages CONTROL, # Control/command messages TOOL_CALL, # Tool invocations EVENT, # Events for workflow engine METRIC, # Telemetry data
# Task assignment modes SELF_ASSIGN, # Creator handles the task TARGETED, # Assign to specific agent POOL, # Load-balanced assignment
# KV operations KV_GET, KV_PUT, KV_LIST, KV_DELETE,
# KV scopes KV_SCOPE_GLOBAL, KV_SCOPE_WORKSPACE, KV_SCOPE_USER, KV_SCOPE_USER_WORKSPACE, KV_SCOPE_GLOBAL_EXCLUSIVE, # Per-agent, tenant-wide KV_SCOPE_WORKSPACE_EXCLUSIVE, # Per-agent, per-workspace KV_SCOPE_USER_SHARED, # Shared across agents, per-user KV_SCOPE_USER_WORKSPACE_SHARED, # Shared across agents, per-user+workspace)Examples
Section titled “Examples”See the example.py and example_async.py files for comprehensive examples including:
- Agent client with messaging, KV operations, and task creation
- Orchestrator for managing agent lifecycle
- Workflow engine for event processing
- Metrics bridge for telemetry collection
- Concurrent async clients
- Context manager usage
Run examples:
# Sync examplespython example.py agentpython example.py orchestratorpython example.py workflowpython example.py metrics
# Async examplespython example_async.py agentpython example_async.py concurrentpython example_async.py contextKey Architectural Principle
Section titled “Key Architectural Principle”Connection = Lock = Heartbeat: The active gRPC stream IS the distributed lock AND the liveness proof for the connected identity. No separate heartbeat API exists. When the stream closes, the lock is released automatically.
Requirements
Section titled “Requirements”- Python 3.11+
- grpcio >= 1.76.0
- protobuf >= 5.29.0
See Also
Section titled “See Also”- Quick Start — connect your first agent
- Go SDK — full-featured Go client
- TypeScript SDK — browser and Node.js clients