Skip to content

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.

  • 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
Terminal window
go get github.com/scitrera/aether/sdk/go
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)
}
}
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
}

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

For task execution. Supports both unique (named) and non-unique (pooled) tasks.

// Unique task
uniqueTask, _ := 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
})

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

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

For processing events and triggering downstream actions.

client, err := aether.NewWorkflowEngineClient(aether.WorkflowEngineOptions{
ClientOptions: aether.ClientOptions{
ServerAddr: "localhost:50051",
},
})

For collecting telemetry data from agents and tasks.

client, err := aether.NewMetricsBridgeClient(aether.MetricsBridgeOptions{
ClientOptions: aether.ClientOptions{
ServerAddr: "localhost:50051",
},
})

All clients support the following handlers:

HandlerDescriptionSignature
OnMessageIncoming message receivedfunc(ctx context.Context, msg *Message) error
OnConfigConfiguration snapshot receivedfunc(ctx context.Context, config *ConfigSnapshot) error
OnSignalSignal receivedfunc(ctx context.Context, signal *Signal) error
OnErrorError occurredfunc(ctx context.Context, err *ErrorInfo) error
OnKVResponseKV operation responsefunc(ctx context.Context, resp *KVResponse) error
OnCheckpointResponseCheckpoint operation responsefunc(ctx context.Context, resp *CheckpointResponse) error
OnTaskAssignmentTask assigned (Orchestrators)func(ctx context.Context, task *TaskAssignment) error
OnConnectConnection establishedfunc(ctx context.Context, ack *ConnectionAck) error
OnDisconnectConnection lostfunc(ctx context.Context, reason string) error
OnReconnectingReconnection attemptfunc(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
})
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 messages
client.SendToAgentWithType("ws", "impl", "spec", payload, pb.MessageType_CONTROL)
// TOOL_CALL - Tool invocation messages
client.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)
// To a specific agent
client.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 session
client.SendToUser("user-123", "window-abc", []byte("Notification"))
// To a user's workspace scope
client.SendToUserWorkspace("user-123", "default", []byte("Workspace message"))
// Broadcast to all agents in workspace
client.BroadcastToAgents("default", []byte("Broadcast message"))
// Broadcast to all users in workspace
client.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 bytes
metric := aether.NewMetric().Add("latency", "gauge", 42).Build()
client.SendMetric(metric)

The KV store supports multiple scopes:

ScopeDescriptionConstants
globalGlobal configurationaether.KVScopeGlobal
workspaceWorkspace-specificaether.KVScopeWorkspace
userUser-specificaether.KVScopeUser
user-workspaceUser + workspace scopedaether.KVScopeUserWorkspace
global-exclusivePer-agent, tenant-wideaether.KVScopeGlobalExclusive
workspace-exclusivePer-agent, per-workspaceaether.KVScopeWorkspaceExclusive
user-sharedShared across agents, per-useraether.KVScopeUserShared
user-workspace-sharedShared across agents, per-user+workspaceaether.KVScopeUserWorkspaceShared
kv := client.KV()
// Store a value (fire-and-forget)
kv.Put("config/setting", []byte("value"), aether.KVScopeGlobal, "", "", 0)
// Store with workspace scope
kv.PutWorkspace("team/setting", []byte("team-value"), "default")
// Get a value (response via OnKVResponse handler)
kv.Get("config/setting", aether.KVScopeGlobal, "", "")
// List keys
kv.ListGlobal("config/")
// Delete a key
kv.DeleteGlobal("config/old")
ctx := context.Background()
kv := client.KV()
// Get with timeout
resp, 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 TTL
resp, 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 keys
listResp, 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)
}
}
// Global scope shortcuts
kv.GetGlobal("key")
kv.PutGlobal("key", []byte("value"))
kv.DeleteGlobal("key")
kv.ListGlobal("prefix/")
// Workspace scope shortcuts
kv.GetWorkspace("key", "my-workspace")
kv.PutWorkspace("key", []byte("value"), "my-workspace")
// User scope shortcuts
kv.GetUser("key", "user-123")
kv.PutUser("key", []byte("value"), "user-123")
// User-workspace scope shortcuts
kv.GetUserWorkspace("key", "user-123", "my-workspace")
kv.PutUserWorkspace("key", []byte("value"), "user-123", "my-workspace")

Save and restore agent/task state:

cp := client.Checkpoint()
// Save checkpoint (fire-and-forget)
cp.Save([]byte("state data"), "my-checkpoint", -1) // -1 = server default TTL
// Save with specific TTL
cp.SaveWithTTL([]byte("state"), "checkpoint-key", time.Hour)
// Save with no expiration
cp.SavePermanent([]byte("state"), "permanent-checkpoint")
// Load checkpoint (response via OnCheckpointResponse handler)
cp.Load("my-checkpoint")
// List checkpoints
cp.List()
// Delete checkpoint
cp.Delete("my-checkpoint")
ctx := context.Background()
cp := client.Checkpoint()
// Save and wait for confirmation
resp, 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 checkpoint
resp, 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 checkpoints
listResp, 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,
})

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,
})

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",
})
opts := aether.DefaultConnectionOptions()
aether.ApplyConnectionOptions(&opts,
aether.WithMaxRetries(10),
aether.WithInitialBackoff(2*time.Second),
aether.WithMaxBackoff(time.Minute),
aether.WithAutoReconnect(true),
)
  • 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

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,
},
client, err := aether.NewAgentClient(aether.AgentOptions{
ClientOptions: aether.ClientOptions{
ServerAddr: "secure.example.com:50051",
TLS: &aether.TLSConfig{
Enabled: true,
ServerName: "secure.example.com",
},
},
// ...
})
// Load certificates from files
tlsConfig, 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,
},
// ...
})
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",
},
},
// ...
})
creds := aether.NewCredentials().
WithAPIKey("your-api-key").
WithTenant("tenant-id")
client, err := aether.NewAgentClient(aether.AgentOptions{
ClientOptions: aether.ClientOptions{
ServerAddr: "localhost:50051",
Credentials: creds,
},
// ...
})

Agents and tasks can switch workspaces:

err := client.SwitchWorkspace("new-workspace")

The SDK provides typed errors for specific error conditions:

import "github.com/scitrera/aether/sdk/go/aether"
// Connection errors
var connErr *aether.ConnectionError
var closedErr *aether.ConnectionClosedError
var reconnErr *aether.ReconnectionError
// Auth errors
var authErr *aether.AuthenticationError
var permErr *aether.PermissionDeniedError
// Identity errors
var dupErr *aether.DuplicateIdentityError
// Timeout errors
var timeoutErr *aether.TimeoutError
// Request errors
var argErr *aether.InvalidArgumentError
var notFoundErr *aether.NotFoundError
// Example error handling
if 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)
}
}
// Check if error is recoverable (should trigger reconnection)
if aether.IsRecoverable(err) {
// Client will auto-reconnect
}
// Check if error is connection-related
if aether.IsConnectionError(err) {
// Handle connection issues
}
// Check if error is a timeout
if aether.IsTimeoutError(err) {
// Handle timeout
}

The SDK provides helpers for constructing topic addresses:

// Agent topics
topic := aether.AgentTopic("workspace", "impl", "spec")
// Result: "ag.workspace.impl.spec"
// Task topics
uniqueTopic := 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 topics
userTopic := aether.UserTopic("user-id", "window-id")
// Result: "us.user-id.window-id"
userWsTopic := aether.UserWorkspaceTopic("user-id", "workspace")
// Result: "uw.user-id.workspace"
// Broadcast topics
agentBroadcast := aether.GlobalAgentsTopic("workspace")
// Result: "ga.workspace"
userBroadcast := aether.GlobalUsersTopic("workspace")
// Result: "gu.workspace"
// Event and metric topics
eventTopic := aether.EventWildcardTopic()
// Result: "event.*"
metricTopic := aether.MetricWildcardTopic()
// Result: "metric.*"
import "github.com/scitrera/aether/sdk/go/aether"
// Message types
aether.MessageTypeOpaque // "OPAQUE" — default for generic send helpers
aether.MessageTypeChat // "CHAT"
aether.MessageTypeControl // "CONTROL"
aether.MessageTypeToolCall // "TOOL_CALL"
aether.MessageTypeEvent // "EVENT"
aether.MessageTypeMetric // "METRIC"
// Task assignment modes
aether.TaskAssignmentSelfAssign // "SELF_ASSIGN"
aether.TaskAssignmentTargeted // "TARGETED"
aether.TaskAssignmentPool // "POOL"
// KV scopes
aether.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 types
aether.SignalForceDisconnect

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.

  • Go 1.25+
  • gRPC 1.78.0+
  • Protobuf 1.36.0+