Documentation
¶
Overview ¶
Package agentsdk defines the contract between the Boxy server and agents.
An agent is the communications layer for one or more provider drivers. The server talks to agents — never to drivers directly. Whether the agent is embedded (in-process) or remote (gRPC) is transparent to the server; both implement the same Agent interface.
Lifecycle:
- Agent starts and registers with the server (token-based auth)
- Agent advertises which provider types it supports
- Server routes CRUD requests to agents based on provider type
- Agent dispatches to the appropriate local driver
Index ¶
- func Run(ctx context.Context, dial Dialer, cfg RemoteClientConfig) error
- func RunSession(ctx context.Context, stream boxyagentv1.AgentTransportService_ConnectClient, ...) error
- type Agent
- type AgentInfo
- type AvailabilityReportingAgent
- type AvailabilitySnapshot
- type Dialer
- type DriverSet
- type EmbeddedAgent
- func (a *EmbeddedAgent) Allocate(ctx context.Context, provider providersdk.Type, id string) (map[string]any, error)
- func (a *EmbeddedAgent) Create(ctx context.Context, provider providersdk.Type, cfg any) (*providersdk.Resource, error)
- func (a *EmbeddedAgent) Delete(ctx context.Context, provider providersdk.Type, id string) error
- func (a *EmbeddedAgent) Info() AgentInfo
- func (a *EmbeddedAgent) List(ctx context.Context, provider providersdk.Type) ([]providersdk.ResourceStatus, error)
- func (a *EmbeddedAgent) PersonalizeGuest(ctx context.Context, provider providersdk.Type, id string) (*providersdk.GuestPersonalizationResult, error)
- func (a *EmbeddedAgent) Read(ctx context.Context, provider providersdk.Type, id string) (*providersdk.ResourceStatus, error)
- func (a *EmbeddedAgent) Update(ctx context.Context, provider providersdk.Type, id string, ...) (*providersdk.Result, error)
- func (a *EmbeddedAgent) UpdateStream(ctx context.Context, provider providersdk.Type, id string, ...) (*providersdk.Result, error)
- type GuestPersonalizingAgent
- type LogSink
- type RemoteAgent
- func (a *RemoteAgent) Allocate(ctx context.Context, provider providersdk.Type, id string) (map[string]any, error)
- func (a *RemoteAgent) Availability() (AvailabilitySnapshot, bool)
- func (a *RemoteAgent) Close()
- func (a *RemoteAgent) Create(ctx context.Context, provider providersdk.Type, cfg any) (*providersdk.Resource, error)
- func (a *RemoteAgent) Delete(ctx context.Context, provider providersdk.Type, id string) error
- func (a *RemoteAgent) HasHeartbeat() bool
- func (a *RemoteAgent) Info() AgentInfo
- func (a *RemoteAgent) LastSeen() time.Time
- func (a *RemoteAgent) List(ctx context.Context, provider providersdk.Type) ([]providersdk.ResourceStatus, error)
- func (a *RemoteAgent) PersonalizeGuest(ctx context.Context, provider providersdk.Type, id string) (*providersdk.GuestPersonalizationResult, error)
- func (a *RemoteAgent) Read(ctx context.Context, provider providersdk.Type, id string) (*providersdk.ResourceStatus, error)
- func (a *RemoteAgent) RequestLogs(ctx context.Context, since time.Time, limit int) (string, error)
- func (a *RemoteAgent) Serve() error
- func (a *RemoteAgent) SetLogSink(sink LogSink)
- func (a *RemoteAgent) Update(ctx context.Context, provider providersdk.Type, id string, ...) (*providersdk.Result, error)
- func (a *RemoteAgent) UpdateStream(ctx context.Context, provider providersdk.Type, id string, ...) (*providersdk.Result, error)
- type RemoteClientConfig
- type ResourceListingAgent
- type StreamingAgent
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func Run ¶
func Run(ctx context.Context, dial Dialer, cfg RemoteClientConfig) error
Run dials, registers, and serves indefinitely, reconnecting with capped exponential backoff (10s base, doubling, capped at 5 minutes — the same shape as internal/pool/manager.go's provisionBackoffState) whenever a session ends for any reason other than ctx being done. Only the first attempt uses cfg.Token; every reconnect after a successful registration clears it, since the agent's identity is carried by its TLS client certificate from that point on.
func RunSession ¶
func RunSession(ctx context.Context, stream boxyagentv1.AgentTransportService_ConnectClient, cfg RemoteClientConfig) error
RunSession drives one already-open stream to completion: sends the initial RegisterRequest, then runs a heartbeat sender and a command-dispatch receiver concurrently until the stream ends or ctx is cancelled. Returns the first error from either.
Types ¶
type Agent ¶
type Agent interface {
// Info returns the agent's identity and the providers it supports.
Info() AgentInfo
// Create provisions a resource through the named provider.
Create(ctx context.Context, provider providersdk.Type, cfg any) (*providersdk.Resource, error)
// Read returns the current status of a resource.
Read(ctx context.Context, provider providersdk.Type, id string) (*providersdk.ResourceStatus, error)
// Update performs an operation on an existing resource.
Update(ctx context.Context, provider providersdk.Type, id string, op providersdk.Operation) (*providersdk.Result, error)
// Delete destroys a resource. It follows the providersdk.Driver Delete
// contract: deleting an already-missing provider resource is successful.
Delete(ctx context.Context, provider providersdk.Type, id string) error
// Allocate runs allocation-time hooks on an existing resource and returns
// additional Properties to merge. Returns nil, nil if the provider has no
// allocation work to perform.
Allocate(ctx context.Context, provider providersdk.Type, id string) (map[string]any, error)
}
Agent is the interface the server uses to communicate with any agent, whether embedded or remote. It wraps one or more provider drivers and routes CRUD operations to them.
type AgentInfo ¶
type AgentInfo struct {
// ID is a unique identifier for this agent instance.
ID string
// Name is a human-readable label (e.g. "docker-host-1", "lab-hypervisor").
Name string
// Providers lists the provider types this agent can handle.
Providers []providersdk.Type
}
AgentInfo describes an agent and the providers it hosts.
type AvailabilityReportingAgent ¶ added in v0.1.40
type AvailabilityReportingAgent interface {
// Availability returns the latest snapshot and whether one has been
// received yet. Each heartbeat wholly replaces the previous snapshot,
// including replacing it with an empty one — a reporter that starts
// erroring is reflected as missing data on the next heartbeat, not
// stale-but-plausible leftover numbers.
Availability() (AvailabilitySnapshot, bool)
}
AvailabilityReportingAgent is an optional agent capability for exposing the latest AvailabilitySnapshot received over an agent's heartbeat stream. Only RemoteAgent implements this today: EmbeddedAgent has no heartbeat to carry a snapshot on, and querying its local drivers live is a different (ctx-bound, error-returning) operation a future caller can add separately if it needs one — see #178 and #179. Callers must type-assert for this capability rather than relying on it being part of Agent.
type AvailabilitySnapshot ¶ added in v0.1.40
type AvailabilitySnapshot struct {
Data map[providersdk.Type]providersdk.ResourceAvailability
At time.Time
}
AvailabilitySnapshot is one agent's most recently received per-provider providersdk.AvailabilityReporter sample, plus when the server received it. At is stamped by the server on receipt, not taken from the agent's self-reported clock — the same trust boundary that keeps liveness keyed to the authenticated connection rather than a claimed value in the message.
A provider type absent from Data means "no reporter, a reporter error, or a sampling timeout" on the most recent heartbeat — never "zero availability". Callers must not conflate the two.
type Dialer ¶
type Dialer func(ctx context.Context) (boxyagentv1.AgentTransportService_ConnectClient, error)
Dialer opens one new AgentTransportService.Connect stream. Supplied by the caller (internal/cli's `boxy agent serve`, see Phase 5/6a) so this package stays transport/TLS-setup agnostic and independently testable.
type DriverSet ¶
type DriverSet map[providersdk.Type]providersdk.Driver
DriverSet maps provider type to the local driver instance that serves it, mirroring EmbeddedAgent's internal driver map, for the agent-side (client) half of the remote agent protocol.
type EmbeddedAgent ¶
type EmbeddedAgent struct {
// contains filtered or unexported fields
}
EmbeddedAgent is an in-process agent that dispatches directly to drivers. No network involved — the server calls driver methods in the same process. Used when boxy serve hosts providers locally.
func NewEmbeddedAgent ¶
func NewEmbeddedAgent(id, name string, drivers ...providersdk.Driver) (*EmbeddedAgent, error)
NewEmbeddedAgent creates an agent backed by the given drivers. Each driver must have a unique Type — one driver per provider type per agent.
func (*EmbeddedAgent) Allocate ¶
func (a *EmbeddedAgent) Allocate(ctx context.Context, provider providersdk.Type, id string) (map[string]any, error)
func (*EmbeddedAgent) Create ¶
func (a *EmbeddedAgent) Create(ctx context.Context, provider providersdk.Type, cfg any) (*providersdk.Resource, error)
func (*EmbeddedAgent) Delete ¶
func (a *EmbeddedAgent) Delete(ctx context.Context, provider providersdk.Type, id string) error
func (*EmbeddedAgent) Info ¶
func (a *EmbeddedAgent) Info() AgentInfo
func (*EmbeddedAgent) List ¶ added in v0.1.29
func (a *EmbeddedAgent) List(ctx context.Context, provider providersdk.Type) ([]providersdk.ResourceStatus, error)
List satisfies ResourceListingAgent for drivers that implement providersdk.ResourceLister. Returns an error for drivers that don't — callers must treat "list not supported" and "list failed" the same way (see docs/adr/0005-remote-agent-transport-and-registration.md's discussion of #133).
func (*EmbeddedAgent) PersonalizeGuest ¶
func (a *EmbeddedAgent) PersonalizeGuest(ctx context.Context, provider providersdk.Type, id string) (*providersdk.GuestPersonalizationResult, error)
func (*EmbeddedAgent) Read ¶
func (a *EmbeddedAgent) Read(ctx context.Context, provider providersdk.Type, id string) (*providersdk.ResourceStatus, error)
func (*EmbeddedAgent) Update ¶
func (a *EmbeddedAgent) Update(ctx context.Context, provider providersdk.Type, id string, op providersdk.Operation) (*providersdk.Result, error)
func (*EmbeddedAgent) UpdateStream ¶ added in v0.1.34
func (a *EmbeddedAgent) UpdateStream(ctx context.Context, provider providersdk.Type, id string, op providersdk.Operation, sink eventstream.Sink) (*providersdk.Result, error)
type GuestPersonalizingAgent ¶
type GuestPersonalizingAgent interface {
PersonalizeGuest(ctx context.Context, provider providersdk.Type, id string) (*providersdk.GuestPersonalizationResult, error)
}
GuestPersonalizingAgent is an optional agent capability for providers that expose the typed guest-personalization contract.
type LogSink ¶ added in v0.1.63
LogSink receives sanitized agent events with the authenticated agent ID. Implementations should be bounded and best effort; a logging failure must not terminate the agent transport.
type RemoteAgent ¶
type RemoteAgent struct {
// contains filtered or unexported fields
}
RemoteAgent is the server-side proxy for one connected remote agent. It implements Agent by sending Commands down the agent's gRPC stream and correlating asynchronous CommandResults back to the caller via a command_id. See docs/adr/0005-remote-agent-transport-and-registration.md.
One RemoteAgent instance corresponds to exactly one live stream. A reconnect from the same agent identity after a drop creates a *new* RemoteAgent (a fresh stream, fresh pending map) — callers holding a reference to the old instance will see every in-flight call fail once Close is called on it.
func NewRemoteAgent ¶
func NewRemoteAgent(info AgentInfo, stream boxyagentv1.AgentTransportService_ConnectServer) *RemoteAgent
NewRemoteAgent wraps a server-side stream handle for one connected agent. The caller must run Serve in its own goroutine to pump incoming frames.
func (*RemoteAgent) Allocate ¶
func (a *RemoteAgent) Allocate(ctx context.Context, provider providersdk.Type, id string) (map[string]any, error)
Allocate carries only generic JSON properties over the wire. Callers that want typed guest personalization should prefer PersonalizeGuest (below) via a GuestPersonalizingAgent type-assertion, falling back to Allocate when the remote driver doesn't implement providersdk.GuestPersonalizer.
func (*RemoteAgent) Availability ¶ added in v0.1.40
func (a *RemoteAgent) Availability() (AvailabilitySnapshot, bool)
Availability implements AvailabilityReportingAgent, returning the latest snapshot received over this connection's heartbeat stream and whether one has arrived yet.
func (*RemoteAgent) Close ¶
func (a *RemoteAgent) Close()
Close tears down this agent's view of the connection: every call currently blocked waiting on a CommandResult fails immediately rather than hanging until its context deadline. Safe to call multiple times.
It deliberately does not close the individual per-command channels in pending/streamPending — only clears the maps. Closing a.closed already unblocks every call()/UpdateStream waiter via their `case <-a.closed` select arm, so closing the per-command channels too would be redundant, and for streamPending it's actively unsafe: deliver() reads a channel reference out of the map, releases a.mu, and only then sends on it, so a concurrent Close() (called for real from Server.Revoke on a different goroutine than the one running Serve()) could delete-and-close that same channel in the gap, panicking deliver()'s send with "send on closed channel". Leaving the channels open removes that race entirely: a late send from deliver() after Close() just lands in an unread buffer that's garbage collected once the waiter (already gone via a.closed) drops its reference.
It also never closes a streamWaiter's done channel — that's exclusively UpdateStream's own deferred cleanup's job, exactly once per call. a.closed already gives deliver() a connection-teardown signal; closing done here too would risk a double-close panic if UpdateStream's own cleanup runs concurrently.
func (*RemoteAgent) Create ¶
func (a *RemoteAgent) Create(ctx context.Context, provider providersdk.Type, cfg any) (*providersdk.Resource, error)
func (*RemoteAgent) Delete ¶
func (a *RemoteAgent) Delete(ctx context.Context, provider providersdk.Type, id string) error
func (*RemoteAgent) HasHeartbeat ¶ added in v0.1.41
func (a *RemoteAgent) HasHeartbeat() bool
HasHeartbeat reports whether this connection has received at least one heartbeat. LastSeen is initialized at connection time for liveness timeout handling, so callers that need to distinguish a real heartbeat sample from that initialization must use this method.
func (*RemoteAgent) Info ¶
func (a *RemoteAgent) Info() AgentInfo
func (*RemoteAgent) LastSeen ¶
func (a *RemoteAgent) LastSeen() time.Time
LastSeen returns the time of the most recent Heartbeat (or connection start, if none has arrived yet).
func (*RemoteAgent) List ¶ added in v0.1.29
func (a *RemoteAgent) List(ctx context.Context, provider providersdk.Type) ([]providersdk.ResourceStatus, error)
List satisfies ResourceListingAgent by sending a ListCommand. The remote agent's executeCommand returns an AgentError if its driver doesn't implement providersdk.ResourceLister — same error path as any other command failure, deliberately not distinguished from a transient failure (see docs/adr/0005-remote-agent-transport-and-registration.md's discussion of #133).
func (*RemoteAgent) PersonalizeGuest ¶ added in v0.1.32
func (a *RemoteAgent) PersonalizeGuest(ctx context.Context, provider providersdk.Type, id string) (*providersdk.GuestPersonalizationResult, error)
PersonalizeGuest satisfies GuestPersonalizingAgent by sending a PersonalizeGuestCommand. An empty PersonalizeGuestResult (zero properties) means either the remote driver doesn't implement providersdk.GuestPersonalizer or it does but had nothing to report — both collapse to nil, nil here so callers fall back to the generic Allocate path exactly as EmbeddedAgent's callers do (see internal/pool/provisioner_agent.go's AgentProvisioner.Allocate).
func (*RemoteAgent) Read ¶
func (a *RemoteAgent) Read(ctx context.Context, provider providersdk.Type, id string) (*providersdk.ResourceStatus, error)
func (*RemoteAgent) RequestLogs ¶ added in v0.1.64
RequestLogs asks this connected agent for a bounded slice of its local diagnostic history. The agent identity is established by the mTLS stream; the request ID is only a correlation value for the returned LogBatch.
func (*RemoteAgent) Serve ¶
func (a *RemoteAgent) Serve() error
Serve reads AgentMessages off the stream until it ends for any reason, dispatching Heartbeats to LastSeen and CommandResults to pending callers. It must be run in its own goroutine, one per connection. When it returns, Close has already been called, failing every still-pending call.
func (*RemoteAgent) SetLogSink ¶ added in v0.1.63
func (a *RemoteAgent) SetLogSink(sink LogSink)
SetLogSink attaches the server-side destination for authenticated agent log batches. It must be called before Serve starts.
func (*RemoteAgent) Update ¶
func (a *RemoteAgent) Update(ctx context.Context, provider providersdk.Type, id string, op providersdk.Operation) (*providersdk.Result, error)
func (*RemoteAgent) UpdateStream ¶ added in v0.1.34
func (a *RemoteAgent) UpdateStream(ctx context.Context, provider providersdk.Type, id string, op providersdk.Operation, sink eventstream.Sink) (*providersdk.Result, error)
type RemoteClientConfig ¶
type RemoteClientConfig struct {
AgentName string
Token string
// AgentVersion is this agent binary's version string, sent on every
// RegisterRequest so the server can refuse a version-mismatched
// connection (see #167) rather than let skewed agent/server builds
// talk an ambiguous protocol to each other.
AgentVersion string
ProviderTypes []providersdk.Type
Drivers DriverSet
HeartbeatInterval time.Duration // default 15s if zero; overridden by the server's RegisterResponse if set
// AvailabilitySampleTimeout bounds each individual
// providersdk.AvailabilityReporter query sampled before a Heartbeat is
// sent — see sampleAvailability. Zero uses
// defaultAvailabilitySampleTimeout; tests override it to avoid real
// sleeps, mirroring hyperv.Driver.memoryQueryTimeout's pattern.
AvailabilitySampleTimeout time.Duration
// LogShipper optionally buffers safe agent diagnostics. RunSession flushes
// one bounded batch over the authenticated agent stream after each
// heartbeat; a failed flush is retained for retry.
LogShipper *diagnostics.Shipper
// LogStore is the agent-local retained diagnostics history. It is queried
// only when the authenticated server sends a LogRequest; no history is
// transmitted during ordinary heartbeats.
LogStore diagnostics.Store
// OnRegistered is invoked once per successful registration (both the
// first, token-based registration and any later cert-based reconnect)
// with the server's RegisterResponse. The caller is responsible for
// persisting ClientCertificatePem/CaCertificatePem to disk on the
// first, token-based registration so future process restarts can
// reconnect without a token.
OnRegistered func(*boxyagentv1.RegisterResponse)
Logger *slog.Logger
}
RemoteClientConfig configures one boxy agent process's connection to a server. Token is the single-use registration token; it should only be set for the very first connection attempt of a process's lifetime — every subsequent reconnect (whether from Run's own backoff loop or a future process restart) authenticates via the client certificate issued in OnRegistered instead.
type ResourceListingAgent ¶ added in v0.1.29
type ResourceListingAgent interface {
List(ctx context.Context, provider providersdk.Type) ([]providersdk.ResourceStatus, error)
}
ResourceListingAgent is an optional agent capability for providers whose underlying driver implements providersdk.ResourceLister. Not every driver supports enumeration, so callers must type-assert for this rather than relying on it being part of Agent.
type StreamingAgent ¶ added in v0.1.34
type StreamingAgent interface {
UpdateStream(ctx context.Context, provider providersdk.Type, id string, op providersdk.Operation, sink eventstream.Sink) (*providersdk.Result, error)
}
StreamingAgent is an optional capability for agents that can carry live provider events back to the server.