agentapi

package
v0.6.1 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 2, 2026 License: AGPL-3.0 Imports: 69 Imported by: 0

Documentation

Index

Constants

View Source
const ActionStringPreviewBytes = 8 * 1024 // 8 KiB

ActionStringPreviewBytes caps verbose stream-like fields kept in the runs.actions audit log.

View Source
const MaxBufferedResponseBytes = 20 << 20 // 20 MiB

MaxBufferedResponseBytes is the hard cap for any response a buffered SDK method accumulates fully in memory before returning to the caller. Overflow surfaces as agentsdk.ErrOutputTooLarge with no partial result.

The cap also applies Airlock-side as defense in depth on integration proxies, so a misbehaving SDK cannot exhaust Airlock memory.

Sized so structured small responses (JSON API replies, HTML pages, CLI tool summaries) pass through and larger data uses streaming storage.

Variables

View Source
var ErrMCPForbidden = errors.New("mcp: access denied for target agent")

ErrMCPForbidden means the caller is authenticated but the target's access ladder rejects them (e.g. non-member on a non-member-closed target). The handler maps this to HTTP 403.

View Source
var ErrMCPUnauthenticated = errors.New("mcp: unauthenticated and target disallows public mcp")

ErrMCPUnauthenticated means the caller presented no credentials and the target agent does not allow public MCP. The handler maps this to HTTP 401.

Functions

func DiscoverMCPAuth

func DiscoverMCPAuth(ctx context.Context, httpClient *http.Client, serverURL string) (*oauth.DiscoveryResult, error)

DiscoverMCPAuth runs RFC 9728/8414 discovery on an MCP server URL.

func DiscoverMCPTools

func DiscoverMCPTools(ctx context.Context, httpClient *http.Client, serverURL string, authInjection []byte, creds string) ([]mcpToolInfo, string, error)

DiscoverMCPTools connects to a remote MCP server and returns its tool schemas plus the server-level `instructions` it advertised in the initialize result (empty when the server set none).

func ExtractCanonicalKeys

func ExtractCanonicalKeys(partsJSON []byte, agentID string) []string

ExtractCanonicalKeys reads `s3ref:K` sentinels from the stored goai-shaped parts JSON (image.image / file.data fields) and returns the canonical `llm/agents/<agentID>/K` keys. The sentinel survives the goai.Content marshal roundtrip since it's just a string in Image/Data.

func ExtractTextSummary

func ExtractTextSummary(parts []wire.DisplayPart) string

ExtractTextSummary builds a text summary from display parts for the content column.

func InjectAuth

func InjectAuth(req *http.Request, authInjectionJSON []byte, creds string)

InjectAuth adds credentials to the upstream request based on the auth injection config.

func PostToConversation

func PostToConversation(ctx context.Context, deps PostDeps, opts PostOpts) error

PostToConversation stores a message, delivers it via the appropriate channel (WebSocket or bridge), and optionally triggers an LLM turn.

func PublishRunEvents

func PublishRunEvents(
	ctx context.Context,
	body io.ReadCloser,
	pubsub *realtime.PubSub,
	db *pgxpool.Pool,
	agentID, runID uuid.UUID,
	conversationID string,
	userID string,
	parentInfo *ParentRunInfo,
	logger *zap.Logger,
) (responseText string, newMessages []message.Message, tokensIn, tokensOut int32)

PublishRunEvents reads NDJSON from body, publishes typed proto events to WS, accumulates the assistant response text, and returns it along with token usage. This runs in a goroutine after the HTTP response has been sent.

userID is the conversation owner — applied to every emitted envelope for user-id-based delivery gating. Pass empty for system-level (no-conversation) runs; an empty UserID delivers to every subscriber on the topic (legacy behaviour).

parentInfo, when non-nil, causes every typed event to be mirrored onto the parent agent's topic with a Subagent tag — the chat UI uses this to attach live sub-run progress to the parent's tool-call card.

func PublishRunTerminal

func PublishRunTerminal(ctx context.Context, pubsub *realtime.PubSub, agentID, runID uuid.UUID, status, errMsg string)

PublishRunTerminal publishes the appropriate WebSocket event for a run's terminal state. Mirrors what PublishRunEvents emits from the agent's NDJSON stream so frontends and bridges see one event regardless of which path the run took to terminal — happy path (NDJSON), cancel (CancelRun closes the stream early), the agent's detached r.Complete POST after the stream died, or the sweeper for stuck rows.

Maps airlock-side status strings to WS event types:

  • "error" / "failed" / "timeout" → run.error
  • everything else (success, tool_errors, cancelled) -> run.complete

Idempotent at the client: chat.ts ignores events for runIDs already finalized locally, so a duplicate from the happy-path NDJSON + this helper is harmless.

func ResolveMediaPartsJSON

func ResolveMediaPartsJSON(ctx context.Context, s3Client *storage.S3Client, logger *zap.Logger, partsJSON []byte) []byte

ResolveMediaPartsJSON walks a JSON array of parts and presigns S3 source keys on media entries. Non-media parts (tool-call, tool-result) are passed through verbatim — deserializing into DisplayPart would strip unknown fields like toolCallId, toolName, and args.

func ServeStoragePath

func ServeStoragePath(w http.ResponseWriter, r *http.Request, database *db.DB, s3 *storage.S3Client, files *agentstoragesvc.Service, agentID uuid.UUID, path, jwtSecret, publicURL string, logger *zap.Logger)

ServeStoragePath serves reads from a registered directory on the subdomain proxy's /__air/storage{path} endpoint, gating by the directory's read_access and exact scope identity:

  • "public" — served unauthenticated when the directory is unscoped
  • "user" — requires a valid __air_session cookie + agent membership
  • "admin" — requires admin role on the agent
  • unknown — 404

Missing cookies on user/admin dirs get redirected through the relay flow (rejectOrRedirect) so a click-from-chat triggers login. Once logged in, the user lands back on the same URL with a session cookie set and the file streams.

Denied and unknown paths return the same response after authentication so URL guessing does not reveal whether an object exists.

func SynthesizeOrphanToolResults

func SynthesizeOrphanToolResults(ctx context.Context, q *dbq.Queries, runID uuid.UUID, status string, logger *zap.Logger)

SynthesizeOrphanToolResults inserts a synthetic role=tool message for every tool-call this run emitted that doesn't have a paired tool-result row. Required for the next LLM turn: provider APIs (Anthropic, OpenAI) reject inputs where an assistant's tool_use isn't followed by a tool_result with the matching id. Common after cancel, deadline-exceeded, panic mid-tool, or any path where the agent didn't get to write the tool's result before terminating.

The synthesized output text is derived from the run's terminal status so the LLM has some signal about why the tool didn't complete:

  • "cancelled" → "Cancelled by user."
  • "timeout" → "Tool timed out."
  • else → "Tool execution failed."

Best-effort: failures are logged but don't block the caller. The SessionLoad lazy-synthesis path is the safety net if this misses.

func ValidateSuspendedCheckpoint

func ValidateSuspendedCheckpoint(checkpoint []byte) error

ValidateSuspendedCheckpoint validates the wire fields Airlock needs before a suspended run can become resumable. It intentionally does not depend on Sol.

Types

type BridgePartsDeliverer

type BridgePartsDeliverer interface {
	SendParts(ctx context.Context, bridgeID uuid.UUID, externalID string, parts []wire.DisplayPart) error
}

BridgePartsDeliverer is the subset of trigger.BridgeManager needed for message delivery.

type Config

type Config struct {
	DB                     *db.DB
	Encryptor              secrets.Store
	OAuthClient            *oauth.Client
	S3                     *storage.S3Client
	Files                  *agentstoragesvc.Service
	Jobs                   *jobssvc.Service
	ConnectorJobs          *connectorjobssvc.Service
	ConnectorDirectories   *connectordirectoriessvc.Service
	ConnectorOrchestration *connectororchestrationsvc.Service
	Builder                *builder.BuildService
	PubSub                 *realtime.PubSub
	BridgeMgr              BridgePartsDeliverer
	Scheduler              scheduleReconciler
	PublicURL              string
	AgentBaseURL           func(slug string) string
	LLMProxyURL            string
	ForceInlineAttachments bool
	JWTSecret              string
	Dispatcher             *trigger.Dispatcher
	HTTPNetwork            *networkpolicy.Policy
	Logger                 *zap.Logger
}

Config bundles the dependencies New requires. Mirrors the struct fields of Handler one-for-one; api/router.go's RouterConfig translates its own merged config into this on wire-up.

type Handler

type Handler struct {
	// contains filtered or unexported fields
}

func New

func New(c Config) *Handler

New constructs the agent-internal HTTP surface. Fail-loud on nil deps — every required field is mandatory (airlock fail-loud rule).

func (*Handler) AgentHTTP

func (h *Handler) AgentHTTP(w http.ResponseWriter, r *http.Request)

AgentHTTP handles POST /api/agent/http — raw HTTP proxy for permitted URLs. No auth injection (use proxy/{slug} for authenticated connections).

func (*Handler) CallMCPTool

func (h *Handler) CallMCPTool(ctx context.Context, agentID uuid.UUID, slug string, req wire.MCPToolCallRequest) (wire.MCPToolCallResponse, error)

func (*Handler) CancelConnectorJob added in v0.6.1

func (h *Handler) CancelConnectorJob(w http.ResponseWriter, r *http.Request)

func (*Handler) CancelConnectorOrchestration added in v0.6.1

func (h *Handler) CancelConnectorOrchestration(w http.ResponseWriter, r *http.Request)

func (*Handler) CancelJob added in v0.5.0

func (h *Handler) CancelJob(w http.ResponseWriter, r *http.Request)

func (*Handler) ConnectorCommand added in v0.6.1

func (h *Handler) ConnectorCommand(w http.ResponseWriter, r *http.Request)

func (*Handler) ConnectorDirectoryDelete added in v0.6.1

func (h *Handler) ConnectorDirectoryDelete(w http.ResponseWriter, r *http.Request)

func (*Handler) ConnectorDirectoryExport added in v0.6.1

func (h *Handler) ConnectorDirectoryExport(w http.ResponseWriter, r *http.Request)

func (*Handler) ConnectorDirectoryImport added in v0.6.1

func (h *Handler) ConnectorDirectoryImport(w http.ResponseWriter, r *http.Request)

func (*Handler) ConnectorDirectoryList added in v0.6.1

func (h *Handler) ConnectorDirectoryList(w http.ResponseWriter, r *http.Request)

func (*Handler) ConnectorDirectoryMove added in v0.6.1

func (h *Handler) ConnectorDirectoryMove(w http.ResponseWriter, r *http.Request)

func (*Handler) ConnectorDirectoryRead added in v0.6.1

func (h *Handler) ConnectorDirectoryRead(w http.ResponseWriter, r *http.Request)

func (*Handler) ConnectorDirectoryStat added in v0.6.1

func (h *Handler) ConnectorDirectoryStat(w http.ResponseWriter, r *http.Request)

func (*Handler) ConnectorDirectoryWrite added in v0.6.1

func (h *Handler) ConnectorDirectoryWrite(w http.ResponseWriter, r *http.Request)

func (*Handler) CreateConnectorOrchestration added in v0.6.1

func (h *Handler) CreateConnectorOrchestration(w http.ResponseWriter, r *http.Request)

func (*Handler) CreateRun

func (h *Handler) CreateRun(w http.ResponseWriter, r *http.Request)

CreateRun handles POST /api/agent/run/create. Called by agent containers to create a run record for programmatic runs (e.g. from route handlers).

func (*Handler) Embed

func (h *Handler) Embed(w http.ResponseWriter, r *http.Request)

Embed handles POST /api/agent/llm/embedding.

func (*Handler) EnqueueJob added in v0.5.0

func (h *Handler) EnqueueJob(w http.ResponseWriter, r *http.Request)

func (*Handler) GetCheckpoint

func (h *Handler) GetCheckpoint(w http.ResponseWriter, r *http.Request)

GetCheckpoint handles GET /api/agent/run/{runID}/checkpoint.

func (*Handler) GetConnectorJob added in v0.6.1

func (h *Handler) GetConnectorJob(w http.ResponseWriter, r *http.Request)

func (*Handler) GetConnectorOrchestration added in v0.6.1

func (h *Handler) GetConnectorOrchestration(w http.ResponseWriter, r *http.Request)

func (*Handler) GetEnvVarValue

func (h *Handler) GetEnvVarValue(w http.ResponseWriter, r *http.Request)

GetEnvVarValue handles GET /api/agent/env-vars/{slug} — runtime read of a configured value. Returns 404 if the slot exists but has no configured value (so the agent can fall back to default or fail loudly per its own policy).

func (*Handler) GetJob added in v0.5.0

func (h *Handler) GetJob(w http.ResponseWriter, r *http.Request)

func (*Handler) ImageGenerate

func (h *Handler) ImageGenerate(w http.ResponseWriter, r *http.Request)

ImageGenerate handles POST /api/agent/llm/image.

func (*Handler) LLMStream

func (h *Handler) LLMStream(w http.ResponseWriter, r *http.Request)

LLMStream handles POST /api/agent/llm/stream.

func (*Handler) ListJobs added in v0.5.0

func (h *Handler) ListJobs(w http.ResponseWriter, r *http.Request)

func (*Handler) ListMCPTools

func (h *Handler) ListMCPTools(ctx context.Context, agentID uuid.UUID, slug string) (integrationservice.MCPTools, error)

func (*Handler) MCPToolCall

func (h *Handler) MCPToolCall(w http.ResponseWriter, r *http.Request)

MCPToolCall handles POST /api/agent/mcp/{slug}/tools/call. Stateless: connect → initialize → tools/call → disconnect.

func (*Handler) Print

func (h *Handler) Print(w http.ResponseWriter, r *http.Request)

Print handles POST /api/agent/print. Processes display parts (upload bytes, copy tmp files to permanent media), then routes to the target conversation(s) — either direct (the `output` JS binding) or via topic subscriptions (TopicHandle.Publish).

func (*Handler) RequestConnection

func (h *Handler) RequestConnection(ctx context.Context, agentID uuid.UUID, slug string, req wire.ProxyRequest) (integrationservice.ConnectionResult, error)

func (*Handler) RunComplete

func (h *Handler) RunComplete(w http.ResponseWriter, r *http.Request)

RunComplete handles POST /api/agent/run/complete.

func (*Handler) Seal

func (h *Handler) Seal(w http.ResponseWriter, r *http.Request)

Seal handles POST /api/agent/seal.

func (*Handler) Search

func (h *Handler) Search(w http.ResponseWriter, r *http.Request)

Search handles POST /api/agent/search — proxies web search requests from agent containers without exposing API keys.

func (*Handler) ServiceProxy

func (h *Handler) ServiceProxy(w http.ResponseWriter, r *http.Request)

ServiceProxy handles POST /api/agent/proxy/{slug}.

func (*Handler) SessionAppend

func (h *Handler) SessionAppend(w http.ResponseWriter, r *http.Request)

SessionAppend handles POST /api/agent/session/{convID}/messages. Appends new messages to the conversation.

func (*Handler) SessionCompact

func (h *Handler) SessionCompact(w http.ResponseWriter, r *http.Request)

SessionCompact handles POST /api/agent/session/{convID}/compact. Non-destructive compaction: inserts a checkpoint marker row + the summary messages, then advances agent_conversations.context_checkpoint_message_id to point at the first summary message. Pre-checkpoint history stays in the DB for UI display; Sol's SessionStore filters it out on the next Load.

func (*Handler) SessionLoad

func (h *Handler) SessionLoad(w http.ResponseWriter, r *http.Request)

SessionLoad handles GET /api/agent/session/{convID}/messages. Returns conversation history as []session.Message.

func (*Handler) SpeechGenerate

func (h *Handler) SpeechGenerate(w http.ResponseWriter, r *http.Request)

SpeechGenerate handles POST /api/agent/llm/speech.

func (*Handler) StartConnectorJob added in v0.6.1

func (h *Handler) StartConnectorJob(w http.ResponseWriter, r *http.Request)

func (*Handler) StorageCopy

func (h *Handler) StorageCopy(w http.ResponseWriter, r *http.Request)

StorageCopy handles POST /api/agent/storage/copy. Both src and dst are absolute paths.

func (*Handler) StorageDelete

func (h *Handler) StorageDelete(w http.ResponseWriter, r *http.Request)

StorageDelete handles DELETE /api/agent/storage/*.

func (*Handler) StorageInfo

func (h *Handler) StorageInfo(w http.ResponseWriter, r *http.Request)

StorageInfo handles POST /api/agent/storage/info. Body: {path}. Returns FileInfo with the original filename surfaced from S3 metadata when the upload set it (chat uploads do; raw writeFile does not).

func (*Handler) StorageList

func (h *Handler) StorageList(w http.ResponseWriter, r *http.Request)

StorageList handles GET /api/agent/storage. Query params:

  • path=uploads/ → list under this storage path (slashless; trailing '/' optional)
  • path= (empty) → list the agent root
  • recursive=true|false (default false; one level only)

func (*Handler) StorageLoad

func (h *Handler) StorageLoad(w http.ResponseWriter, r *http.Request)

StorageLoad handles GET /api/agent/storage/*.

func (*Handler) StorageShare

func (h *Handler) StorageShare(w http.ResponseWriter, r *http.Request)

StorageShare handles POST /api/agent/storage/share. Returns a presigned, unauthenticated, time-limited URL for the given storage path. Used by the agent's shareFileURL JS binding (and ShareFileURL Go method) to hand out a link the user — or a third party — can fetch directly, without going through agent membership / login on the subdomain proxy.

The URL is signed for the public S3 endpoint when configured (so browsers, LLM providers, and external tools can resolve it from outside the docker network). Falls back to the internal endpoint otherwise.

404s on missing files via HeadObject pre-check so the LLM gets a clear signal instead of a working URL that 404s when followed.

func (*Handler) StorageStore

func (h *Handler) StorageStore(w http.ResponseWriter, r *http.Request)

StorageStore handles PUT /api/agent/storage/*. Path-based: the wildcard is the absolute path under the agent's storage root. Original filename can ride on `X-Filename` and is persisted as S3 object metadata.

func (*Handler) Sync

func (h *Handler) Sync(w http.ResponseWriter, r *http.Request)

func (*Handler) TopicSubscribe

func (h *Handler) TopicSubscribe(w http.ResponseWriter, r *http.Request)

TopicSubscribe handles POST /api/agent/topic/{slug}/subscribe. Subscribes the given conversation to the agent's topic.

func (*Handler) TopicUnsubscribe

func (h *Handler) TopicUnsubscribe(w http.ResponseWriter, r *http.Request)

TopicUnsubscribe handles DELETE /api/agent/topic/{slug}/subscribe. Unsubscribes the given conversation from the agent's topic.

func (*Handler) Transcribe

func (h *Handler) Transcribe(w http.ResponseWriter, r *http.Request)

Transcribe handles POST /api/agent/llm/transcription.

func (*Handler) Unseal

func (h *Handler) Unseal(w http.ResponseWriter, r *http.Request)

Unseal handles POST /api/agent/unseal. A decrypt failure is a 400, not a 500: the usual cause is a value sealed by (bound to) a different agent or a corrupted blob — a bad request, not a server fault.

func (*Handler) UpdateJobProgress added in v0.5.0

func (h *Handler) UpdateJobProgress(w http.ResponseWriter, r *http.Request)

func (*Handler) Upgrade

func (h *Handler) Upgrade(w http.ResponseWriter, r *http.Request)

Upgrade handles POST /api/agent/upgrade.

func (*Handler) UpsertEnvVar

func (h *Handler) UpsertEnvVar(w http.ResponseWriter, r *http.Request)

UpsertEnvVar handles PUT /api/agent/env-vars/{slug}. The agent declares its required env vars at sync time; we upsert the registration row (slot definition, no value yet).

func (*Handler) WebFetch

func (h *Handler) WebFetch(w http.ResponseWriter, r *http.Request)

WebFetch handles POST /api/agent/webfetch — proxies URL fetching for agent containers. Currently a pass-through; future: logging, URL allowlists, rate limiting.

type MCPPrincipal

type MCPPrincipal struct {
	Kind          MCPPrincipalKind
	UserID        uuid.UUID // anon: uuid.Nil; user/agent/oauth: original user
	CallerAgentID uuid.UUID // agent only
	ParentRunID   uuid.UUID // agent only
	ClientID      string    // oauth only — for audit / logging
}

MCPPrincipal carries the caller identity for the MCP server endpoint. Built by the MCP handler from request headers / JWT and threaded into computeA2ACallerAccess.

type MCPPrincipalKind

type MCPPrincipalKind int

MCPPrincipalKind discriminates how a caller hit the A2A MCP endpoint.

const (
	// MCPPrincipalAnon means no JWT was presented.
	MCPPrincipalAnon MCPPrincipalKind = iota
	// MCPPrincipalUser means a user JWT was presented (web SPA path).
	// UserID is set.
	MCPPrincipalUser
	// MCPPrincipalAgent means an agent JWT was presented (sibling A2A
	// caller). CallerAgentID + ParentRunID are set; UserID is derived
	// from the parent run's conversation and also populated.
	MCPPrincipalAgent
	// MCPPrincipalOAuthClient means an OAuth-issued access token was
	// presented (external MCP client — Claude Desktop, Codex CLI,
	// etc.). UserID + ClientID are set. The audience binding (token
	// `aud` matches the target agent's canonical resource URL) is
	// verified upstream in resolvePrincipal, so this principal is
	// treated identically to MCPPrincipalUser for the access ladder.
	MCPPrincipalOAuthClient
)

type MCPServer

type MCPServer struct {
	// contains filtered or unexported fields
}

MCPServer is the JSON-RPC 2.0 server exposed at /api/agent/{identifier}/mcp. Each Airlock agent becomes an MCP server: its registered tools surface plus a built-in `prompt` meta-tool that drives the agent's full LLM loop.

func NewMCPServer

func NewMCPServer(dispatcher *trigger.Dispatcher, pubsub *realtime.PubSub, logger *zap.Logger) *MCPServer

NewMCPServer wires the MCP handler to its dependencies. Mounted at POST /api/agent/{identifier}/mcp by router.go. pubsub is required: when an A2A child run streams its events back through here, the MCP server also mirrors them to the parent agent's WS topic with a SubagentInfo tag so the parent's chat UI shows sub-run progress.

func (*MCPServer) ServeHTTP

func (s *MCPServer) ServeHTTP(w http.ResponseWriter, r *http.Request, h *Handler)

ServeHTTP handles one JSON-RPC message. Notification methods (no id) get a 202 with no body; request methods return either a single JSON-RPC response (most methods) or an SSE stream of notifications terminated by a final response (the `prompt` meta-tool).

The handler reads its dependencies — DB pool, JWT secret — from the Handler stored on the request via a closure in router.go.

func (*MCPServer) ServePublicHTTP

func (s *MCPServer) ServePublicHTTP(w http.ResponseWriter, r *http.Request, h *Handler)

ServePublicHTTP is the no-auth /public-mcp route. Mounted always; 404s unless the target agent has allow_public_mcp = true. Skips principal resolution entirely — every caller is Anon, which the access ladder maps to AccessPublic on agents that opted in.

type ParentRunInfo

type ParentRunInfo struct {
	AgentID        uuid.UUID
	ConvID         string
	UserID         string
	ChildAgentID   uuid.UUID
	ChildRunID     uuid.UUID
	ChildAgentSlug string
}

ParentRunInfo carries the parent run's coordinates when an A2A child run's events should mirror onto the parent's WS topic with a SubagentInfo tag, so the parent's chat UI can render sub-run progress under its active tool-call card. nil for top-level runs.

type PostDeps

type PostDeps struct {
	DB         *db.DB
	PubSub     *realtime.PubSub
	BridgeMgr  BridgePartsDeliverer // nil if no bridge configured
	Dispatcher *trigger.Dispatcher  // nil if TriggerLLM not used
	S3         *storage.S3Client    // for resolving presigned URLs
	Logger     *zap.Logger
}

PostDeps holds shared dependencies for PostToConversation.

type PostOpts

type PostOpts struct {
	AgentID        uuid.UUID
	ConversationID uuid.UUID
	RunID          uuid.UUID          // zero = no run linkage
	Role           string             // "assistant", "system"
	Text           string             // plain text content
	Parts          []wire.DisplayPart // rich content (optional)
	Source         string             // "notification", "system", etc.
	Ephemeral      bool               // stored for UI but excluded from LLM context
	TriggerLLM     bool               // forward to agent for a response turn
	LLMMessage     string             // message text for the LLM turn (if TriggerLLM)
}

PostOpts configures a message post to a conversation.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL