Go SDK
Go client SDK for the Scitrera 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”- Idiomatic Go API: Context-based cancellation, functional options, and error wrapping
- 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)
- Checkpointing: Persist and restore agent/task state across restarts
- Auto-Reconnection: Configurable exponential backoff with automatic reconnection
- Callback Handlers: Event-driven message handling with context support
- TLS/mTLS Support: Secure connections with optional mutual TLS authentication
- Concurrency Safe: All clients are safe for concurrent use
Installation
Section titled “Installation”go get github.com/scitrera/aether/sdk/goQuick Start
Section titled “Quick Start”Agent Client
Section titled “Agent Client”package main
import ( "context" "fmt" "log" "os" "os/signal"
"github.com/scitrera/aether/sdk/go/aether")
func main() { // Create an agent client client, err := aether.NewAgentClient(aether.AgentOptions{ ClientOptions: aether.ClientOptions{ ServerAddr: "localhost:50051", }, Workspace: "default", Implementation: "my-agent", Specifier: "agent-01", }) if err != nil { log.Fatal(err) }
// Set up message handler client.OnMessage(func(ctx context.Context, msg *aether.Message) error { fmt.Printf("Received from %s: %s\n", msg.SourceTopic, msg.Payload) return nil })
// Set up connection handlers client.OnConnect(func(ctx context.Context, ack *aether.ConnectionAck) error { log.Printf("Connected with session %s", ack.SessionID) return nil })
client.OnDisconnect(func(ctx context.Context, reason string) error { log.Printf("Disconnected: %s", reason) return nil })
// Create cancellable context ctx, cancel := context.WithCancel(context.Background()) defer cancel()
// Handle interrupt signal go func() { sigCh := make(chan os.Signal, 1) signal.Notify(sigCh, os.Interrupt) <-sigCh cancel() }()
// Connect to the gateway if err := client.Connect(ctx); err != nil { log.Fatal(err) } defer client.Close()
// Send a message to another agent err = client.SendToAgent("default", "other-agent", "01", []byte("Hello!")) if err != nil { log.Printf("Failed to send: %v", err) }
// Run the message loop (blocks until disconnect) if err := client.Run(ctx); err != nil { log.Printf("Client stopped: %v", err) }}Task Client
Section titled “Task Client”package main
import ( "context" "log"
"github.com/scitrera/aether/sdk/go/aether")
func main() { // Unique task (named) uniqueTask, err := aether.NewTaskClient(aether.TaskOptions{ ClientOptions: aether.ClientOptions{ ServerAddr: "localhost:50051", }, Workspace: "default", Implementation: "data-processor", Specifier: "job-123", // Unique task }) if err != nil { log.Fatal(err) }
// Non-unique task (pooled, server assigns ID) pooledTask, err := aether.NewTaskClient(aether.TaskOptions{ ClientOptions: aether.ClientOptions{ ServerAddr: "localhost:50051", }, Workspace: "default", Implementation: "worker", // Specifier omitted - server assigns ID }) if err != nil { log.Fatal(err) }
// Non-unique tasks subscribe to both their specific topic // and the broadcast topic for work claiming log.Printf("Pooled task broadcast topic: %s", pooledTask.BroadcastTopic())
_ = uniqueTask}Client Types
Section titled “Client Types”AgentClient
Section titled “AgentClient”For long-running agent processes that need unique identities.
client, err := aether.NewAgentClient(aether.AgentOptions{ ClientOptions: aether.ClientOptions{ ServerAddr: "localhost:50051", }, Workspace: "default", Implementation: "python-worker", Specifier: "worker-01",})TaskClient
Section titled “TaskClient”For task execution. Supports both unique (named) and non-unique (pooled) tasks.
// Unique taskuniqueTask, _ := aether.NewTaskClient(aether.TaskOptions{ ClientOptions: aether.ClientOptions{ServerAddr: "localhost:50051"}, Workspace: "default", Implementation: "processor", Specifier: "task-123", // Named task})
// Non-unique task (server assigns ID)pooledTask, _ := aether.NewTaskClient(aether.TaskOptions{ ClientOptions: aether.ClientOptions{ServerAddr: "localhost:50051"}, Workspace: "default", Implementation: "worker", // Specifier omitted})UserClient
Section titled “UserClient”For user session connections (e.g., from browser windows).
client, err := aether.NewUserClient(aether.UserOptions{ ClientOptions: aether.ClientOptions{ ServerAddr: "localhost:50051", }, UserID: "user-123", WindowID: "window-abc",})OrchestratorClient
Section titled “OrchestratorClient”For managing agent/task lifecycle and compute resources.
client, err := aether.NewOrchestratorClient(aether.OrchestratorOptions{ ClientOptions: aether.ClientOptions{ ServerAddr: "localhost:50051", }, Implementation: "kubernetes-orchestrator", SupportedProfiles: []string{"docker", "kubernetes"},})WorkflowEngineClient
Section titled “WorkflowEngineClient”For processing events and triggering downstream actions.
client, err := aether.NewWorkflowEngineClient(aether.WorkflowEngineOptions{ ClientOptions: aether.ClientOptions{ ServerAddr: "localhost:50051", },})MetricsBridgeClient
Section titled “MetricsBridgeClient”For collecting telemetry data from agents and tasks.
client, err := aether.NewMetricsBridgeClient(aether.MetricsBridgeOptions{ ClientOptions: aether.ClientOptions{ ServerAddr: "localhost:50051", },})Handlers
Section titled “Handlers”All clients support the following handlers:
| Handler | Description | Signature |
|---|---|---|
OnMessage | Incoming message received | func(ctx context.Context, msg *Message) error |
OnConfig | Configuration snapshot received | func(ctx context.Context, config *ConfigSnapshot) error |
OnSignal | Signal received | func(ctx context.Context, signal *Signal) error |
OnError | Error occurred | func(ctx context.Context, err *ErrorInfo) error |
OnKVResponse | KV operation response | func(ctx context.Context, resp *KVResponse) error |
OnCheckpointResponse | Checkpoint operation response | func(ctx context.Context, resp *CheckpointResponse) error |
OnTaskAssignment | Task assigned (Orchestrators) | func(ctx context.Context, task *TaskAssignment) error |
OnConnect | Connection established | func(ctx context.Context, ack *ConnectionAck) error |
OnDisconnect | Connection lost | func(ctx context.Context, reason string) error |
OnReconnecting | Reconnection attempt | func(ctx context.Context, attempt int) error |
Example handler registration:
client.OnMessage(func(ctx context.Context, msg *aether.Message) error { // Process message fmt.Printf("From %s: %s\n", msg.SourceTopic, msg.Payload) return nil})
client.OnConnect(func(ctx context.Context, ack *aether.ConnectionAck) error { log.Printf("Connected, session: %s, resumed: %v", ack.SessionID, ack.Resumed) return nil})
client.OnReconnecting(func(ctx context.Context, attempt int) error { log.Printf("Reconnection attempt %d", attempt) if attempt > 10 { return errors.New("too many reconnection attempts") } return nil})Messaging
Section titled “Messaging”Message Types
Section titled “Message Types”import pb "github.com/scitrera/aether/api/proto"
// OPAQUE - Default type; forwarded verbatim by Aether (callers own the payload schema)client.SendToAgent("ws", "impl", "spec", payload)
// CONTROL - Control/command messagesclient.SendToAgentWithType("ws", "impl", "spec", payload, pb.MessageType_CONTROL)
// TOOL_CALL - Tool invocation messagesclient.SendToolCallMessage(topic, payload)
// EVENT - Events for workflow engine (agents/tasks only)client.SendEvent(payload)
// METRIC - Structured telemetry (agents/tasks only) — build via NewMetric()metric := aether.NewMetric().Add("cpu", "gauge", 0.75).Tag("host", "worker-01").Build()client.SendMetric(metric)Sending Messages
Section titled “Sending Messages”// To a specific agentclient.SendToAgent("default", "worker", "01", []byte("Hello!"))
// To a task (unique or broadcast)client.SendToTask("default", "processor", "task-123", []byte("Process this"))client.SendToTask("default", "worker", "", []byte("Broadcast to pool"))
// To a user sessionclient.SendToUser("user-123", "window-abc", []byte("Notification"))
// To a user's workspace scopeclient.SendToUserWorkspace("user-123", "default", []byte("Workspace message"))
// Broadcast to all agents in workspaceclient.BroadcastToAgents("default", []byte("Broadcast message"))
// Broadcast to all users in workspaceclient.BroadcastToUsers("default", []byte("User broadcast"))
// Send event (agents/tasks only)client.SendEvent([]byte(`{"event": "completed"}`))
// Send metric (agents/tasks only) — build a Metric protobuf, not raw bytesmetric := aether.NewMetric().Add("latency", "gauge", 42).Build()client.SendMetric(metric)KV Operations
Section titled “KV Operations”The KV store supports multiple scopes:
| Scope | Description | Constants |
|---|---|---|
global | Global configuration | aether.KVScopeGlobal |
workspace | Workspace-specific | aether.KVScopeWorkspace |
user | User-specific | aether.KVScopeUser |
user-workspace | User + workspace scoped | aether.KVScopeUserWorkspace |
global-exclusive | Per-agent, tenant-wide | aether.KVScopeGlobalExclusive |
workspace-exclusive | Per-agent, per-workspace | aether.KVScopeWorkspaceExclusive |
user-shared | Shared across agents, per-user | aether.KVScopeUserShared |
user-workspace-shared | Shared across agents, per-user+workspace | aether.KVScopeUserWorkspaceShared |
Asynchronous Operations
Section titled “Asynchronous Operations”kv := client.KV()
// Store a value (fire-and-forget)kv.Put("config/setting", []byte("value"), aether.KVScopeGlobal, "", "", 0)
// Store with workspace scopekv.PutWorkspace("team/setting", []byte("team-value"), "default")
// Get a value (response via OnKVResponse handler)kv.Get("config/setting", aether.KVScopeGlobal, "", "")
// List keyskv.ListGlobal("config/")
// Delete a keykv.DeleteGlobal("config/old")Synchronous Operations
Section titled “Synchronous Operations”ctx := context.Background()kv := client.KV()
// Get with timeoutresp, err := kv.GetSync(ctx, aether.KVGetOptions{ Key: "config/setting", Scope: aether.KVScopeGlobal, Timeout: 5 * time.Second,})if err != nil { log.Fatal(err)}fmt.Printf("Value: %s\n", resp.Value)
// Put with TTLresp, err = kv.PutSync(ctx, aether.KVPutOptions{ Key: "session/data", Value: []byte("session-value"), Scope: aether.KVScopeWorkspace, Workspace: "default", TTL: time.Hour, Timeout: 5 * time.Second,})
// List keyslistResp, err := kv.ListSync(ctx, aether.KVListOptions{ KeyPrefix: "config/", Scope: aether.KVScopeGlobal, Timeout: 5 * time.Second,})if err == nil { for _, key := range listResp.Keys { fmt.Println(key) }}Convenience Methods
Section titled “Convenience Methods”// Global scope shortcutskv.GetGlobal("key")kv.PutGlobal("key", []byte("value"))kv.DeleteGlobal("key")kv.ListGlobal("prefix/")
// Workspace scope shortcutskv.GetWorkspace("key", "my-workspace")kv.PutWorkspace("key", []byte("value"), "my-workspace")
// User scope shortcutskv.GetUser("key", "user-123")kv.PutUser("key", []byte("value"), "user-123")
// User-workspace scope shortcutskv.GetUserWorkspace("key", "user-123", "my-workspace")kv.PutUserWorkspace("key", []byte("value"), "user-123", "my-workspace")Checkpointing
Section titled “Checkpointing”Save and restore agent/task state:
Asynchronous Operations
Section titled “Asynchronous Operations”cp := client.Checkpoint()
// Save checkpoint (fire-and-forget)cp.Save([]byte("state data"), "my-checkpoint", -1) // -1 = server default TTL
// Save with specific TTLcp.SaveWithTTL([]byte("state"), "checkpoint-key", time.Hour)
// Save with no expirationcp.SavePermanent([]byte("state"), "permanent-checkpoint")
// Load checkpoint (response via OnCheckpointResponse handler)cp.Load("my-checkpoint")
// List checkpointscp.List()
// Delete checkpointcp.Delete("my-checkpoint")Synchronous Operations
Section titled “Synchronous Operations”ctx := context.Background()cp := client.Checkpoint()
// Save and wait for confirmationresp, err := cp.SaveSync(ctx, aether.CheckpointSaveOptions{ Data: []byte("state data"), Key: "my-checkpoint", TTL: time.Hour, // 0 = no expiration, -1 = server default Timeout: 5 * time.Second,})
// Load checkpointresp, err = cp.LoadSync(ctx, aether.CheckpointLoadOptions{ Key: "my-checkpoint", Timeout: 5 * time.Second,})if err == nil && resp.Success { fmt.Printf("Restored state: %s\n", resp.Data)}
// List all checkpointslistResp, err := cp.ListSync(ctx, 5*time.Second)if err == nil { fmt.Printf("Checkpoints: %v\n", listResp.Keys)}
// Delete checkpoint_, err = cp.DeleteSync(ctx, aether.CheckpointDeleteOptions{ Key: "my-checkpoint", Timeout: 5 * time.Second,})Task Creation
Section titled “Task Creation”Create tasks with different assignment modes:
// Self-assigned task (creator handles it)client.CreateTask(aether.CreateTaskOptions{ TaskType: "data-processing", Workspace: "default", AssignmentMode: aether.TaskAssignmentSelfAssign, Metadata: map[string]string{"priority": "high"},})
// Targeted task (assign to specific agent)client.CreateTask(aether.CreateTaskOptions{ TaskType: "specialized-work", Workspace: "default", AssignmentMode: aether.TaskAssignmentTargeted, TargetAgentID: "ag.default.worker.specialist-01", LaunchParamOverrides: map[string]string{"memory": "4G"},})
// Pool task (load-balanced to available workers)client.CreateTask(aether.CreateTaskOptions{ TaskType: "batch-job", Workspace: "default", AssignmentMode: aether.TaskAssignmentPool,})Connection Configuration
Section titled “Connection Configuration”All clients support configurable connection behavior:
client, err := aether.NewAgentClient(aether.AgentOptions{ ClientOptions: aether.ClientOptions{ ServerAddr: "localhost:50051", Connection: aether.ConnectionOptions{ MaxRetries: 5, // Max connection attempts (0 = infinite) InitialBackoff: time.Second, // Initial retry delay MaxBackoff: 30 * time.Second, // Maximum retry delay BackoffMultiplier: 2.0, // Exponential backoff multiplier AutoReconnect: true, // Auto-reconnect on connection loss ConnectTimeout: 30 * time.Second, // Connection timeout KeepAliveInterval: 30 * time.Second, // Keepalive ping interval }, }, Workspace: "default", Implementation: "worker", Specifier: "01",})Functional Options
Section titled “Functional Options”opts := aether.DefaultConnectionOptions()aether.ApplyConnectionOptions(&opts, aether.WithMaxRetries(10), aether.WithInitialBackoff(2*time.Second), aether.WithMaxBackoff(time.Minute), aether.WithAutoReconnect(true),)Reconnection Behavior
Section titled “Reconnection Behavior”- On connection loss, clients automatically attempt to reconnect (if
AutoReconnect=true) - Exponential backoff with jitter prevents thundering herd
- Non-recoverable errors (authentication failures, duplicate identity, etc.) stop reconnection attempts
- Session IDs are preserved for session resumption when possible
WithRetryOnDuplicate
Section titled “WithRetryOnDuplicate”By default, a DuplicateIdentityError is treated as non-recoverable and stops reconnection. Enable RetryOnDuplicate to treat it as recoverable instead — useful when restarting a container or process before the previous connection’s Redis lock (30s TTL) has expired:
opts := aether.DefaultConnectionOptions()aether.ApplyConnectionOptions(&opts, aether.WithRetryOnDuplicate(true), // retry until lock expires (~30 s) aether.WithMaxRetries(0), // infinite retries aether.WithMaxBackoff(10*time.Second),)
client, err := aether.NewAgentClient(aether.AgentOptions{ ClientOptions: aether.ClientOptions{ ServerAddr: "localhost:50051", Connection: opts, }, Workspace: "default", Implementation: "my-agent", Specifier: "agent-01",})Or set it directly in the struct:
Connection: aether.ConnectionOptions{ AutoReconnect: true, RetryOnDuplicate: true, MaxRetries: 0, // 0 = infinite InitialBackoff: time.Second, MaxBackoff: 10 * time.Second, BackoffMultiplier: 2.0,},TLS Configuration
Section titled “TLS Configuration”Simple TLS (Server Authentication)
Section titled “Simple TLS (Server Authentication)”client, err := aether.NewAgentClient(aether.AgentOptions{ ClientOptions: aether.ClientOptions{ ServerAddr: "secure.example.com:50051", TLS: &aether.TLSConfig{ Enabled: true, ServerName: "secure.example.com", }, }, // ...})mTLS (Mutual Authentication)
Section titled “mTLS (Mutual Authentication)”// Load certificates from filestlsConfig, err := aether.LoadTLSConfigFromFiles( "/path/to/ca.pem", "/path/to/client.pem", "/path/to/client-key.pem",)if err != nil { log.Fatal(err)}
client, err := aether.NewAgentClient(aether.AgentOptions{ ClientOptions: aether.ClientOptions{ ServerAddr: "secure.example.com:50051", TLS: tlsConfig, }, // ...})TLS Config from Bytes
Section titled “TLS Config from Bytes”client, err := aether.NewAgentClient(aether.AgentOptions{ ClientOptions: aether.ClientOptions{ ServerAddr: "secure.example.com:50051", TLS: &aether.TLSConfig{ Enabled: true, RootCAs: caCertPEM, // []byte ClientCert: clientCertPEM, // []byte ClientKey: clientKeyPEM, // []byte ServerName: "secure.example.com", }, }, // ...})Credentials
Section titled “Credentials”creds := aether.NewCredentials(). WithAPIKey("your-api-key"). WithTenant("tenant-id")
client, err := aether.NewAgentClient(aether.AgentOptions{ ClientOptions: aether.ClientOptions{ ServerAddr: "localhost:50051", Credentials: creds, }, // ...})Workspace Switching
Section titled “Workspace Switching”Agents and tasks can switch workspaces:
err := client.SwitchWorkspace("new-workspace")Error Types
Section titled “Error Types”The SDK provides typed errors for specific error conditions:
import "github.com/scitrera/aether/sdk/go/aether"
// Connection errorsvar connErr *aether.ConnectionErrorvar closedErr *aether.ConnectionClosedErrorvar reconnErr *aether.ReconnectionError
// Auth errorsvar authErr *aether.AuthenticationErrorvar permErr *aether.PermissionDeniedError
// Identity errorsvar dupErr *aether.DuplicateIdentityError
// Timeout errorsvar timeoutErr *aether.TimeoutError
// Request errorsvar argErr *aether.InvalidArgumentErrorvar notFoundErr *aether.NotFoundError
// Example error handlingif err := client.Connect(ctx); err != nil { if errors.As(err, &dupErr) { log.Printf("Identity already connected: %s", dupErr.Identity) } else if errors.As(err, &authErr) { log.Printf("Authentication failed: %s", authErr.Message) } else if aether.IsRecoverable(err) { log.Printf("Recoverable error, will retry: %v", err) } else { log.Fatal(err) }}Error Classification Helpers
Section titled “Error Classification Helpers”// Check if error is recoverable (should trigger reconnection)if aether.IsRecoverable(err) { // Client will auto-reconnect}
// Check if error is connection-relatedif aether.IsConnectionError(err) { // Handle connection issues}
// Check if error is a timeoutif aether.IsTimeoutError(err) { // Handle timeout}Topic Helpers
Section titled “Topic Helpers”The SDK provides helpers for constructing topic addresses:
// Agent topicstopic := aether.AgentTopic("workspace", "impl", "spec")// Result: "ag.workspace.impl.spec"
// Task topicsuniqueTopic := aether.UniqueTaskTopic("workspace", "impl", "spec")// Result: "tu.workspace.impl.spec"
nonUniqueTopic := aether.TaskTopic("workspace", "impl", "id")// Result: "ta.workspace.impl.id"
broadcastTopic := aether.TaskBroadcastTopic("workspace", "impl")// Result: "tb.workspace.impl"
// User topicsuserTopic := aether.UserTopic("user-id", "window-id")// Result: "us.user-id.window-id"
userWsTopic := aether.UserWorkspaceTopic("user-id", "workspace")// Result: "uw.user-id.workspace"
// Broadcast topicsagentBroadcast := aether.GlobalAgentsTopic("workspace")// Result: "ga.workspace"
userBroadcast := aether.GlobalUsersTopic("workspace")// Result: "gu.workspace"
// Event and metric topicseventTopic := aether.EventWildcardTopic()// Result: "event.*"
metricTopic := aether.MetricWildcardTopic()// Result: "metric.*"Constants
Section titled “Constants”import "github.com/scitrera/aether/sdk/go/aether"
// Message typesaether.MessageTypeOpaque // "OPAQUE" — default for generic send helpersaether.MessageTypeChat // "CHAT"aether.MessageTypeControl // "CONTROL"aether.MessageTypeToolCall // "TOOL_CALL"aether.MessageTypeEvent // "EVENT"aether.MessageTypeMetric // "METRIC"
// Task assignment modesaether.TaskAssignmentSelfAssign // "SELF_ASSIGN"aether.TaskAssignmentTargeted // "TARGETED"aether.TaskAssignmentPool // "POOL"
// KV scopesaether.KVScopeGlobal // "global"aether.KVScopeWorkspace // "workspace"aether.KVScopeUser // "user"aether.KVScopeUserWorkspace // "user-workspace"aether.KVScopeGlobalExclusive // "global-exclusive"aether.KVScopeWorkspaceExclusive // "workspace-exclusive"aether.KVScopeUserShared // "user-shared"aether.KVScopeUserWorkspaceShared // "user-workspace-shared"
// Signal typesaether.SignalForceDisconnectKey Architectural Principle
Section titled “Key Architectural Principle”Connection = Lock = Heartbeat: The connection itself IS the distributed lock AND the heartbeat. When the gRPC stream closes, the identity lock is immediately released. No separate heartbeat API exists.
Requirements
Section titled “Requirements”- Go 1.25+
- gRPC 1.78.0+
- Protobuf 1.36.0+
See Also
Section titled “See Also”- Quick Start — connect your first agent
- Python SDK — sync and async Python clients
- TypeScript SDK — browser and Node.js clients