Runtime
Architecture Overview
The Goa-AI runtime orchestrates the plan/execute/resume loop, enforces policies, manages state, and coordinates with engines, planners, tools, memory, hooks, and feature modules.
| Layer | Responsibility |
|---|---|
| DSL + Codegen | Produce agent registries, tool specs/codecs, completion specs/codecs, workflows, MCP adapters |
| Runtime Core | Orchestrates plan/start/resume loop, policy enforcement, hooks, memory, streaming |
| Workflow Engine Adapter | Temporal adapter implements engine.Engine; other engines can plug in |
| Host Runtime Store | Stores session scope, run state, continuation checkpoints, and run records that cannot change after insertion together |
| Feature Modules | Optional integrations (MCP, Pulse, memory and prompt stores, model providers) |
High-Level Agentic Architecture
At runtime, Goa-AI organizes your system around a small set of composable constructs:
Agents: Long-lived orchestrators identified by
agent.Ident(for example,service.chat). Each agent owns a planner, run policy, generated workflows, and tool registrations.Runs: A single execution of an agent. Runs are identified by a
RunIDand tracked viarun.Contextandrun.Handle. Sessionful runs are grouped bySessionIDandTurnIDto form conversations; one-shot runs are explicitly sessionless.Toolsets & tools: Named collections of capabilities, identified by
tools.Ident(service.toolset.tool). Service-backed toolsets call APIs; agent-backed toolsets run other agents as tools.Completions: Service-owned typed direct assistant-output contracts generated under
gen/<service>/completions. Completion helpers attach provider-enforced structured output to unary and direct-streaming model requests, then decode the canonical typed payload through generated codecs.Planners: Your LLM-driven strategy layer implementing
PlanStart/PlanResume. Planners decide when to call tools versus answer directly; the runtime enforces caps and time budgets around those decisions.Run tree & agent-as-tool: When an agent calls another agent as a tool, the runtime starts a real child run with its own
RunID. The parentToolResultcarries aRunLink(*run.Handle) pointing to the child, and a correspondingchild_run_linkedstream event is emitted so UIs can correlate parent tool calls with child run IDs without guessing.Session-owned streams & profiles: Goa-AI publishes typed
stream.Eventvalues into a session-owned stream (session/<session_id>). Events carry bothRunIDandSessionID, and include an explicit boundary marker (run_stream_end) so consumers can close SSE/WebSocket deterministically without timers.stream.StreamProfileselects which event kinds are visible for a given audience (chat UI, debug, metrics).
Quick Start
package main
import (
"context"
"time"
chat "example.com/assistant/gen/orchestrator/agents/chat"
"goa.design/goa-ai/runtime/agent/model"
"goa.design/goa-ai/runtime/agent/runtime"
storageinmem "goa.design/goa-ai/runtime/agent/storage/inmem"
)
func main() {
// In-memory engine is the default; pass WithEngine for Temporal or custom engines.
store := storageinmem.New()
rt := runtime.New(store)
ctx := context.Background()
err := chat.RegisterChatAgent(ctx, rt, chat.ChatAgentConfig{Planner: newChatPlanner()})
if err != nil {
panic(err)
}
// Sessions are first-class: create a session before starting runs under it.
if _, err := store.CreateSession(ctx, "session-1", time.Now().UTC()); err != nil {
panic(err)
}
client := chat.NewClient(rt)
out, err := client.Run(ctx, "session-1", []*model.Message{{
Role: model.ConversationRoleUser,
Parts: []model.Part{model.TextPart{Text: "Summarize the latest status."}},
}})
if err != nil {
panic(err)
}
// Use out.RunID, out.Final (the assistant message), etc.
}
Typed Direct Completions
Not every structured interaction should be modeled as a tool call. When your
service needs a typed final assistant answer, declare Completion(...) in the
DSL and regenerate.
goa gen emits gen/<service>/completions with:
- typed result and union types
- private result schemas and generated codecs
- generated
Complete<Name>(ctx, client, req)helpers - typed
StreamComplete<Name>(ctx, client, req)helpers <Name>Example()when the root result has an authoredExample(...)
Services may declare completions without declaring any Agent(...). Agent
quickstart/example scaffolding is emitted only for services that actually own
agents.
Those helpers clone the request, attach provider-neutral structured output
metadata, call the underlying model.Client, and decode the canonical typed
payload through the generated codec:
resp, err := taskcompletion.CompleteDraftFromTranscript(ctx, modelClient, &model.Request{
Messages: []*model.Message{{
Role: model.ConversationRoleUser,
Parts: []model.Part{model.TextPart{Text: "Create a startup investigation task."}},
}},
})
if err != nil {
panic(err)
}
fmt.Println(resp.Value.Name)
Every low-level model.StructuredOutput requires a nonempty name. Generated
helpers derive it from the validated completion DSL. Unary completion makes
exactly one model call. Invalid JSON returns a non-retryable
planner.OutputContractError and a nil response; it never triggers a correction
request. On success, resp.ModelResponse contains the exact provider response
and token usage.
Streaming completions return completion.Streamer[T]. Recv exposes preview
fragments, while Value() stays unavailable until the stream ends and the
terminal response agrees with the final completion:
stream, err := taskcompletion.StreamCompleteDraftFromTranscript(ctx, modelClient, &model.Request{
Messages: []*model.Message{{
Role: model.ConversationRoleUser,
Parts: []model.Part{model.TextPart{Text: "Create a startup investigation task."}},
}},
})
if err != nil {
panic(err)
}
defer stream.Close()
for {
chunk, err := stream.Recv()
if errors.Is(err, io.EOF) {
break
}
if err != nil {
panic(err)
}
// Render preview completion_delta chunks here when useful.
_ = chunk
}
value, ok := stream.Value()
if !ok {
panic("completion stream ended without a typed value")
}
fmt.Println(value.Name)
Typed completion helpers are intentionally strict:
- Unary helpers accept unary requests only.
- Completion names are validated at the DSL boundary: 1-64 ASCII characters,
letters/digits/
_/-only, and must start with a letter or digit. - Unary and streaming helpers reject tool-enabled requests and caller-supplied
StructuredOutput. - Streaming providers may emit
completion_delta*previews plus exactly one finalcompletion, or reject the request explicitly. - The typed wrapper releases
Value()only after clean end-of-stream and full validation. There is no public decoder that accepts an unchecked chunk. - Completion streams use their generated typed wrapper directly; planner streaming helpers are for assistant transcript text and tool calls.
- Providers that do not implement structured output surface
model.ErrStructuredOutputUnsupported. - Generated schemas are canonical and provider-neutral; provider adapters may normalize them to a supported subset, but must fail explicitly when they cannot preserve the declared contract.
Client-Only vs Worker
Two roles use the runtime:
- Client-only (submit runs): Constructs a runtime with a client-capable engine and does not register agents. Use the generated
<agent>.NewClient(rt), which carries the generatedAgentDefinitionshared with remote workers. - Worker (execute runs): Constructs a runtime with a worker-capable engine, registers toolsets and agents, then seals registration so polling starts only after the local runtime registry is complete.
Each generated AgentDefinition is the complete, fixed contract for one agent.
It owns the workflow name, default task queue, generated tool contracts,
required labels, completion policy, and definitions of every reachable child
agent. Callers use it to validate and route work before the engine accepts a
workflow; workers use the same value when registering that workflow. An
individual start may select another task queue with WithTaskQueue, but
handwritten registrations must not define a second route or child-agent graph.
Client-Only Example
rt := runtime.New(runtimeStore, runtime.WithEngine(temporalClient)) // engine client
// The host session service has already created "s1".
// No agent registration is needed in a caller-only process.
client := chat.NewClient(rt)
out, err := client.Run(ctx, "s1", msgs)
Sessionless One-Shot Runs
Use StartOneShot and OneShotRun when you want durable work that is not attached to an existing session.
Start/Runare sessionful: they require a concreteSessionID, participate in session lifecycle, and emit session-scoped stream events.StartOneShot/OneShotRunare sessionless: they take noSessionID, do not create one, and store normal run metadata and records for introspection byRunID.StartOneShotreturns anengine.WorkflowHandleimmediately.OneShotRunis the blocking convenience wrapper that callshandle.Wait(ctx)for you.
client := chat.NewClient(rt)
handle, err := client.StartOneShot(ctx, msgs,
runtime.WithRunID("run-123"),
runtime.WithLabels(map[string]string{"tenant": "acme"}),
)
if err != nil {
panic(err)
}
out, err := handle.Wait(ctx)
if err != nil {
panic(err)
}
fmt.Println(out.RunID)
The lower-level Runtime.RunOneShot stores the run before it calls application
code. After the callback returns, it records prompt renders and the terminal
result even when the callback canceled its context. Temporary storage failures
retry the prepared records without invoking the callback again.
Worker Example
eng, err := temporal.NewWorker(temporal.Options{
ClientOptions: &client.Options{HostPort: "temporal:7233", Namespace: "default"},
WorkerOptions: temporal.WorkerOptions{TaskQueue: "orchestrator.chat"},
})
if err != nil {
panic(err)
}
defer eng.Close()
rt := runtime.New(runtimeStore, runtime.WithEngine(eng))
if err := chat.RegisterUsedToolsets(ctx, rt /* executors... */); err != nil {
panic(err)
}
if err := chat.RegisterChatAgent(ctx, rt, chat.ChatAgentConfig{Planner: myPlanner}); err != nil {
panic(err)
}
if err := rt.Seal(ctx); err != nil {
panic(err)
}
Plan → Execute → Resume Loop
- The engine accepts a workflow for the agent (in-memory or Temporal).
- The workflow’s first activity records the run identity and first permanent
record through
StartRootRun,StartChildRun,StartOneShotRun, orStartOneShotChildRun. Every accepted workflow storesRunStarted. A sessionful run proceeds only when its session is active; when the session has ended, the store followsRunStartedwith a canceledRunCompletedand the workflow does no planner or tool work. Memory & Sessions defines the three valid ways cancellation reasons and intent records are stored. - The runtime calls your planner’s
PlanStartwithPrepareMessagesand arun.ContextcontainingRunID,SessionID,TurnID, labels, and policy caps. - It schedules tool calls returned by the planner (planner passes canonical JSON payloads; the runtime handles encoding/decoding using generated codecs).
- It calls
PlanResumewith the surviving planner-visible tool results. Budgeted tools are visible by default. A failed bookkeeping tool schedules another planner turn according toToolFailure.Recovery.Action: correction, replanning without that tool, or finalization. The loop repeats until the planner returns a final response, a final tool result, or a successfulTerminalRuntool completes the run. If caps or deadlines force finalization, the planner may close through terminal bookkeeping tools instead of prose. As execution progresses, the run advances throughrun.Phasevalues (prompted,planning,executing_tools,synthesizing, terminal phases). - The runtime store saves immutable records as the run changes. Hooks and stream subscribers emit planner thoughts, tool start/update/end, awaits, usage, workflow, and agent-run links. An optional memory store persists transcript entries.
Run Phases
As a run progresses through the plan/execute/resume loop, it transitions through a series of lifecycle phases. These phases provide fine-grained visibility into where a run is in its execution, enabling UIs to show high-level progress indicators.
Phase Values
| Phase | Description |
|---|---|
prompted | Input has been received and the run is about to begin planning |
planning | The planner is deciding whether and how to call tools or answer directly |
executing_tools | Tools (including nested agents) are currently executing |
synthesizing | The planner is synthesizing a final answer without scheduling additional tools |
completed | The run has completed successfully |
failed | The run has failed |
canceled | The run was canceled |
Phase Transitions
A typical successful run follows this progression:
prompted → planning → executing_tools → planning → synthesizing → completed
↑__________________|
(loop while tools needed)
The runtime emits RunPhaseChanged hook events for non-terminal phases (e.g., planning, executing_tools, synthesizing) so stream subscribers can track progress in real time.
Phase vs Status
Phases are distinct from run.Status:
- Status (
running,suspended,completed,failed,canceled) is the coarse-grained lifecycle state stored in durable run metadata. There is no pre-admissionpendingstate - Phase provides finer-grained visibility into the execution loop, intended for streaming/UX surfaces
Lifecycle events: phase changes vs terminal completion
The runtime emits:
RunPhaseChangedfor non-terminal phase transitions.RunCompletedonce per run for terminal lifecycle (success / failed / canceled).
Stream subscribers translate both into workflow stream events (stream.WorkflowPayload):
- Non-terminal updates (from
RunPhaseChanged):phaseonly. - Terminal update (from
RunCompleted):status+ terminalphase, plus structured error fields on failures.
Terminal status mapping
status="success"→phase="completed"status="failed"→phase="failed"status="canceled"→phase="canceled"
Cancellation is not an error
For status="canceled", the stream payload must not include a user-facing error. Consumers should treat cancellation as a terminal, non-error end state.
Failures are structured
For status="failed", the stream payload includes:
error_kind: stable classifier for UX/decisioning (provider kinds likerate_limited,unavailable, or runtime kinds liketimeout/internal)retryable: whether retrying may succeed without changing inputerror: user-safe message suitable for direct displaydebug_error: diagnostic error text; the application decides who may see it
Terminal identity
RunCompleted carries Labels: the run-scoped labels provided when the run
started (RunInput.Labels, set via runtime.WithLabels(...)), nil when the run
had none. Completion subscribers can attribute the terminal outcome — success,
failed, or canceled — without maintaining their own run-ID-to-identity map. The
same labels are exposed on run.Snapshot.Labels for polling readers, replayed
from the durable RunStarted record, so run identity survives process restarts
on both engines. Labels merged by policy decisions mid-run are not included;
they remain observable via PolicyDecision events.
Error Diagnostics
Goa-AI preserves complete valid UTF-8 diagnostic messages and typed provider error text without per-field length cutoffs. Applications own disclosure: they choose what their instrumentation stores and who can read or display it. Planner and Temporal activity spans receive the original error before workflow transport or error conversion. Workflow replay does not re-emit those activity diagnostics. Display summaries, failure classification, retryability, and model recovery behavior remain unchanged; diagnostic text is not model correction guidance.
Local model-request rejections
Use model.NewRequestValidationError(cause) only when application-side
validation rejects a model request before a provider accepts it. A remote model
adapter can restore this type from its service’s explicit request-validation
error. The cause is required: Unwrap() exposes the original error and Error()
returns its complete diagnostic. This type carries no provider name, HTTP
status, retry setting, or recovery instruction.
Do not use it for network, observer, cancellation, or provider failures.
Existing validators and adapters keep their behavior unless their owner
explicitly marks a local request rejection. Genuine provider rejections still
use model.ProviderError; invalid model or planner output keeps its separate
output-validation contract.
The run ends with kind model_request, Retryable: false, and no provider,
operation, provider code, or HTTP status. Its default summary is “The AI request
could not be prepared.” Applications can override hooks.PublicErrorModelRequest
at process startup; DebugMessage keeps the complete diagnostic. Tool execution
and nested agents do not turn this error into a retry or an extra model call for
correction. Text already published by an earlier model call in the same planner
activity is preserved before the run ends. A directly returned custom Temporal
ApplicationError retains the application’s existing classification and retry
policy; this type does not override that explicit outer error.
Temporal saves the error as goa_ai.request_validation_error, with
NonRetryable: true and the complete valid diagnostic in its message, without
details or a cause object. Readers reject a saved value that is retryable,
contains details or a cause, or has invalid UTF-8. Existing diagnostic encoding
and external failure-size limits still apply; there is no new text limit.
Previously saved provider, output, generic, and cancellation failures are not
reinterpreted.
Upgrade workers before adapters start producing this type. Older workers cannot read its saved classification and must not process histories containing it, including during rollback. Route those histories only to upgraded workers. There is no database migration, generated API change, or compatibility mode, and upgrading does not relabel earlier saved failures.
Saved error formats
New OutputContractFailure, ModelOutputRejected, and PlannerOutputRejected
records use ReasonVersion="goa_ai.rejection_reason.v2". Reason retains the
exact selected cause text, identified by ReasonSHA256 and ReasonSize.
Valid text has an empty ReasonOmitted; invalid UTF-8 produces an empty
Reason and ReasonOmitted="invalid_utf8". New records do not use
size_limit to omit a long reason.
New provider, generic, output-contract, and invalid-reserved Temporal failures use these four private application types:
goa_ai.provider_error.v3goa_ai.generic_error.v3goa_ai.output_contract_error.v3goa_ai.invalid_reserved_error.v3
The type selects the saved details format. Provider and generic details retain their owned text as plain strings, separately from the outer diagnostic message. Invalid UTF-8 becomes an explicit unavailable-text notice with the original hash and byte count, not silent replacement characters. These formats do not serialize arbitrary Go causes, SDK error objects, or custom application details. Exact valid text is not a promise to store arbitrary bytes.
Transport and application limits
The existing whole workflow argument/result limit still applies to the complete encoded value, including accompanying fields. A planner result that does not fit returns an explicit transport-budget failure; it does not persist the oversized diagnostic by silently shortening it.
Native Temporal failure objects use a separate SDK failure converter, not the
workflow argument/result admission path. Goa-AI adds no native failure-size
cap or preflight check. An application’s configured FailureConverter,
including its rejection behavior, remains application-owned.
Temporal request and history limits can reject large failures; pending activity retry state may retain a server-shortened failure. Instrumentation also has sampling, exporter, and backend limits. Preserving text in the framework does not guarantee unlimited storage, delivery, or recovery of text that was previously omitted.
Worker upgrades and saved histories
New readers retain the existing behavior for versionless and v1 rejection records and historical Temporal types with v1/v2 details. Decoding and replay do not rewrite those records, change their published bytes, or restore missing text. Old omission and validation rules still apply to the old formats.
Already stored terminal failures retain their original bytes. A workflow that reads old rejection metadata but closes for the first time after the upgrade writes the current terminal failure type. Successful replay does not prove that old and new terminal failure commands have identical encoded details.
Upgrade validating hook consumers and workflow/activity workers before they receive the new formats. Do not mix new writers with incompatible old readers on the same task queues; use the application’s verified worker-version routing or drain-and-replacement procedure. Rollback must retain readers that understand every format already written. Keep historical decoders while supported run records or workflow histories still require them; replacing workers alone does not remove that requirement.
Policies, Caps, and Labels
Design-Time RunPolicy
At design time, you configure per-agent policies with RunPolicy:
Agent("chat", "Conversational runner", func() {
RunPolicy(func() {
DefaultCaps(
MaxToolCalls(8),
MaxRecoveryTurns(3),
)
TimeBudget("2m")
InterruptsAllowed(true)
})
})
This becomes a runtime.RunPolicy attached to the agent’s registration:
Caps:
MaxToolCallsis the total budgeted tool calls per run.MaxRecoveryTurnslimits replacement planner calls after rejected tool or model output. A successful budgeted tool call starts a fresh recovery allowance. Tools declaredBookkeeping()consume neither budget. Model-authored batches stay atomic: bookkeeping calls add zero cost, but the runtime never removes individual calls to make a mixed batch fit. Successful bookkeeping results stay out of compact futureToolOutputs.Time budget:
TimeBudget– wall-clock budget for the run.FinalizerGrace(runtime-only) – optional reserved window for finalization.Interrupts:
InterruptsAllowed– opt-in for pause/resume.Missing fields behavior:
OnMissingFields– governs what happens when validation indicates missing fields.Terminal tools: Tools declared
TerminalRun()automatically become bookkeeping and complete the run once they succeed—no follow-upPlanResumeturn is scheduled. A terminal commit can therefore be admitted with no retrieval budget remaining. During forced finalization, the runtime admits only terminal bookkeeping calls, executes them inside the remaining hard-deadline window, and closes the run only if every terminal side effect succeeds. Before execution, the runtime writes the exactplanner.TerminationReasontoruntime.FinalizationReasonLabel(goa-ai.finalization_reason). Run labels, policy labels, planner output, and model output cannot choose or replace this value. Ordinary calls do not receive it.Consumers of fixed-limit or planner-authored terminal calls, including
tool_failure, useruntime.FinalizationReasonLabel. Deploy a change to this execution contract across consumers and runtime workers together.
Runtime Policy Overrides
In some environments you may want to tighten or relax policies without changing the design. The rt.OverridePolicy API allows process-local policy adjustments:
err := rt.OverridePolicy(chat.AgentID, runtime.RunPolicy{
MaxToolCalls: 3,
MaxRecoveryTurns: 1,
InterruptsAllowed: true,
})
Scope: Overrides are local to the current runtime instance and affect only subsequent runs. They do not persist across process restarts or propagate to other workers.
Overridable Fields:
| Field | Description |
|---|---|
MaxToolCalls | Maximum total tool calls per run |
MaxRecoveryTurns | Replacement planner calls after rejected tool or model output |
TimeBudget | Wall-clock budget for the run |
FinalizerGrace | Reserved window for finalization |
InterruptsAllowed | Enable pause/resume capability |
Only non-zero fields are applied (and InterruptsAllowed when true). This allows selective overrides without affecting other policy settings.
MaxRecoveryTurns counts only replacement attempts. If the allowance is
exhausted, the runtime may make one separate finalization call so the run can
return or persist a terminal outcome.
Recovering a Rejected Model Answer
A planner that rejects one completed model answer can return
planner.NewRecoverableModelOutputError. The error includes the rejected
FinalResponse and clear correction text. The workflow records the rejected
answer and its token usage, then spends one recovery turn on a replacement
answer with tools disabled. Ordinary OutputContractError values remain
terminal because they do not promise that another model call can correct the
output.
Use Cases:
- Temporary backoffs during provider throttling
- A/B testing different policy configurations
- Development/debugging with relaxed constraints
- Per-tenant policy customization at runtime
Labels and Policy Engines
Goa-AI integrates with pluggable policy engines via policy.Engine. Policies
receive tool metadata (IDs, tags), run context (SessionID, TurnID, labels), and
the structured ToolFailure after failed execution.
Labels flow into:
run.Context.Labels– available to planners during a run- tool activity input (
api.ToolInput.Labels) – cloned into dispatched tool executions so activities observe the same run-scoped metadata; terminal finalization calls also receive the runtime-owned reason underruntime.FinalizationReasonLabel - runtime records that cannot change after insertion (
storage.Store) – persisted with lifecycle changes for audit/search/dashboards (where indexed) - terminal completion and snapshots – the start labels come back out at the end of the run on
hooks.RunCompletedEvent.Labelsandrun.Snapshot.Labels, so completion hooks andGetRunSnapshotreaders recover run identity without out-of-band tracking
Per-Run Tool Filtering
Design-time tags and runtime options let callers narrow the tool surface before planner prompting and again before execution:
out, err := client.Run(ctx, "session-1", messages,
runtime.WithAllowedTags([]string{"read", "safe"}),
runtime.WithDeniedTags([]string{"destructive"}),
runtime.WithTagPolicyClauses([]runtime.TagPolicyClause{
{AllowedAny: []string{"docs", "search"}},
{DeniedAny: []string{"external"}},
}),
)
Use WithRestrictToTool when a repair flow should expose exactly one tool:
out, err := client.Run(ctx, "session-1", messages,
runtime.WithRestrictToTool(searchspecs.Search),
)
This is a run-wide caller policy. Tool failures use a separate contract:
ToolFailure.Recovery.Action selects correction, replanning, or finishing and the
runtime enforces the resulting tool catalog on the next planner turn.
Tool Execution
- Native toolsets: You write implementations; runtime handles decoding typed args using generated codecs
- Agent-as-tool: Generated agent-tool toolsets run provider agents as child runs (inline from the planner’s perspective) and adapt their
RunOutputinto aplanner.ToolResultwith aRunLinkhandle back to the child run - MCP toolsets: Runtime forwards canonical JSON to generated callers; callers handle transport
Tool payload defaults
Tool payload decoding follows Goa’s decode-body → transform pattern and applies Goa-style defaults deterministically for tool payloads.
See Tool Payload Defaults for the contract and codegen invariants.
Bounded tool results
Tools that return partial views of larger datasets should declare BoundedResult(...)
in the DSL. The runtime contract for those tools is:
- generated
tools.ToolSpec.Boundsdeclares the canonical bounded-result schema - successful executions must populate
planner.ToolResult.Bounds - the runtime projects provider-owned bounds into emitted
tool_resultJSON, result-hint template data under.Bounds, hook payloads, and stream events - for paged tools, provider code sets
Bounds.NextCursorto the opaque next-page cursor
tools.ToolSpec.Bounds uses model-facing JSON names. A DSL declaration may
refer to lower-camel Goa attributes such as NextCursor("nextCursor"), but
generated specs, schemas, runtime projection, and result codecs use
next_cursor.
Canonical projected fields:
returned(required)truncated(required)total(optional)refinement_hint(optional)next_cursor(optional, when declared viaNextCursor(...)and exposed by a self-pagingCursorcontract)
planner.ToolResult.Bounds remains the single machine-readable provider contract.
Authored Go result types stay semantic and domain-specific; they do not need to
duplicate the canonical bounded fields just so models can see them.
ContinueWith("continue_tool", "cursor") declares mechanical continuation as
a separate action. The runtime advertises that action only when result history
contains exactly one live successful chain head with another cursor. Exact
cursor lineage advances sequential pages. Parallel source invocations remain
valid, but multiple live heads make the no-argument action unavailable because
the contract exposes no result-chain selector. Its model-facing schema is an
empty object; when the model chooses the action with {}, the runtime binds the
cursor and any retained canonical query fields before execution.
Cursor("cursor") on a self-paging tool is the open contract: the runtime
projects Bounds.NextCursor into next_cursor, and the model repeats the
unchanged query arguments with that opaque cursor. Use it when pagination
itself requires a model-visible choice rather than mechanical continuation.
Generated result codecs accept the same canonical bounded fields projected by the runtime and reject unknown fields outside the semantic result plus those runtime-owned fields. This keeps method results, tool results, and transcript JSON aligned without handwritten schema walking.
For method-backed BindTo tools, the bound service method result still needs to
carry the canonical bounded fields so the generated executor can build
planner.ToolResult.Bounds before projection. Explicit tool-facing Return(...)
shapes must not duplicate those canonical fields. Within the bound method
result, only returned and truncated may be required; total,
refinement_hint, and next_cursor remain optional and are omitted from emitted
JSON whenever runtime bounds omit them.
When a service boundary must assemble canonical result JSON outside
ExecuteToolActivity, use runtime.EncodeCanonicalToolResult(...) rather than
calling the generated result codec and bounded-result projection helpers
separately.
Prompt Runtime Contracts
Prompt management is runtime-native and versioned:
runtime.PromptRegistrystores immutable baselineprompt.PromptSpecregistrations.runtime.WithPromptStore(prompt.Store)enables scoped override resolution (session->facility->org-> global).- Planners call
PlannerContext.RenderPrompt(ctx, id, data)to resolve and render prompt content. - Rendered content includes
prompt.PromptRefmetadata for provenance; planners can attach these tomodel.Request.PromptRefs.
messages, err := input.PrepareMessages()
if err != nil {
return nil, err
}
content, err := input.Agent.RenderPrompt(ctx, "assistant.system", map[string]any{
"AssistantName": "Ops Assistant",
})
if err != nil {
return nil, err
}
resp, err := modelClient.Complete(ctx, &model.Request{
Messages: messages,
PromptRefs: []prompt.PromptRef{content.Ref},
})
PromptRefs record which rendered prompt versions influenced a model request;
they are not provider wire payload fields. The runtime derives a run’s effective
prompt references by walking prompt_rendered and parent/child link records,
which cannot change after insertion. The store does not maintain a separate
prompt-reference list that could disagree with those records.
Rendering itself never writes runtime storage. Every path uses
prompt.RenderRecorder to create the same prompt.RenderEvent with the
resolved prompt ID, version, and scope:
- application code that renders initial messages passes
recorder.Events()throughruntime.WithRenderedPromptswith those messages; - planner activities return their recorded events with the planner result;
- child-agent prompt preparation runs in an activity and returns its rendered text and events in the child input;
RunOneShotrecords renders performed by its callback.
The workflow stores each accepted event as the same PromptRendered record.
The initial path is not a different rendering rule; it only delivers an event
that was created before workflow start. Child preparation runs in an activity
so Temporal replay reuses the text and events already in workflow history
instead of reading a newer prompt from storage.
RenderRecorder.Events returns completed renders in stable prompt ID, version,
session, and scope order. Concurrent completion order therefore cannot change
the exact workflow start request.
Memory, Streaming, Telemetry
Hook bus publishes structured hook events for the full agent lifecycle: run start/completion, phase changes,
prompt_rendered, tool scheduling/results/updates, planner notes and thinking blocks, awaits,ToolFailurerecovery directives, and agent-as-tool links.Memory stores (
memory.Store) subscribe and append durable memory events (user/assistant messages, tool calls, tool results, planner notes, thinking) per(agentID, RunID).The runtime store (
storage.Store) appends records that cannot change after insertion for eachRunID, for audit/debug UIs and run introspection. Lifecycle methods store their status, checkpoint, or cancellation change with the matching record in one operation.Stream sinks (
stream.Sink, for example Pulse or custom SSE/WebSocket) receive typedstream.Eventvalues produced by thestream.Subscriber. AStreamProfilecontrols which event kinds are emitted.The durable transcript preserves each selected provider response exactly. When an assistant message contains a tool call, its text stays in that transcript for provider replay but is not emitted as a user-facing assistant answer. Tool and await events present the nonterminal step. Only assistant messages without tool calls produce committed assistant-text events.
Telemetry: OTEL-aware logging, metrics, and tracing instrument workflows and activities end to end.
Tool Call Display Hints (DisplayHint)
Tool calls may carry a user-facing DisplayHint (for example for UIs).
Contract:
- Hook constructors do not render hints. Tool call scheduled events default to
DisplayHint=="". - The runtime enriches and persists a durable default call hint at publish time from the typed template when payload decoding succeeds.
- Tool registration requires a non-empty metadata title. When typed decoding fails or no template is registered, the runtime uses that title as the display hint. Malformed payloads still fail at the tool boundary; the metadata title only keeps the attempted work renderable. Hints are never rendered against raw JSON bytes.
- If a producer explicitly sets
DisplayHint(non-empty) before publishing the hook event, the runtime treats it as authoritative and does not overwrite it. - For per-consumer wording changes, configure
runtime.WithHintOverrideson the runtime. Overrides take precedence over DSL-authored templates for streamedtool_startevents.
Consuming a Session Stream (Pulse)
In production, the common pattern is:
- publish runtime stream events to Pulse (Redis Streams) using a
stream.Sink - subscribe to the session stream (
session/<session_id>) from your UI fan-out (SSE/WebSocket) - stop streaming a run when you observe
type=="run_stream_end"for the active run ID
import (
pulsestream "goa.design/goa-ai/features/stream/pulse"
"goa.design/goa-ai/runtime/agent/runtime"
"goa.design/goa-ai/runtime/agent/stream"
)
streams, err := pulsestream.NewRuntimeStreams(pulsestream.RuntimeStreamsOptions{
Client: pulseClient,
})
if err != nil {
panic(err)
}
rt := runtime.New(
runtimeStore,
runtime.WithEngine(eng),
runtime.WithStream(streams.Sink()),
)
sub, err := streams.NewSubscriber(pulsestream.SubscriberOptions{SinkName: "ui"})
if err != nil {
panic(err)
}
events, errs, cancel, err := sub.Subscribe(ctx, "session/session-123")
if err != nil {
panic(err)
}
defer cancel()
activeRunID := "run-123"
for {
select {
case evt, ok := <-events:
if !ok {
return
}
if evt.Type() == stream.EventRunStreamEnd && evt.RunID() == activeRunID {
return
}
// evt.SessionID(), evt.RunID(), evt.Type(), evt.Payload()
case err := <-errs:
panic(err)
}
}
Engine Abstraction
- In-memory: Fast dev loop, no external deps
- Temporal: Durable execution, replay, activity retries, signals, workers; adapters wire activities and context propagation
Goa-AI agent workflows are single-attempt. The runtime retries individual activities when their contracts allow it, but it never restarts an entire agent workflow after failure. A whole-workflow retry could repeat tool side effects or conflict with the final lifecycle record already saved by the first attempt.
Custom engines use engine/contract.NormalizeRootRequest and
NormalizeChildRequest to validate and privately retain accepted requests.
Before every initial or retry attempt, they call CopyRunInput so the workflow
handler receives a fresh input. After success they retain one private
CopyRunOutput result and copy it again for every wait, query, or other
caller-facing read. Goa-AI also converts search fields to values that every
supported engine can store and computes the identity used to recognize an exact
retry. Each engine adapter still translates and submits those values through
its own backend.
Semantic timing vs Temporal liveness
Goa-AI keeps the public runtime contract engine-agnostic:
RunPolicy.Timing.PlanandRunPolicy.Timing.Toolsare semantic attempt budgetsruntime.WithTiming(...)overrides those semantic budgets for a run- Generated clients use the agent’s default task queue. Pass
runtime.WithTaskQueue("orchestrator.chat")to oneStartorRuncall when that run must use a different queue
If you use the Temporal adapter and need queue-wait or liveness tuning, configure it on the Temporal engine itself:
eng, err := temporal.NewWorker(temporal.Options{
ClientOptions: &client.Options{
HostPort: "temporal:7233",
Namespace: "default",
},
WorkerOptions: temporal.WorkerOptions{
TaskQueue: "orchestrator.chat",
},
ActivityDefaults: temporal.ActivityDefaults{
Planner: temporal.ActivityTimeoutDefaults{
QueueWaitTimeout: 30 * time.Second,
LivenessTimeout: 20 * time.Second,
},
Tool: temporal.ActivityTimeoutDefaults{
QueueWaitTimeout: 2 * time.Minute,
LivenessTimeout: 20 * time.Second,
},
},
})
if err != nil {
panic(err)
}
This split keeps workflow mechanics behind the Temporal boundary while the generic runtime stays honest across both Temporal and the in-memory engine.
Storage and completion adapter contracts
The runtime registers one typed activity named runtime.store. Every
StorageActivityCommand sets exactly one of Append, RootStart,
ChildStart, OneShotStart, OneShotChildStart, Cancellation, Suspension,
or Terminal.
The returned StorageActivityResult sets exactly the matching field and no
other field. Custom stores return storage.ContractError when repeating the
same command cannot succeed. Temporary database and network failures remain
ordinary errors and can retry. runtime.WithStorageActivityTimeout sets the
activity’s Start-to-Close timeout and requires a value greater than zero.
Engine.QueryRunCompletion returns the current run Status. Once the run is
closed, the same result also contains its stable CompletedAt time and final
Output or WorkflowError. EnsureRunCompletion uses CompletedAt as the
record timestamp, so every retry submits the same value. The method’s separate
error means the engine could not retrieve those facts. There is no separate
status query.
Child prompt preparation returns exactly one Success or Failure. Success
contains only the messages and rendered prompt facts. The workflow derives the
child run, session, parent, tool, and label identity from the original recorded
tool call. The in-memory engine copies and bounds the input and output and
applies the same retry policy as Temporal.
Run Contracts
SessionIDis required for sessionful starts.StartandRunfail fast whenSessionIDis empty or whitespaceStartOneShotandOneShotRunare explicitly sessionless. They do not require or create a session and do not emit session-scoped stream events- The host creates sessionful sessions before submitting work. Agent runtimes do not create, end, or delete sessions
- The workflow engine accepts a root workflow before its first activity records the run. The runtime creates no
pendingrow before engine admission - Repeating a start with the same run ID and the exact same submitted request returns the accepted workflow while the engine can still query that workflow history. Reusing the ID with different input is rejected. After the engine’s history-retention window expires, Goa-AI no longer promises exact-ID lookup; a product that needs permanent command identity must store it in its own service
- Every accepted workflow has one
RunStartedrecord. An ended session does not erase that accepted start; it adds a canceledRunCompletedimmediately afterward - Root, child, sessionless root, and sessionless child starts are separate storage calls. Child starts store the parent link with the child start; sessionless starts store full run metadata without a session
- Temporal terminates a child workflow if its parent workflow closes first
- New children require a running parent.
StartChildRunandStartOneShotChildRuneach store the parent link and child start together. An exact retry of an accepted start remains valid after the parent stops, while a changed retry or a new child is rejected - Cancellation reasons are write-once. An exact retry succeeds; a different reason for the same run is a conflict
- Suspension and terminal changes store the new status with the matching record, which cannot later be changed
- Durable
RunStarted,RunSuspended,RunCompleted, andChildRunLinkedpayloads must contain exactly one matching JSON value. Unknown fields and trailing JSON values are rejected - Agents must be registered before the first run. The runtime rejects registration after the first run submission with
ErrRegistrationClosedto keep engine workers deterministic - Tool executors receive explicit per-call metadata (
ToolCallMeta) rather than fishing values fromcontext.Context. Its labels contain cloned run and policy labels plusruntime.FinalizationReasonLabelonly when that call is executing terminal finalization - Do not rely on implicit fallbacks; all domain identifiers (run, session, turn, correlation) must be passed explicitly
Ensuring a final record and its delivery
Normal workflows retry suspension and terminal writes until the runtime store accepts them. A host can use two explicit commands after engine history closes:
Runtime.EnsureRunCompletion(ctx, runID)stores a missing suspension or terminal result for a run that is still active in storage. If the run is already closed, or another final result wins while the command runs, it validates and delivers the exact stored result instead.Runtime.EnsureChildRunLink(ctx, runID)validates and delivers only a session-backed child’s exact stored parent link. Hosts can call it in parent-first order before they deliver the final results of nested children.
EnsureRunCompletion delivers the parent link before a child’s final event.
Stable event keys make repeated stream delivery safe, and an already stored
result does not produce another local lifecycle notification. Neither ensure
command changes the result that storage already accepted.
Both commands require Runtime.WithStream when the Session status used for
delivery is active. EnsureChildRunLink reads the current status through
LoadSessionStatus. EnsureRunCompletion instead uses the SessionStatus
returned with the final-record write or exact retry. A newly checked ended
Session keeps its stored records and suppresses delivery. If storage accepted
the event while the Session was active, that event remains due: ending the
Session while the same delivery call retries does not cancel it.
EnsureRunCompletion returns ErrRunCompletionNotReady when the engine still
reports a running workflow. It returns ErrRunCompletionCorrupt when engine
history or stored lifecycle data cannot form one valid result. An error loading
engine history is returned to the caller and is never stored as the workflow’s
failure.
Run listing and snapshot methods are read-only and never call either command. These commands add no database schema migration and do not change a public wire format. They do change the Go source interface for custom stores, and existing durable lifecycle records must match the strict typed JSON shapes documented in Memory & Sessions.
External Input and Workflow Continuations
Each accepted user input starts one top-level workflow for that turn. The workflow ends with either that turn’s final result or an external-input suspension. Nested agents still run as linked child workflows.
Clarifications, structured questions, external tool results, and confirmations
end the current workflow successfully. The returned RunOutput.Suspension
contains visible Pending requests and a private Checkpoint. The application
keeps the complete suspension in trusted server storage and sends only
Suspension.Pending to the UI or external system that must answer it. No
Temporal workflow remains open while a person is deciding.
Before the workflow completes, Goa-AI stores its private checkpoint under the completed run ID. When the application must record acceptance of the answer in its own storage, continuation has three explicit steps:
PrepareContinuationloads the saved checkpoint, validates the typed answer and the current generatedAgentDefinition, and returns aPreparedRuncontaining an immutable copy of the complete successor request. It does not write runtime state or call the workflow engine.- The application calls
MarshalBinaryand atomically stores those bytes with its accepted answer. - Any process can load those bytes, call
ParsePreparedRun, and pass the restored value toStartPrepared.
This lets concurrent requests compete for one application-owned obligation without starting two workflows. If the engine response is uncertain, the application retries the same prepared value; it does not prepare a replacement with different fields.
prepared, err := client.PrepareContinuation(
ctx,
"session-1",
previous.RunID,
"run-124",
"turn-2",
&api.PendingInputResponse{
Clarification: &api.ClarificationAnswer{
ID: "clarify-device",
Answer: "Device ID is ABC-123",
},
},
runtime.WorkflowOptions{},
)
if err != nil {
return err
}
preparedBytes, err := prepared.MarshalBinary()
if err != nil {
return err
}
// This application-owned method uses one database transaction to accept the
// answer and store the prepared workflow ID with preparedBytes.
if err := workflowStarts.AcceptContinuation(ctx, previous.RunID, prepared.RunID(), preparedBytes); err != nil {
return err
}
preparedBytes, err = workflowStarts.Load(ctx, "run-124")
if err != nil {
return err
}
stored, err := runtime.ParsePreparedRun(preparedBytes)
if err != nil {
return err
}
handle, err := client.StartPrepared(ctx, stored)
if err != nil {
return err
}
next, err := handle.Wait(ctx)
When no application write is needed between validation and engine submission,
Continue is the convenience method that prepares, starts, and waits. In both
forms the runtime store keeps the authoritative checkpoint. The complete
suspension and prepared bytes may also contain that checkpoint and the complete
transcript, so keep them in trusted, access-controlled application storage and
never send them to an untrusted client. Goa-AI validates the checkpoint’s
version and pending request, restores saved payloads through the current
generated codecs, and resumes planning.
The accepted checkpoint format is exactly goa-ai.run-suspension.v8. Version
eight saves the advertised tool names when an accepted recovery plan waits for
input: failed tool names alone cannot reconstruct the other choices offered
on that turn. Continuation preserves those choices while still checking the
current agent definition and execution policy.
Goa-AI rejects every earlier checkpoint version. Before upgrading, finish old-format saved work under its owning runtime. If any must remain unfinished, the host must explicitly decide how to preserve it and whether it will remain resumable; this runtime cannot resume it. The framework supplies no conversion command and never automatically deletes, cancels, or rewrites saved work.
When an answer completes a model-authored tool call from the earlier workflow,
the new tool_end event has two distinct run identities:
- its normal run ID names the new workflow that received the answer; and
call_run_idnames the earlier workflow that emitted the matchingtool_start.
Stream consumers must pair those events using call_run_id and the tool call
ID. They must not search previous runs or assume the call and result belong to
the same workflow. See Transparent rollouts
for the worker and deployment requirements that preserve this boundary during
a release.
Tool Confirmation
Goa-AI supports runtime-enforced confirmation gates for sensitive tools (writes, deletes, commands).
You can enable confirmation in two ways:
- Design-time (common case): declare
Confirmation(...)inside the tool DSL. Codegen stores the policy intools.ToolSpec.Confirmation. - Runtime (override/dynamic): pass
runtime.WithToolConfirmation(...)when constructing the runtime to require confirmation for additional tools or override design-time behavior.
At execution time, the workflow emits a confirmation request and completes with a suspension. The accepted decision starts a new workflow. That continuation executes the tool only when approved. When denied, the runtime synthesizes a schema-compliant tool result so the transcript remains valid and the planner can react deterministically.
Confirmation protocol
At runtime, confirmation is implemented as a dedicated await/decision protocol:
Await payload (streamed as
await_confirmation):{ "id": "...", "title": "...", "prompt": "...", "tool_name": "facility.commands.change_setpoint", "tool_call_id": "toolcall-1", "payload": { "...": "canonical tool arguments (JSON)" } }
Contract:
payloadalways contains the canonical JSON tool arguments for the pending call. If approved, those are the arguments the runtime executes.Confirmation overrides may customize the prompt and denied-result rendering, but they do not introduce a separate display-payload channel or change the meaning of
payload.Products that need a richer confirmation UI should materialize it in the application layer from the canonical payload plus application-owned reads.
Continuation response:
response := &api.PendingInputResponse{ Confirmation: &api.ConfirmationDecision{ ID: "await-1", Approved: true, // or false RequestedBy: "user:123", Labels: map[string]string{"source": "front-ui"}, Metadata: map[string]any{"ticket_id": "INC-42"}, }, }
Tool authorization events
When a decision is provided, the runtime emits a first-class authorization event:
- Hook event:
hooks.ToolAuthorization - Stream event type:
tool_authorization
This event is the canonical “who/when/what” record for a confirmed tool call:
tool_name,tool_call_idapproved(true/false)summary(deterministic runtime-rendered summary)approved_by(copied fromapi.ConfirmationDecision.RequestedBy, intended to be a stable principal identifier)
The event is emitted immediately after the decision is received (before tool execution when approved, and before the denied tool result is synthesized when denied).
Notes:
- Consumers should treat confirmation as a runtime protocol:
- Render the first pending item when its kind is
confirmation, then submit the decision throughContinue. - Do not couple UI behavior to a specific confirmation tool name; treat it as an internal transport detail.
- Render the first pending item when its kind is
- Confirmation templates (
PromptTemplateandDeniedResultTemplate) are Gotext/templatestrings executed withmissingkey=error. In addition to the standard template functions (e.g.printf), Goa-AI provides:json v→ JSON encodesv(useful for optional pointer fields or embedding structured values).quote s→ returns a Go-escaped quoted string (likefmt.Sprintf("%q", s)).
Runtime validation
The runtime validates confirmation interactions at the boundary:
- The confirmation
IDmatches the pending await identifier when provided. - The continuation contains exactly one response variant and a well-formed decision.
Planner Contract
Planners implement:
type Planner interface {
PlanStart(ctx context.Context, input *planner.PlanInput) (*planner.PlanResult, error)
PlanResume(ctx context.Context, input *planner.PlanResumeInput) (*planner.PlanResult, error)
}
PlanResult contains tool requests, a final response, a final tool result,
annotations, and the selected post-tool transition. PlanResumeInput tells the
planner why it is being called.
Planner-authored requests contain domain intent only. Use
planner.NewToolRequest(typedTool, payload) to encode one. When forwarding a
validated provider call, use planner.ToolRequestFromModelCall(call) so the
provider correlation ID is preserved without becoming the runtime execution
ID. The runtime validates the complete plan before it assigns execution IDs or
publishes tool events.
These contracts are separate:
| Contract | Scope | Plain-English meaning |
|---|---|---|
ToolSpec.Tags | A tool, for every run | Flat labels available to generic policy and UI filtering. |
ToolSpec.Meta | A tool, for every run | Inert generated annotations whose semantics belong to the named consumer; metadata alone changes no runtime behavior. |
ToolSpec.Bookkeeping | A tool, for every run | The call is a durable control record whose success does not require another planner turn. It consumes no retrieval or consecutive-failure budget. |
ToolSpec.TerminalRun | A tool, for every run | Successful execution itself ends the run. It automatically implies bookkeeping. |
ToolFailure.Recovery.Action | One failed result | Selects correction with the failed tool still available, replanning without it, or finalization. |
PlanResult.SynthesizeAfterTools | One selected batch | If the batch has no recoverable failure, the next planner turn must answer. |
PlanResumeInput.SynthesisOnly | One planner activity | Return a final answer; tool calls are invalid. |
PlanResumeInput.Finalize | Runtime-forced termination | A cap or deadline has prohibited normal work. |
The runtime chooses one next state in this order:
| Completed step | Next state |
|---|---|
| A cap or deadline requires finalization | Finalize turn |
A successful TerminalRun tool completed | End immediately |
Any failed result has AllowsToolTurn() == true | Normal repair turn |
SynthesizeAfterTools is true | SynthesisOnly turn |
| Otherwise | Normal continuation turn |
This keeps planner intent from becoming a second retry policy. A recoverable failure is repaired first; a successful or terminally failed final batch proceeds to synthesis. The runtime rejects tool calls returned from a SynthesisOnly turn.
Each recoverable ToolFailure also selects a Recovery.Action:
correct_callkeeps the failed tool available and gives the next planner turn the original model-authored input, generated validation issues, field guidance, and example. It does not require one replacement call per failure. The planner may combine work, make any number of valid calls to advertised tools, wait for input, or answer from the evidence already collected.replanremoves the failed tool from the next planner turn. The planner may use another advertised tool, wait for input, or answer.finishremoves all tools and requires a final answer from the available evidence.
An ordinary correct_call turn combines the current agent’s executable tools
with the exact failed-tool contracts. Matching names are deduplicated;
conflicting contracts, missing executable registrations, and revoked tools
fail before a model call. Caller restrictions, run tag restrictions, and
recovery exclusions still apply. A denied correction tool causes an error; the
runtime does not silently drop it or restore it after filtering. Downstream
authorization still checks every executed call.
Unfinished queries retain their runtime-generated continuation actions; failed requests do not create continuations. Forced finalization offers only the exact failed terminal tool for correction, and synthesis-only turns remain tool-free. After correction, normal turns return to the current agent’s tools.
The workflow owns model-facing correction evidence. It replaces executor-
supplied prior input and examples with the original provider call and the
registered tool specification before the failure enters run history. A
runtime-created continuation has no model-authored input, so it cannot request
correct_call; this prevents private cursors or injected execution fields from
appearing in a later model request. Model transcripts correlate results with
ModelToolCallID, while activities, retries, and stored execution records use
the separate runtime ToolCallID.
The runtime records the exact tool catalog shown on a recovery turn and rejects every executable call outside it, including a call embedded in a request for user or external input. Generated codecs still validate every payload, and the run’s tool, failure, and time limits still stop repeated invalid work. If a recovery turn waits for input, its failure evidence remains available when the run resumes; choosing a tool call or final answer clears that evidence.
Recovery activity inputs and their advertised catalog are part of durable workflow history. Production deployments must use pinned Temporal Worker Deployment Versioning and retain each old worker version until Temporal reports it drained. Starting a new worker does not make it safe to replay an existing workflow on new code. A continuation is a new workflow and may use the current version after its saved checkpoint passes validation.
When PlanResumeInput.Finalize is set, planners may return terminal bookkeeping tools; those calls are not replayed into a later planner turn and must durably finish finalization.
Planners also receive a PlannerContext via input.Agent that exposes runtime services:
AdvertisedToolDefinitions()- get the runtime-filtered tool definitions visible to the model for this turnModelClient(id string)- get a raw provider-agnostic model clientPlannerModelClient(id string)- get a planner-scoped model client with runtime-owned event emissionRenderPrompt(ctx, id, data)- resolve and render prompt content for the current run scopeAddReminder(r reminder.Reminder)- register run-scoped system remindersRemoveReminder(id string)- clear reminders when preconditions no longer holdMemory()- access conversation history
Preparing conversation messages
Both planner inputs provide PrepareMessages func() ([]*model.Message, error)
instead of a Messages field. Call it before inspecting or transforming the
conversation, including prompt construction and message-dependent checks:
messages, err := input.PrepareMessages()
if err != nil {
return nil, err
}
The first call applies the registered history policy to this activity’s messages and advertised tools, using the activity context and deadline. Repeated or concurrent calls return the same prepared slice, message pointers, and error without rerunning the policy. Callers must coordinate mutations as with any shared slice. Complete all calls before the planner returns, and do not retain the function for another invocation. Each new activity attempt gets its own preparation; completed-activity replay does not prepare history again.
A decision based only on run context and typed tool outputs need not call it, so unused history causes no token counting or summarization. A code-only decision that reads earlier messages must still prepare them. The runtime always supplies the callback, including when no history policy is configured; there is no raw-history alternative or nil-as-bypass mode.
Preparation failure fails the activity even if planner code ignores the error. The original history error and its retry classification take precedence over later planner or model-output errors; they do not become model-output recovery. Cancellation reaches token counting and summarization through the activity context. Repeated access does not retry a failed preparation. History limits, retention rules, native tool/result pairing, and stored transcripts are unchanged.
Upgrade custom planners and direct test fixtures when upgrading goa-ai: replace reads of the removed input fields with the checked call above and supply the callback in fixtures whose planners read history. This is a Go source change, not a new wire format or stored-history migration.
Feature Modules
runtime/mcp– MCP callers for HTTP and stdio transports; HTTP accepts JSON and event-stream responsesruntime/agent/storage/inmem– integrated in-memory runtime store for local examples and testsfeatures/memory/mongo– durable memory storefeatures/prompt/mongo– Mongo-backed prompt override storefeatures/stream/pulse– Pulse sink/subscriber helpersfeatures/model/{anthropic,bedrock,openai,vertex}– provider adapters that return validated model clientsfeatures/model/gateway– remote provider server and validated transport clientsfeatures/model/middleware– provider middleware installed beneath client validation, including exact-token adaptive rate limitingfeatures/policy/basic– simple policy engine with allow/block lists andToolFailurehandling
Model Client Throughput & Rate Limiting
Goa-AI ships an adaptive input-token limiter under
features/model/middleware. It asks the wrapped client for the exact request
token count, reserves that capacity before the call, and adjusts its effective
input-tokens-per-minute budget when providers report throttling.
import (
"github.com/aws/aws-sdk-go-v2/service/bedrockruntime"
"goa.design/goa-ai/runtime/agent/runtime"
"goa.design/goa-ai/features/model/bedrock"
mdlmw "goa.design/goa-ai/features/model/middleware"
)
awsClient := bedrockruntime.NewFromConfig(cfg)
bed, err := bedrock.New(awsClient, bedrock.Options{
DefaultModel: "us.anthropic.claude-4-5-sonnet-20251120-v1:0",
})
if err != nil {
panic(err)
}
rl := mdlmw.NewAdaptiveRateLimiter(
ctx,
throughputMap, // *rmap.Map joined earlier (nil for process-local)
"bedrock:sonnet", // key for this model family
80_000, // initial input tokens per minute
1_000_000, // maximum input tokens per minute
)
limited, err := rl.Middleware()(bed)
if err != nil {
panic(err)
}
rt := runtime.New(runtimeStore)
if err := rt.RegisterModel("bedrock", limited); err != nil {
panic(err)
}
Middleware construction does not test token-count support. If the selected
provider or request cannot be counted exactly, the first Complete or Stream
call returns model.ErrTokenCountingUnsupported before inference. Vertex
Gemini supports exact counting; Bedrock supports it only for requests and
models accepted by Runtime CountTokens. OpenAI has no native counter.
The limiter meters input tokens only. A unary success or clean stream end probes upward; a unary or terminal streaming rate-limit error backs off. Merely opening or closing a stream does not count as success.
LLM Integration
Goa-AI planners interact with large language models through a provider-agnostic interface. This design lets you swap providers—AWS Bedrock, OpenAI, Google Vertex AI (Gemini and Claude-on-Vertex), or custom endpoints—without changing your planner code.
The Validated Model Client
All planner interactions go through an opaque model.Client:
resp, err := client.Complete(ctx, req)
stream, err := client.Stream(ctx, req) // *model.ValidatedStream
Provider integrations implement model.Provider, which produces raw transport
responses and chunks. Goa-AI constructs model.Client with
model.NewClient(provider) and validates requests and complete responses around
that provider. External packages cannot implement model.Client or expose raw
provider chunks to planners.
Before a provider call, the client validates tool names and schemas, message parts, thinking options, structured-output metadata, and the request’s dynamic values. Requests and unary responses are limited to 16 MiB and 100,000 visited values. Nested dynamic metadata is limited to depth 64. Streaming applies one cumulative budget across chunks and the terminal response. These limits reject the whole operation; Goa-AI never truncates, repairs, or coerces model data.
ValidatedStream must be drained to io.EOF. Only then does Response()
return the accepted canonical response. An incomplete, malformed, or
contradictory stream returns an error and no accepted response.
Complete tool calls must satisfy their advertised schema and any attached generated
decoder before planner code can observe them. When a schema rejection has one
unique deepest cause matching generated array field metadata, correction
guidance names that array’s inclusive minimum or maximum, for example:
Field "items" must contain at most 3 items. Ambiguous or unsupported causes
retain generic guidance. These are bounds on that one submitted array, not a
limit on the complete run. The arguments remain rejected before execution;
Goa-AI does not split, truncate, or rewrite them. The original validation error,
acceptance rules, and configured recovery-turn limit remain unchanged.
Provider Adapters
Goa-AI ships with adapters for popular LLM providers:
AWS Bedrock
import (
"github.com/aws/aws-sdk-go-v2/service/bedrockruntime"
"goa.design/goa-ai/features/model/bedrock"
)
awsClient := bedrockruntime.NewFromConfig(cfg)
modelClient, err := bedrock.New(awsClient, bedrock.Options{
DefaultModel: "anthropic.claude-3-5-sonnet-20241022-v2:0",
HighModel: "anthropic.claude-sonnet-4-20250514-v1:0",
SmallModel: "anthropic.claude-3-5-haiku-20241022-v1:0",
MaxTokens: 4096,
Temperature: 0.7,
})
if err != nil {
panic(err)
}
OpenAI
import (
"os"
"goa.design/goa-ai/runtime/agent/runtime"
)
rt := runtime.New(runtimeStore) // host-owned runtime storage
modelClient, err := rt.NewOpenAIModelClient(runtime.OpenAIConfig{
APIKey: os.Getenv("OPENAI_API_KEY"),
DefaultModel: "gpt-5-mini",
HighModel: "gpt-5",
SmallModel: "gpt-5-nano",
})
if err != nil {
panic(err)
}
Google Vertex AI (Gemini and Claude-on-Vertex)
The features/model/vertex package ships two constructors that both satisfy
model.Client: a native Gemini adapter, and a pure-construction helper that
points the Anthropic adapter at Claude models hosted on Vertex.
import "goa.design/goa-ai/runtime/agent/runtime"
// Gemini on Vertex, using Application Default Credentials.
geminiClient, err := rt.NewVertexGeminiModelClient(ctx, runtime.VertexConfig{
ProjectID: "my-gcp-project",
Location: "us-central1",
DefaultModel: "gemini-2.5-flash",
HighModel: "gemini-3-pro-preview",
SmallModel: "gemini-2.5-flash-lite",
MaxTokens: 4096,
ThinkingBudget: 10000,
})
// Claude on Vertex. This is pure construction: it builds an Anthropic SDK
// client against the SDK's Vertex transport and hands it to
// features/model/anthropic, which owns Messages translation and
// HTTP-status error classification for every Anthropic-hosted adapter
// (direct API and Vertex-hosted alike) — no separate translation layer.
claudeOnVertexClient, err := rt.NewVertexAnthropicModelClient(ctx, runtime.VertexConfig{
ProjectID: "my-gcp-project",
Location: "us-east5",
DefaultModel: "claude-sonnet-4-5@20250929",
})
Gemini 3-class models attach an opaque thought signature to functionCall
parts (not only to thought/thinking parts) to authenticate the reasoning chain
behind a tool call. The Vertex adapter round-trips this signature through
model.ToolCall.ThoughtSignature / model.ToolUsePart.ThoughtSignature using
the same base64 convention as ThinkingPart.Signature. The runtime captures
this signature at the model-client boundary — before either integration style
below ever produces a planner.ToolRequest — and reattaches it by tool-call
ID when rebuilding the provider transcript. planner.ToolRequest never
carries a signature field; planner code does not need to know signatures
exist.
Provider Capability Differences
The shared request type is provider-neutral, but adapters reject combinations their APIs cannot preserve:
| Provider | Contract enforced before or during the call |
|---|---|
| OpenAI | Structured output uses strict schema projection. Tools and structured output cannot be combined. Overlapping oneOf branches are rejected rather than widened. Strict schemas are bounded to 5,000 properties, 1,000 enum values, 10 object levels, and 120,000 aggregate name/enum characters; fine-tuned models reject additional unsupported keywords. Thinking requests reject temperature. |
| Anthropic | Current Claude models use native structured output where supported. Adaptive thinking permits tools and normal forced choice, while older manual thinking rejects forced any or named-tool choice. Current-generation Claude models omit deprecated sampling parameters. Streams must close every content block and report a stop reason. |
| Bedrock | Claude 4.5 and 4.6 use native OutputConfig; other Claude models use one private forced tool and validate its result against the same contract. Event-stream exceptions retain their provider error kind. Runtime CountTokens rejects models that require the separate Mantle endpoint and cannot count structured-output requests. |
| Vertex Gemini | Gemini 3 uses thinking levels, rejects numeric thinking budgets and explicit thinking disablement, and forwards API-valid temperatures. Tool-call thought signatures are retained and replayed by the runtime. Streams require exactly one candidate and a finish reason. |
For Claude Opus 4.7+, Sonnet 5+, Haiku 5+, Fable, and Mythos, the Anthropic and
Bedrock adapters omit temperature, top_p, and top_k because those models
reject the parameters. Older Claude generations continue to receive configured
sampling values.
Provider adapters are stateless with respect to conversation history. Every
request must carry the complete provider-ready Messages transcript; a
RunID does not ask an adapter to load prior messages.
Common sentinel errors include:
model.ErrStructuredOutputUnsupportedwhen an adapter cannot express the requested output contractmodel.ErrTokenCountingUnsupportedwhen exact provider counting is unavailablemodel.ErrEmptyStreamwhen a provider closes without model outputmodel.ErrRateLimitedfor retryable provider throttling
*planner.OutputContractError is a structured error, not a sentinel. Detect it
with errors.As and inspect its origin to distinguish invalid model, planner,
or tool output. It is non-retryable because another request must not hide a
contract violation.
Canonical Message Metadata and Citation Replay
model.Message.Meta contains provider-authored data needed to replay a
response exactly. Boundaries that persist or transport metadata should use
model.MarshalMetadata and model.UnmarshalMetadata. These codecs require one
JSON object, preserve decoded numbers as json.Number, reject trailing data,
and canonicalize nil or an empty object to nil.
Citation replay is provider-specific and must never flatten citations into
ordinary text. The Bedrock adapter can replay assistant CitationsPart values
as native citation content blocks, preserving source identity, excerpts, and
document character, chunk, or page locations. Bedrock system citations remain
unsupported because its system-content union has no citation member. Anthropic
and Vertex reject citation replay when the canonical part lacks fields their
provider protocol requires.
Using Model Clients in Planners
Planners obtain model clients through the runtime’s PlannerContext. There are
two explicit integration styles:
PlannerModelClient(id)for planner-scoped streaming with runtime-owned event emissionModelClient(id)when you need direct validated model access and will drain the returned stream withplanner.ConsumeStream
PlannerModelClient (Recommended)
PlannerContext.PlannerModelClient(id) returns a planner-scoped client that
owns AssistantChunk, PlannerThinkingBlock, and UsageDelta emission. Its
Stream(...) method drains the underlying provider stream and returns a
planner.StreamSummary. PlannerModelClient permits exactly one Complete or
Stream invocation in that planner turn:
func (p *MyPlanner) PlanStart(ctx context.Context, input *planner.PlanInput) (*planner.PlanResult, error) {
mc, ok := input.Agent.PlannerModelClient("anthropic.claude-3-5-sonnet-20241022-v2:0")
if !ok {
return nil, errors.New("model not configured")
}
messages, err := input.PrepareMessages()
if err != nil {
return nil, err
}
req := &model.Request{
Messages: messages,
Tools: input.Agent.AdvertisedToolDefinitions(),
Stream: true,
}
sum, err := mc.Stream(ctx, req)
if err != nil {
return nil, err
}
if len(sum.ToolCalls) > 0 {
return &planner.PlanResult{ToolCalls: sum.ToolCalls}, nil
}
final := sum.FinalResponse()
if final == nil {
return nil, errors.New("model stream ended without a canonical response")
}
return &planner.PlanResult{
FinalResponse: final,
Streamed: true, // Assistant text was already streamed
}, nil
}
This is the simplest integration style because the planner-scoped client drains
and summarizes the validated stream itself. Returning sum.FinalResponse()
also selects the exact
provider response captured for that invocation; rebuilding a text-only message
would discard thinking, citations, signatures, metadata, and message boundaries.
Validated Client + ConsumeStream
When you need direct model.Client access, fetch it from
PlannerContext.ModelClient and pair its validated stream with
planner.ConsumeStream:
mc, ok := input.Agent.ModelClient("anthropic.claude-3-5-sonnet-20241022-v2:0")
if !ok {
return nil, errors.New("model not configured")
}
messages, err := input.PrepareMessages()
if err != nil {
return nil, err
}
req := &model.Request{
Messages: messages,
Tools: input.Agent.AdvertisedToolDefinitions(),
Stream: true,
}
stream, err := mc.Stream(ctx, req)
if err != nil {
return nil, err
}
sum, err := planner.ConsumeStream(ctx, stream)
if err != nil {
return nil, err
}
if len(sum.ToolCalls) > 0 {
return &planner.PlanResult{ToolCalls: sum.ToolCalls}, nil
}
final := sum.FinalResponse()
if final == nil {
return nil, errors.New("model stream ended without a canonical response")
}
return &planner.PlanResult{
FinalResponse: final,
Streamed: true,
}, nil
This helper only drains the stream and returns a StreamSummary with
accumulated text and tool calls. The runtime model-invocation journal publishes
accepted presentation and usage events later.
Generated tool definitions use model.ToolDefinitionFromSpec, which retains
the generated payload decoder. Caller-authored tools use
model.AdvertisedToolInputFromSchema. Both paths reject unknown tools and
invalid payloads before planner code receives a provider tool call.
Use the direct client path when planner logic needs to inspect validated preview
chunks or make multiple model calls in one planner turn. Drain every selected
stream to its terminal result; closing early does not produce an accepted
response. The returned PlanResult must forward one exact selected result:
either that summary’s complete ToolCalls set or its FinalResponse(). The
runtime rejects modified, mixed, or ambiguous results. Do not mix
PlannerModelClient.Stream(...) with planner.ConsumeStream; choose one stream
owner per planner turn.
Remote Model Gateways
features/model/gateway carries model requests to a separately deployed
provider process without weakening validation. The server operates on a raw
model.Provider, so provider-side middleware runs before the transport:
server, err := gateway.NewServer(
gateway.WithProvider(provider),
gateway.WithUnary(unaryMiddleware...),
gateway.WithStream(streamMiddleware...),
)
The consumer constructs a validated client from its transport functions:
client, err := gateway.NewRemoteClient(completeRemote, streamRemote)
countingClient, err := gateway.NewCountingRemoteClient(
completeRemote,
streamRemote,
countRemote,
)
Use NewCountingRemoteClient only when the remote endpoint implements exact
token counting. NewRemoteClient deliberately returns
model.ErrTokenCountingUnsupported for count requests instead of estimating.
History Policies
History policies run on the first PrepareMessages call in each planner
activity, not before every planner invocation. See
Preparing conversation messages for the
caller contract and the decisions that need no history work.
History compression separates the condition that starts summarization from the amount of exact recent history retained:
CompressAtTurnsandCompressAtMaxInputTokensare ORed triggers.KeepMaxTurnsandKeepMaxInputTokensboth bound which newest complete turns are eligible to remain exact; the runtime never cuts a turn in half.- Token policies require a
HistoryModelwhose client’sCountTokensoperation returns an exact count. The count includes preserved system messages, candidate turns, and currently advertised tools. CompressAtMaxInputTokensis exclusive: a request exactly at the threshold fits; only a larger request triggers compression.
Bedrock Runtime cannot count structured-output requests. Claude Opus 4.7,
Sonnet 5, and Mythos 5 require AWS’s separate Mantle count endpoint, so the
Bedrock adapter returns model.ErrTokenCountingUnsupported for those models.
The generated agent config exposes HistoryCompression for deployment-specific
overrides without changing the design defaults.
Exact retention and summary coverage
KeepMaxTurns and KeepMaxInputTokens establish the eligible complete turns
before summarization. The newest complete turn is mandatory. The older-turn
token allowance measures each candidate request minus the request containing
only newest; both counts include preserved system messages and the current tool
catalog. Equality fits. The summary is not charged against this older-turn
allowance: it counts against the separate CompressAtMaxInputTokens total limit.
When that total limit is positive, one summary receives every turn older than newest, including optional turns eligible to remain exact, even when the turn trigger fires first. The runtime counts the preserved system messages, actual rendered summary, eligible exact turns, and unchanged tool catalog together. If too large, it removes the oldest optional whole turn and counts the next candidate. It returns the longest fitting suffix permitted by the initial eligibility calculation, without assuming that counts are additive or monotonic. Equality with the limit fits. The count covers this history-policy request, not thinking or structured output a planner may choose afterward; the limit is not a summary-model context maximum.
For example, if three turns are eligible, it tests the summary plus all three, then plus the newest two, then plus newest alone, stopping at the first fit. Every removed turn has already reached the summary model. Newest is never summarized or split. Some older turns may occur in both summary prose and exact history; their tools are not executed again. However, the answering model still has to interpret repeated or conflicting facts correctly. Complete evidence delivery and a fitting count do not prove that interpretation.
If the summary plus newest cannot fit, compression returns the original history
with an explicit error to the runtime. PrepareMessages returns that
preparation error without exposing the original history as a fallback.
Counting or summary errors stop immediately rather than
trying another candidate, dropping evidence, or generating another summary.
There is no automatic restart. For K eligible turns, final selection makes at
most K exact count calls, stopping on fit or error. The unchanged trigger and
eligibility checks may also count; these bounds describe logical calls, not
provider HTTP attempts or billing. Broader summary input and extra counts may
increase cost and latency. The existing model, thresholds, single summary
completion, deadlines, and request limits remain in effect.
With no total limit, only the excluded older prefix enters the summary and the eligible exact suffix stays unchanged. There is no added overlap or final counting, even when an older-turn token allowance is configured. Empty, system-only, non-triggered, and nothing-to-summarize returns remain unchanged. Each invocation recomputes fit from its own messages and tools.
Custom prompt upgrade: with a positive total limit, WithSummaryPrompt now
receives all turns older than newest, not only discarded history. Change wording
that assumes “only discarded history” to refer to the supplied older history.
The caller’s chosen focus, %s interpolation, escaped percent signs, model class,
and summary role remain unchanged. No new configuration or stored-history
migration is needed. Without a total limit, the prefix scope is unchanged.
Evidence supplied to the summary model
Compress gives its selected older messages to the summary model as evidence,
not as a conversation to continue. Text, complete tool arguments and results,
call/result IDs, error status and full error messages, and citation fields are
quoted through the canonical model.Message JSON codec. Their original roles,
message/part positions, and order remain explicit. Values are not selected,
rounded, deduplicated, or replaced with tool-result placeholders. The model
decides which supplied facts matter; the runtime does not predict relevance.
WithSummaryPrompt still inserts the complete quoted textual transcript at its
%s placeholder. Images and documents remain native attachments, supplied once
after that prompt in the same completion. Each original user message containing
media produces one user attachment message, with matching original message/part
references. The history policy neither duplicates media bytes as base64 prose
nor extracts or fetches document bodies. Historical tools are quoted data: the
summary request advertises no tools and does not replay historical assistant
turns as new actions.
Original Message.Meta, thinking, cache checkpoints, and tool thought signatures
are not copied into the new summary request. This does not redact or change
original history, exact retained messages, stored diagnostics, or full tool
errors. Put facts that must be available for summarization in canonical text,
tool, citation, or media parts, not an opaque metadata map.
The returned summary keeps its configured WithSummaryRole, single text part,
[Conversation Summary] prefix, and runtime summary metadata. Plain text keeps
its existing representation. Cited sentences and all supplied citation fields
(title, source, location, and excerpts) survive in output order as quoted
records, not native citation blocks replayed under another role. Whitespace-only
generated text still fails as an empty summary, even with attribution fields.
Citation coordinates retain the meaning of the request that produced them.
Historical document indices are not reassigned to the summary’s attachments.
When a cited summary used native documents, its text also describes their
request layout: attachment-message position, document occurrence within that
message, original history position, and supplied name, format, and URI. It
includes no document bodies and does not claim a mapping to a provider’s
DocumentIndex or invent a file link from incomplete attribution. Later
compression treats these records as text, not a document registry to rebuild.
This fuller input may be larger than the former placeholder prompt. The coverage and exact-retention rules determine which messages are summarized and which complete histories are counted. Adapter support and existing client/provider limits still apply; providers may combine consecutive user messages. Unsupported media or an oversized summary request fails explicitly, without dropping evidence, a text-only fallback, or another summary call. Media adds no separate counting step. Complete input does not guarantee that the model preserves every important fact in its prose, or that every history fits the summary model.
Coordinated Generated-System Releases
Compatible releases may roll transparently; incompatible generated changes require a coordinated drain and cutover. See the authoritative Production rollout contract for the checkpoint-version, generated-codec, required-tool-name, and worker-retention requirements.
Bedrock Message Ordering Validation
When using AWS Bedrock with thinking mode enabled, the runtime validates message ordering constraints before sending requests. Bedrock requires:
- Any assistant message containing
tool_usemust start with a thinking block - Each user message containing
tool_resultmust immediately follow an assistant message with matchingtool_useblocks - The number of
tool_resultblocks cannot exceed the priortool_usecount
The Bedrock client validates these constraints early and returns a descriptive error if violated:
bedrock: invalid message ordering with thinking enabled (run=xxx, model=yyy):
bedrock: assistant message with tool_use must start with thinking
This validation ensures that transcript ledger reconstruction produces provider-compliant message sequences.
This ordering check runs before the provider call. Stream validation is a separate boundary: incomplete content blocks, unsigned reasoning, and missing stop reasons fail as output-contract errors before planner code receives an accepted response.
Next Steps
- Learn about Toolsets to understand tool execution models
- Explore Agent Composition for agent-as-tool patterns
- Read about Memory & Sessions for transcript persistence