Skip to content

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.

  • 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. alongside on_message — both fire for matching messages
Terminal window
pip install scitrera-aether-client

For development:

Terminal window
pip install scitrera-aether-client[dev]
from scitrera_aether_client import AgentClient, CHAT
# Create an agent client
client = AgentClient(
workspace="default",
implementation="my-agent",
specifier="agent-01"
)
# Set up message callback
def on_message(msg):
print(f"Received from {msg.source_topic}: {msg.payload.decode()}")
client.on_message = on_message
# Connect to the gateway
client.connect("localhost:50051")
# Send a message to another agent
client.send_message_to_agent(
workspace="default",
implementation="other-agent",
specifier="01",
payload=b"Hello!"
)
# Keep running until interrupted
try:
while True:
import time
time.sleep(1)
except KeyboardInterrupt:
pass
finally:
client.close()
import asyncio
from 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())
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 closed

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"
)

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"
)

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()

All clients support the following callbacks:

CallbackDescriptionSignature
on_messageEvery incoming message (catch-all)(msg: IncomingMessage) -> None
on_chat_messageCHAT-typed messages(msg: IncomingMessage) -> None
on_control_messageCONTROL-typed messages(msg: IncomingMessage) -> None
on_tool_callTOOL_CALL-typed messages(msg: IncomingMessage) -> None
on_eventEVENT-typed messages(msg: IncomingMessage) -> None
on_metricMETRIC-typed messages(msg: IncomingMessage) -> None
on_configConfiguration snapshot received(config: ConfigSnapshot) -> None
on_signalSignal received(signal: Signal) -> None
on_errorError occurred(error: ErrorResponse) -> None
on_kv_responseAsync KV operation response(kv: KVResponse) -> None
on_task_assignmentTask assigned (Orchestrators)(assignment: TaskAssignment) -> None
on_checkpoint_responseAsync checkpoint response(response: CheckpointResponse) -> None
on_connectConnection established() -> None
on_disconnectConnection 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 callback
def on_message(msg):
print(msg.payload.decode())
# Async callback
async def on_message(msg):
await process_message(msg)
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 messages
client.send_message_to_agent(..., message_type=CHAT)
# CONTROL - Control/command messages
client.send_message_to_agent(..., message_type=CONTROL)
# EVENT - Events for workflow engine
client.send_event(payload)
# METRIC - Structured telemetry for metrics bridge — use new_metric() builder
from scitrera_aether_client import new_metric
metric = new_metric().add("latency", "gauge", 42.0).build()
client.send_metric(metric)
# To a specific agent
client.send_message_to_agent(
workspace="default",
implementation="worker",
specifier="01",
payload=b"Hello!"
)
# To a task
client.send_message_to_task(
workspace="default",
implementation="processor",
payload=b"Process this",
unique_specifier="task-123" # Optional for unique tasks
)
# To a user session
client.send_message_to_user_session(
user_id="user-123",
window_id="window-abc",
payload=b"Notification"
)
# Broadcast to all agents in workspace
client.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 bytes
from scitrera_aether_client import new_metric
metric = new_metric().add("latency", "gauge", 42.0).build()
client.send_metric(metric)

The KV store supports multiple scopes:

ScopeDescriptionRequired Parameters
globalGlobal configurationNone
workspaceWorkspace-specificworkspace
userUser-specificuser_id
user-workspaceUser + workspace scopeduser_id, workspace
# Store a value (fire-and-forget)
client.kv_put(
key="config/setting",
value=b"value",
scope="global"
)
# Store with workspace scope
client.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 keys
client.kv_list(key_prefix="config/", scope="global")
# Delete a key
client.kv_delete(key="config/old", scope="global")

Async clients support both fire-and-forget (_nowait) and awaitable operations:

# Fire-and-forget
await client.kv_put_nowait(
key="setting",
value=b"value",
scope="global"
)
# Await response
response = await client.kv_get(
key="setting",
scope="global",
timeout=5.0
)
if response:
print(f"Value: {response.value}")
# Put and await confirmation
response = await client.kv_put(
key="setting",
value=b"new-value",
scope="global",
timeout=5.0
)

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
)

Save and restore agent/task state:

# Save checkpoint
client.checkpoint_save(data=b"state data", key="my-checkpoint")
# Save and wait for confirmation
response = client.checkpoint_save_sync(
data=b"state data",
key="my-checkpoint",
timeout=5.0
)
# Load checkpoint
response = client.checkpoint_load_sync(key="my-checkpoint", timeout=5.0)
if response and response.data:
print(f"Restored state: {response.data}")
# List checkpoints
response = client.checkpoint_list_sync(timeout=5.0)
if response:
print(f"Checkpoints: {response.keys}")
# Delete checkpoint
client.checkpoint_delete_sync(key="my-checkpoint", timeout=5.0)
# Save and wait
response = await client.checkpoint_save(
data=b"state data",
key="my-checkpoint",
timeout=5.0
)
# Load
response = await client.checkpoint_load(key="my-checkpoint", timeout=5.0)
# Fire-and-forget operations
await client.checkpoint_save_nowait(data=b"state", key="quick-save")
await client.checkpoint_delete_nowait(key="old-checkpoint")

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
)
  • 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

Agents and tasks can switch workspaces:

client.switch_workspace("new-workspace")
# Simple TLS (server authentication using system CA)
client = AgentClient(
workspace="default", implementation="worker", specifier="w1",
tls_enabled=True,
)
# Custom CA certificate
client = 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 paths
client = 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,
)
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
)

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:

Terminal window
# Sync examples
python example.py agent
python example.py orchestrator
python example.py workflow
python example.py metrics
# Async examples
python example_async.py agent
python example_async.py concurrent
python example_async.py context

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.

  • Python 3.11+
  • grpcio >= 1.76.0
  • protobuf >= 5.29.0