ランタイム
アーキテクチャ概要
Goa-AI ランタイムは、Plan/Execute/Resume ループをオーケストレーションし、ポリシーを強制し、状態を管理し、エンジン、プランナー、ツール、メモリ、フック、およびフィーチャーモジュールと連携します。
ツールに加えて、ランタイムは最終アシスタント応答を型付きで扱う
Completion(...) コントラクトもサポートします。
これらのコントラクトは gen/<service>/completions に unary と
streaming のヘルパーを生成します。completion 名は 1-64 文字の
ASCII、英字・数字・_・- のみ、先頭は英字または数字という
ルールで DSL 境界で検証されます。streaming では completion_delta
はプレビュー専用で、正規の値は最後の 1 つの completion chunk
だけです。structured output を実装しない provider は
model.ErrStructuredOutputUnsupported を返します。
生成されるスキーマは正規かつ provider 非依存であり、モデル
アダプターは対応するサブセットへ正規化できますが、宣言された
契約を保てない場合は明示的に失敗しなければなりません。
| レイヤー | 責務 |
|---|---|
| DSL + Codegen | エージェント登録、ツール仕様/コーデック、ワークフロー、MCP アダプターを生成する |
| Runtime Core | plan/start/resume ループ、ポリシー強制、フック、メモリ、ストリーミングをオーケストレートする |
| Workflow Engine Adapter | Temporal アダプターが engine.Engine を実装し、他のエンジンも差し替え可能 |
| Feature Modules | 任意の統合(MCP、Pulse、Mongo ストア、モデルプロバイダーなど) |
ハイレベルなエージェントアーキテクチャ
Goa-AI は実行時に、少数の合成可能な構成要素を中心にシステムを組み立てます。
Agents:
agent.Ident(例:service.chat)で識別される長寿命のオーケストレーターです。各エージェントは、プランナー、ランポリシー、生成されたワークフロー、およびツール登録を所有します。Runs: エージェントの 1 回の実行です。
RunIDで識別され、run.Contextとrun.Handleで追跡されます。セッション付き run はSessionIDとTurnIDでグルーピングされて会話を構成し、one-shot run は明示的にセッションレスです。Toolsets & tools:
tools.Ident(service.toolset.tool)で識別される能力の集合です。サービスバックのツールセットは API を呼び出し、エージェントバックのツールセットは他のエージェントをツールとして実行します。Completions:
gen/<service>/completions配下に生成される、サービス所有の型付き直接アシスタント出力コントラクトです。completion helper は unary と direct streaming の model request に provider-enforced structured output を付与し、正規 payload を生成 codec で decode します。Planners: LLM による戦略レイヤで、
PlanStart/PlanResumeを実装します。プランナーは、ツールを呼ぶか直接回答するかを決め、ランタイムはその決定に対して上限(caps)と時間予算(time budget)を強制します。Run tree & agent-as-tool: あるエージェントが別のエージェントをツールとして呼ぶと、ランタイムは独自の
RunIDを持つ実際の子ランを開始します。親のToolResultには子ランへのRunLink(*run.Handle)が格納され、ストリーミングではchild_run_linkedイベントが親ツールコールと子ランを結び付けます。Session-owned streams & profiles: Goa-AI は型付けされた
stream.Eventを セッション所有ストリーム(session/<session_id>)へ発行します。イベントはRunIDとSessionIDを持ち、run_stream_endが SSE/WebSocket をタイマーなしで閉じるための明示マーカーになります。stream.StreamProfileは、対象(チャット UI、デバッグ、メトリクス)に応じてどのイベント種別を可視化するかを選択します。
クイックスタート
package main
import (
"context"
chat "example.com/assistant/gen/orchestrator/agents/chat"
"goa.design/goa-ai/runtime/agent/model"
"goa.design/goa-ai/runtime/agent/runtime"
)
func main() {
// In-memory engine is the default; pass WithEngine for Temporal or custom engines.
rt := runtime.New()
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 := rt.CreateSession(ctx, "session-1"); 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.
}
型付き直接 Completion
すべての構造化されたやり取りをツール呼び出しとして表現する必要はありません。サービスが型付きの最終アシスタント応答を必要とする場合は、DSL で Completion(...) を宣言して再生成します。
goa gen は gen/<service>/completions に次を出力します:
- result schema と型付き result/union 型
- 生成 JSON codec と validation helper
- 型付き
completion.Spec値 - 生成
Complete<Name>(ctx, client, req)helper - 生成
StreamComplete<Name>(ctx, client, req)とDecode<Name>Chunk(chunk)helper
service は Agent(...) を宣言しなくても completion を宣言できます。agent quickstart/example scaffold は、実際に agent を所有する service にだけ出力されます。
helper は request を clone し、provider-neutral structured output metadata を付与し、基盤の model.Client を呼び出し、正規の型付き payload を生成 codec で decode します:
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)
streaming completion は raw model.Streamer surface に留まり、最後の正規 completion chunk だけを decode します:
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)
}
value, ok, err := taskcompletion.DecodeDraftFromTranscriptChunk(chunk)
if err != nil {
panic(err)
}
if ok {
fmt.Println(value.Name)
}
}
型付き completion helper は意図的に厳格です:
- unary helper は unary request だけを受け付けます。
- completion 名は DSL 境界で検証されます。1-64 文字の ASCII、英字/数字/
_/-のみ、先頭は英字または数字です。 - unary と streaming helper は tool-enabled request と caller-supplied
StructuredOutputを拒否します。 - streaming provider は
completion_delta*preview fragment と正確に 1 つの正規completionchunk を emit するか、request を明示的に拒否します。 Decode<Name>Chunkは preview chunk を無視し、最後のcompletionだけを decode します。- completion stream は direct
model.Streamerpath に留まります。assistant transcript text/tool execution event 用の planner streaming helper には通さないでください。 - structured output を実装しない provider は
model.ErrStructuredOutputUnsupportedを表面化します。 - 生成 schema は正規かつ provider-neutral です。provider adapter は対応 subset へ normalize できますが、宣言された contract を保てない場合は明示的に失敗しなければなりません。
クライアント専用 vs ワーカー
ランタイムは大きく 2 つのロールで利用されます。
クライアント専用(run の送信): クライアント機能を持つエンジンでランタイムを構築し、エージェント登録は行いません。生成された
<agent>.NewClient(rt)は、リモートワーカーによって登録されたルート(workflow + queue)を保持しており、これを用いて run を送信します。ワーカー(run の実行): ワーカー機能を持つエンジンでランタイムを構築し、実際のプランナーを使ってエージェントを登録します。その上で、エンジンが workflow/activity をポーリングして実行します。
クライアント専用の例
rt := runtime.New(runtime.WithEngine(temporalClient)) // engine client
// No agent registration needed in a caller-only process
client := chat.NewClient(rt)
if _, err := rt.CreateSession(ctx, "s1"); err != nil {
panic(err)
}
out, err := client.Run(ctx, "s1", msgs)
セッションレス one-shot 実行
既存セッションに紐づかない耐久実行が必要な場合は StartOneShot と OneShotRun を使います。
Start/Runはセッション付きです。具体的なSessionIDが必要で、セッションのライフサイクルに参加し、セッションスコープのストリームイベントを発行します。StartOneShot/OneShotRunはセッションレスです。SessionIDを受け取らず、セッションも作成せず、RunIDによる introspection のための canonical な runlog イベントだけを追記します。StartOneShotはengine.WorkflowHandleを即座に返します。OneShotRunは内部でhandle.Wait(ctx)を呼ぶ blocking な convenience wrapper です。
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)
ワーカーの例
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(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 ループ
- ランタイムはエージェントのワークフロー(インメモリまたは Temporal)を開始し、
RunID、SessionID、TurnID、ラベル、ポリシー上限を含む新しいrun.Contextを記録します。 - 現在のメッセージと run コンテキストを渡して、プランナーの
PlanStartを呼び出します。 - プランナーが返したツール呼び出しをスケジュールします(プランナーは「正規(canonical)JSON」のペイロードを渡し、エンコード/デコードはランタイムが生成済みコーデックで処理します)。
- プランナーから見えるまま残ったツール結果を添えて
PlanResumeを呼び出します。予算対象ツールは既定で可視ですが、bookkeeping ツールはRetryHint.AllowsRetry()が修復を許可する場合にのみ再生されます。プランナーが最終応答、最終ツール結果を返すか、成功したTerminalRunツールが run を完了するまでループします。上限や deadline によって強制 finalization が有効な場合、プランナーは prose ではなく terminal bookkeeping ツールで閉じることができます。進行に応じて run はrun.Phase(prompted/planning/executing_tools/synthesizing/ 終端フェーズ)を遷移します。 - フックとストリームサブスクライバは、イベント(プランナー思考、ツール start/update/end、await、usage、workflow、agent-run links)を発行し、設定に応じてトランスクリプトや run メタデータを永続化します。
Run フェーズ
run が plan/execute/resume ループを進むにつれて、ライフサイクルフェーズを遷移します。フェーズは、run が今どの段階にいるかをきめ細かく可視化し、UI が高レベルの進捗を表示できるようにします。
フェーズ値(Phase Values)
| Phase | 説明 |
|---|---|
prompted | 入力を受け取り、これからプランニングを開始する状態 |
planning | ツールを呼ぶか直接答えるか、どのように進めるかをプランナーが判断している状態 |
executing_tools | ツール(ネストされたエージェントを含む)が実行中の状態 |
synthesizing | 追加ツールをスケジュールせず最終回答を合成している状態 |
completed | 正常に完了した状態 |
failed | 失敗した状態 |
canceled | キャンセルされた状態 |
フェーズ遷移
典型的な成功 run は、次のような経過をたどります。
prompted → planning → executing_tools → planning → synthesizing → completed
↑__________________|
(loop while tools needed)
ランタイムは planning / executing_tools / synthesizing などの 非終端フェーズに対して RunPhaseChanged フックイベントを発行し、ストリーム購読者がリアルタイムに進捗を追跡できるようにします。
Phase と Status の違い
フェーズは run.Status とは異なります。
- Status(
pending,running,completed,failed,canceled,paused)は、耐久化された run メタデータに格納される粗い粒度のライフサイクル状態です。 - Phase は、ストリーミング/UX 向けに実行ループをより細かく可視化するものです。
ライフサイクルイベント: フェーズ遷移 vs 終端完了
ランタイムは次を発行します:
RunPhaseChanged: 非終端フェーズ遷移。RunCompleted: run ごとに 1 回の終端ライフサイクル(success / failed / canceled)。
ストリーム購読者は、両方を workflow ストリームイベント(stream.WorkflowPayload)に変換します:
- 非終端更新(
RunPhaseChanged):phaseのみ。 - 終端更新(
RunCompleted):status+ 終端phase。失敗時は構造化されたエラー情報を含みます。
終端 status のマッピング
status="success"→phase="completed"status="failed"→phase="failed"status="canceled"→phase="canceled"
キャンセルはエラーではありません
status="canceled" の場合、ストリームペイロードにユーザー向け error を含めてはいけません。
失敗は構造化されます
status="failed" の場合、ストリームペイロードに以下が含まれます:
error_kindretryableerror(ユーザー向け)debug_error(診断向け)
終端のアイデンティティ
RunCompleted は Labels を持ちます。run 開始時に指定した run スコープの
ラベル(RunInput.Labels、runtime.WithLabels(...) で設定)で、ラベルが
なかった run では nil です。完了サブスクライバーは、独自の run-ID から
アイデンティティへのマップを維持することなく、終端結果(success / failed /
canceled)を帰属できます。同じラベルはポーリングリーダー向けに
run.Snapshot.Labels にも公開され、永続的な RunStarted レコードから
再構築されるため、run のアイデンティティは両エンジンでプロセス再起動を
またいで保持されます。run 途中にポリシー決定がマージしたラベルは含まれず、
PolicyDecision イベントで引き続き観測できます。
ポリシー、上限(Caps)、ラベル
設計時 RunPolicy
設計時には、RunPolicy でエージェントごとのポリシーを設定します。
Agent("chat", "Conversational runner", func() {
RunPolicy(func() {
DefaultCaps(
MaxToolCalls(8),
MaxConsecutiveFailedToolCalls(3),
)
TimeBudget("2m")
InterruptsAllowed(true)
})
})
これはエージェント登録に紐づく runtime.RunPolicy になります。
- Caps:
MaxToolCallsは run あたりの予算対象 tool call 総数です。DSL でBookkeeping()として宣言されたツールは retrieval budget を消費せず、MaxConsecutiveFailedToolCallsも変更しません。モデルが生成した batch は原子的なままです。bookkeeping call のコストはゼロですが、mixed batch を収めるために個々の call を除去することはありません。成功した bookkeeping 結果は将来の compact なToolOutputsに入りません。 - Time budget:
TimeBudget(run の wall-clock 予算)、FinalizerGrace(ランタイム専用: 最終化のための予約ウィンドウ)。 - Interrupts:
InterruptsAllowed(pause/resume のオプトイン)。 - Terminal tools: DSL で
TerminalRun()として宣言されたツールは自動的に bookkeeping となり、成功すると後続PlanResumeなしで run を終了します。したがって terminal commit は retrieval budget が残っていなくても受け入れられます。強制 finalization 中、ランタイムは terminal bookkeeping call だけを受け入れ、残りの hard deadline 内で実行し、すべての terminal effect が成功した場合にのみ run を閉じます。 - Missing fields behavior:
OnMissingFields(バリデーションが欠落フィールドを示した場合の挙動)。
ランタイムポリシーのオーバーライド
環境によっては、設計を変更せずにポリシーを強化/緩和したい場合があります。rt.OverridePolicy API により、プロセスローカルにポリシーを調整できます。
err := rt.OverridePolicy(chat.AgentID, runtime.RunPolicy{
MaxToolCalls: 3,
MaxConsecutiveFailedToolCalls: 1,
InterruptsAllowed: true,
})
Scope: オーバーライドは現在のランタイムインスタンスにローカルで、以降の run にのみ影響します。プロセス再起動を越えて永続化されず、他ワーカーへも伝播しません。
Overridable Fields:
| Field | 説明 |
|---|---|
MaxToolCalls | run あたりの予算対象ツール呼び出し総数の上限(Bookkeeping() ツールは免除) |
MaxConsecutiveFailedToolCalls | 連続失敗回数の上限 |
TimeBudget | run の wall-clock 予算 |
FinalizerGrace | 最終化のための予約ウィンドウ |
InterruptsAllowed | pause/resume を有効化する |
ゼロ値でないフィールドのみが適用されます(InterruptsAllowed は true の場合に適用)。これにより、他のポリシー設定へ影響を与えず選択的なオーバーライドが可能です。
Use Cases:
- プロバイダスロットリング時の一時的なバックオフ
- ポリシー設定の A/B テスト
- 制約を緩和した開発/デバッグ
- テナントごとのランタイムポリシー調整
ラベルとポリシーエンジン
Goa-AI は policy.Engine を介してプラガブルなポリシーエンジンと統合します。ポリシーは、ツールのメタデータ(ID、タグ)、run コンテキスト(SessionID、TurnID、labels)、そしてツール失敗後の RetryHint 情報を受け取ります。
ラベルは次に流れます。
run.Context.Labels– run 中にプランナーが参照可能- ツールアクティビティ入力(
api.ToolInput.Labels)– dispatch 済みのツール実行へクローンされ、プランナーが特定の呼び出しで上書きしない限り、ツールアクティビティは同じ run スコープ metadata を参照できます - runlog イベント(
runlog.Store)– ライフサイクルイベントとともに永続化され、検索/ダッシュボードに有用(インデックスされる場合) - 終端完了とスナップショット – 開始時のラベルは run の最後に
hooks.RunCompletedEvent.Labelsとrun.Snapshot.Labelsとして戻ってくるため、完了フックやGetRunSnapshotのリーダーは帯域外の追跡なしに run のアイデンティティを取得できます
run ごとのツールフィルタリング
設計時 tag と runtime option により、planner prompting の前と execution の前に tool surface を絞り込めます:
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"}},
}),
)
repair flow で 1 つの tool だけを公開したい場合は WithRestrictToTool を使います:
out, err := client.Run(ctx, "session-1", messages,
runtime.WithRestrictToTool(searchspecs.Search),
)
retry hint が RestrictToTool を設定した repair turn では、runtime は restricted-tool を latch します。次の planner turn は修正が必要な tool だけを見ます。これにより validation repair が focused になり、無関係な tool へ drift することを防ぎます。
ツール実行
- Native toolsets: 実装はアプリ側で書き、ランタイムが生成済みコーデックで型付き引数をデコードします。
- Agent-as-tool: 生成された agent-tool ツールセットはプロバイダーエージェントを子ランとして実行し(プランナー視点ではインライン)、その
RunOutputをplanner.ToolResultに変換し、子ランへのRunLink(ハンドル)を返します。 - MCP toolsets: ランタイムは正規 JSON を生成済み caller へ転送し、caller がトランスポートを扱います。
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
大きな data set の一部だけを返す tool は、DSL で BoundedResult(...) を宣言するべきです。これらの tool の runtime contract は次の通りです:
- 生成
tools.ToolSpec.Boundsが正規 bounded-result schema を宣言する - successful execution は
planner.ToolResult.Boundsを populate する必要がある - runtime は provider-owned bounds を emitted
tool_resultJSON、.Bounds配下の result-hint template data、hook payload、stream event へ project する - paged tool では、provider code は次ページの opaque cursor を
Bounds.NextCursorに設定します
tools.ToolSpec.Bounds は model-facing JSON 名を使います。DSL 宣言が
NextCursor("nextCursor") のような lower-camel Goa attribute を参照しても、
生成 specs、schemas、runtime projection、result codec は next_cursor を使います。
正規 projected field:
returned(required)truncated(required)total(optional)refinement_hint(optional)next_cursor(directCursorcontract がNextCursor(...)を公開する場合は optional)
planner.ToolResult.Bounds が唯一の machine-readable provider contract です。手書きの Go result type は semantic かつ domain-specific のままでよく、model に見せるためだけに正規 bounded field を重複させる必要はありません。
ContinueWith("continue_tool", "cursor") は機械的な continuation を別 action として宣言します。runtime は、次の cursor を持つ live chain head が一つだけある場合に限って action を公開します。正確な cursor lineage によって連続 page を進めます。source call は並列実行できますが、複数の live head がある間は引数なしの action を公開しません。model は {} で action を選び、runtime が実行前に cursor と保持済み query field を bind します。direct Cursor("cursor") は open contract であり、model は next_cursor の opaque cursor と変更していない query argument を次の call に指定します。
method-backed BindTo tool では、生成 executor が projection 前に planner.ToolResult.Bounds を構築できるよう、bound service method result は正規 bounded field を保持する必要があります。明示的な tool-facing Return(...) shape はそれらの正規 field を重複させてはいけません。bound method result の中で required にできるのは returned と truncated だけです。total、refinement_hint、next_cursor は bounds contract の optional part のままで、runtime bounds が省略した場合は emitted JSON からも省略されます。
service boundary が ExecuteToolActivity の外で正規 result JSON を assemble する必要がある場合は、生成 result codec と bounded-result projection helper を別々に呼ぶのではなく runtime.EncodeCanonicalToolResult(...) を使います。
Prompt ランタイムコントラクト
Prompt 管理はランタイムネイティブで、バージョン管理されます。
runtime.PromptRegistryは不変なベースラインprompt.PromptSpec登録を保持するruntime.WithPromptStore(prompt.Store)はスコープ付き override 解決(session->facility->org-> global)を有効化する- プランナーは
PlannerContext.RenderPrompt(ctx, id, data)を呼び、prompt 内容を解決・描画する - 描画済み内容には provenance 用の
prompt.PromptRefが含まれ、プランナーはmodel.Request.PromptRefsに付与できる
content, err := input.Agent.RenderPrompt(ctx, "aura.chat.system", map[string]any{
"AssistantName": "Ops Assistant",
})
if err != nil {
return nil, err
}
resp, err := modelClient.Complete(ctx, &model.Request{
RunID: input.RunContext.RunID,
Messages: input.Messages,
PromptRefs: []prompt.PromptRef{content.Ref},
})
PromptRefs は監査/プロベナンス向けのランタイムメタデータであり、プロバイダー wire payload のフィールドではありません。
メモリ、ストリーミング、テレメトリ
Hook bus は、run の開始/完了、フェーズ変更、
prompt_rendered、ツールのスケジューリング/結果/更新、プランナーノートと思考ブロック、await、retry hints、agent-as-tool links など、エージェントライフサイクル全体の構造化フックイベントを publish します。Memory stores(
memory.Store)は、(agentID, RunID)ごとに耐久化されるメモリイベント(ユーザー/アシスタントメッセージ、ツール呼び出し、ツール結果、プランナーノート、思考)を購読し追記します。Run event stores(
runlog.Store)は、RunIDごとに hook イベントのカノニカルログを追記し、audit/debug UI と run の introspection に利用できます。Stream sinks(
stream.Sink。例: Pulse またはカスタム SSE/WebSocket)は、stream.Subscriberが生成する型付きstream.Eventを受け取ります。StreamProfileは送出するイベント種別を制御します。Telemetry: OTEL 対応のロギング、メトリクス、トレーシングが workflow/activity を end-to-end で計測します。
ツール呼び出しヒント(DisplayHint)
ツール呼び出しには、ユーザー向けの DisplayHint(例: UI 表示用)を含めることができます。
契約:
- hooks のイベントコンストラクタはヒントをレンダリングしません。ツール呼び出しのスケジュールイベントは既定で
DisplayHint==""です。 - ランタイムは、payload のデコードに成功した場合、型付きテンプレートから公開時に 永続的な 呼び出しヒントを付与して保存します。
- ツール登録には空でないメタデータ title が必要です。型付きデコードに失敗する、またはテンプレートが登録されていない場合、ランタイムはその title を display hint として使用します。不正な payload は引き続きツール境界で失敗します。メタデータ title は、試行された作業を表示可能に保つためだけに使われます。ヒントは生の JSON に対してレンダリングされません。
- producer が hook イベントを公開する前に
DisplayHint(非空)を明示的に設定した場合、ランタイムはそれを権威ある値として扱い、上書きしません。 - consumer ごとの文言変更(例: UI の表現)にはランタイムで
runtime.WithHintOverridesを設定します。override は、ストリームのtool_startイベントにおいて DSL テンプレートより優先されます。
セッションストリームの消費(Pulse)
プロダクションでは一般に以下のパターンを取ります:
- 共有バス(Pulse / Redis Streams)からセッションストリーム(
session/<session_id>)を消費する run_idでフィルタして run ごとのカード/レーンを構築する- アクティブ run の
run_stream_endを観測したら SSE/WebSocket を終了する
import "goa.design/goa-ai/runtime/agent/stream"
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
}
case err := <-errs:
panic(err)
}
}
エンジン抽象
- In-memory: 高速な開発ループ、外部依存なし
- Temporal: 耐久実行、リプレイ、リトライ、シグナル、ワーカー。アダプタがアクティビティとコンテキスト伝搬を統合します。
セマンティックな Timing と Temporal の Liveness
Goa-AI は公開ランタイム契約をエンジン非依存に保ちます:
RunPolicy.Timing.PlanとRunPolicy.Timing.Toolsはセマンティックな「試行ごとの予算」runtime.WithTiming(...)は run ごとにそれらのセマンティック予算を上書きするruntime.WithWorker(...)はキュー配置のためのもので、ワークフローエンジン調整ではない
Temporal アダプタを使っていて、キュー待ちや liveness を調整したい 場合は、それらを Temporal エンジン側で設定します:
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)
}
この分離により、ワークフローのメカニクスは Temporal の境界の内側に 保たれ、汎用ランタイムは Temporal とインメモリエンジンの両方に対して 正直なままでいられます。
Run コントラクト
SessionIDはセッション付き開始で必須です。StartとRunはSessionIDが空、または空白のみの場合に fail-fast します。StartOneShotとOneShotRunは明示的にセッションレスです。セッションを要求/作成せず、セッションスコープのストリームイベントも発行しません。- エージェントは最初の run の前に登録されなければなりません。ランタイムは、エンジンワーカーの決定性を保つため、最初の run 送信後の登録を
ErrRegistrationClosedで拒否します。 - ツール実行者は、
context.Contextから値を“釣る”のではなく、呼び出しごとの明示メタデータ(ToolCallMeta)を受け取ります。 - 暗黙のフォールバックには依存しません。すべてのドメイン識別子(run / session / turn / correlation)は明示的に渡します。
一時停止と再開
Human-in-the-loop のワークフローは、ランタイムの interrupt ヘルパーを使って run を一時停止/再開できます。
import "goa.design/goa-ai/runtime/agent/interrupt"
// Pause
if err := rt.PauseRun(ctx, interrupt.PauseRequest{
RunID: "session-1-run-1",
Reason: "human_review",
}); err != nil {
panic(err)
}
// Resume
if err := rt.ResumeRun(ctx, interrupt.ResumeRequest{
RunID: "session-1-run-1",
}); err != nil {
panic(err)
}
内部では pause/resume シグナルが run ストアを更新し、run_paused / run_resumed フックイベントを発行するため、UI レイヤは同期を維持できます。
外部ツール結果の提供
一部の await は、ExecuteToolActivity 自身ではなく 外部アクターが提供したツール結果 によって再開されます。代表例は、構造化質問のような UI 管理ツールや、別システムから結果を集めて run を再開させるブリッジサービスです。
生の提供結果とともに ProvideToolResults を使います:
err := rt.ProvideToolResults(ctx, interrupt.ToolResultsSet{
RunID: "run-123",
ID: "await-1",
Results: []*api.ProvidedToolResult{
{
Name: "chat.ask_question.ask_question",
ToolCallID: "toolcall-1",
Result: rawjson.Message(`{"answers":[{"question_id":"topic","selected_ids":["alarms"]}]}`),
},
},
})
契約:
- 呼び出し側は、正規の生結果 JSON と、任意の
Bounds、Error、RetryHintを渡します。 - 呼び出し側が
api.ToolEventを構築することは ありません。それはランタイム内部の workflow エンベロープです。 - ランタイムは、登録済みツール仕様を使って提供結果をデコードし、型付きの結果 materialization を実行し、必要な server-only sidecar を付与し、正規の
tool_resultを transcript/run log に追記してから、計画処理を再開します。
これにより await パスは通常の実行パスと概念的に揃います。どちらのフローも、公開前に同じ型付き planner.ToolResult 契約へ収束します。
ツール確認(Tool Confirmation)
Goa-AI は、書き込み・削除・コマンド実行などのセンシティブなツールに対して、ランタイム強制の確認ゲートをサポートします。
確認は次の 2 通りで有効化できます。
- 設計時(一般的): ツール DSL 内で
Confirmation(...)を宣言します。Codegen はポリシーをtools.ToolSpec.Confirmationに格納します。 - ランタイム(上書き/動的): ランタイム構築時に
runtime.WithToolConfirmation(...)を渡し、追加ツールに確認を要求したり設計時の挙動を上書きしたりできます。
実行時には、workflow がアウトオブバンドの確認要求を発行し、明示的な承認が与えられた後にのみツールを実行します。拒否された場合、ランタイムはスキーマ準拠のツール結果を合成し、トランスクリプトの整合性(決定性)を保ったままプランナーが反応できるようにします。
確認プロトコル
実行時の確認は、専用の await/decision プロトコルとして実装されます。
- Await payload(
await_confirmationとしてストリームされる):
{
"id": "...",
"title": "...",
"prompt": "...",
"tool_name": "atlas.commands.change_setpoint",
"tool_call_id": "toolcall-1",
"payload": { "...": "canonical tool arguments (JSON)" }
}
契約:
payloadには常に、保留中のツール呼び出しに対する正規の JSON 引数が入ります。承認された場合、ランタイムが実行するのはその引数です。確認のオーバーライドは prompt や拒否結果のレンダリングをカスタマイズできますが、表示専用の別 payload チャネルを導入したり、
payloadの意味を変えたりしてはいけません。よりリッチな確認 UI が必要なプロダクトでは、アプリケーション層で正規 payload とアプリケーション所有の読み取り結果からその表示を materialize してください。
決定の提供(ランタイムの
ProvideConfirmationを通して):
err := rt.ProvideConfirmation(ctx, interrupt.ConfirmationDecision{
RunID: "run-123",
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"},
})
ツール承認イベント
決定が提供されると、ランタイムは第一級の承認イベントを発行します:
- Hook event:
hooks.ToolAuthorization - Stream event type:
tool_authorization
このイベントは、確認が必要なツール呼び出しに対する “who/when/what” の正規レコードです:
tool_name,tool_call_idapproved(true/false)summary(ランタイムが決定論的にレンダリングする要約)approved_by(interrupt.ConfirmationDecision.RequestedByからコピーされる安定 principal ID)
イベントは決定受信直後に発行されます(承認時はツール実行前、拒否時は拒否結果の合成前)。
注意:
- コンシューマは確認を「ランタイムプロトコル」として扱うべきです。
- 付随する
RunPausedの理由(await_confirmation)を見て、確認 UI を出すべきタイミングを判断します。 - 確認 UI の挙動を特定の確認ツール名に結びつけないでください(内部トランスポート詳細として扱います)。
- 付随する
- 確認テンプレート(
PromptTemplateとDeniedResultTemplate)は Go のtext/template文字列で、missingkey=errorで実行されます。標準関数(例:printf)に加えて、Goa-AI は次を提供します。json v→vを JSON エンコード(オプショナルポインタや構造値の埋め込みに便利)quote s→ Go のエスケープ済み引用符付き文字列を返す(fmt.Sprintf("%q", s)相当)
ランタイムバリデーション
ランタイムは境界で確認操作をバリデートします。
- 提供された確認
IDが、保留中の await 識別子と一致すること - decision オブジェクトが整形されていること(空でない
RunID、真偽値のApproved)
プランナー契約
プランナーは次を実装します。
type Planner interface {
PlanStart(ctx context.Context, input *planner.PlanInput) (*planner.PlanResult, error)
PlanResume(ctx context.Context, input *planner.PlanResumeInput) (*planner.PlanResult, error)
}
PlanResult には tool call、最終応答、最終 tool result、注釈、選択された
post-tool transition が含まれます。PlanResumeInput は、プランナーが呼ばれた
理由を示します。
これらの契約は別々です。
| 契約 | スコープ | 意味 |
|---|---|---|
ToolSpec.Tags | すべての run における 1 つの tool | 汎用的な policy と UI filtering に使えるフラットなラベル。 |
ToolSpec.Meta | すべての run における 1 つの tool | 名前付きコンシューマが意味を所有する、不活性な生成アノテーション。メタデータだけでは runtime 動作は変わらない。 |
ToolSpec.Bookkeeping | すべての run における 1 つの tool | 成功後に別の planner turn を必要としない durable な制御記録。retrieval と連続失敗の budget を消費しない。 |
ToolSpec.TerminalRun | すべての run における 1 つの tool | 成功そのものが run を終了し、自動的に bookkeeping を含む。 |
RetryHint.AllowsRetry() | 1 つの失敗 result | この run で別の tool attempt が許可される。timeout hint は terminal failure を表し false を返す。 |
PlanResult.SynthesizeAfterTools | 選択された 1 batch | recoverable failure がなければ、次の planner turn は回答しなければならない。 |
PlanResumeInput.SynthesisOnly | 1 planner activity | 最終回答を返す。tool call は無効。 |
PlanResumeInput.Finalize | runtime が強制する終了 | cap または deadline により通常作業が禁止されている。 |
ランタイムは次の順序で次状態を選びます。
| 完了したステップ | 次の状態 |
|---|---|
| cap または deadline が finalization を要求 | Finalize turn |
TerminalRun tool が成功 | 即時終了 |
失敗 result のいずれかで AllowsRetry() == true | 通常の repair turn |
SynthesizeAfterTools が true | SynthesisOnly turn |
| その他 | 通常の continuation turn |
これにより planner intent が 2 つ目の retry policy になることを防ぎます。
recoverable failure を先に修復し、成功した final batch または terminal failure を
含む final batch は synthesis に進みます。ランタイムは SynthesisOnly turn
から返された tool call を拒否します。
recoverable な ToolFailure は Recovery.Action も 1 つ選択します。
correct_callは失敗した tool を引き続き利用可能にし、拒否された input、 生成済み validation issue、field guidance、example を次の planner turn に 渡します。失敗 1 件につき replacement call 1 件を要求するものではありません。 planner は作業をまとめ、表示された tool を任意の回数だけ正しく呼び出し、 input を待つか、すでに集めた evidence から回答できます。replanは失敗した tool を次の planner turn から除外します。planner は別の 表示済み tool を使うか、input を待つか、回答できます。finishはすべての tool を除外し、利用可能な evidence に基づく最終回答を 要求します。
ランタイムは recovery turn で表示した tool catalog を正確に記録し、その外側の 実行可能な call をすべて拒否します。user または外部 input の要求に埋め込まれた call も対象です。生成済み codec は引き続きすべての payload を検証し、run の tool・failure・time limit は無効な作業の繰り返しを停止します。recovery turn が input を待つ場合、failure evidence は再開後も利用できます。tool call または最終 回答を選ぶと、その evidence は消去されます。
recovery activity の input と表示された catalog は、durable workflow history の 一部です。この contract を変更する deployment では、新しい worker bundle を 開始する前に、古い worker と実行中 workflow を drain または stop する必要が あります。この境界をまたいで worker version を混在させることは安全ではありません。
PlanResumeInput.Finalize が設定されている場合、プランナーは terminal
bookkeeping tool を返せます。これらは後続プランナーターンには再生されず、
finalization を永続的に完了する必要があります。
プランナーは input.Agent 経由でランタイムサービスを提供する PlannerContext も受け取ります。
AdvertisedToolDefinitions()- このターンでモデルに見えている、runtime がフィルタ済みのツール定義を取得するModelClient(id string)- provider-agnostic な生のモデルクライアントを取得するPlannerModelClient(id string)- planner ターン専用で runtime-owned なイベント発行を行うモデルクライアントを取得するRenderPrompt(ctx, id, data)- 現在の run scope で prompt 内容を解決・描画するAddReminder(r reminder.Reminder)- run スコープの system reminder を登録するRemoveReminder(id string)- 前提条件が満たされなくなったときに reminder を削除するMemory()- 会話履歴へアクセスする
フィーチャーモジュール
runtime/mcp– HTTP、SSE、stdio transport 用の MCP callerfeatures/memory/mongo– durable memory storefeatures/prompt/mongo– Mongo-backed prompt override storefeatures/runlog/mongo– run event log store(append-only, cursor pagination)features/session/mongo– session metadata storefeatures/stream/pulse– Pulse sink/subscriber helpersfeatures/model/{anthropic,bedrock,openai}– モデルクライアントアダプター(プランナー向け)features/model/middleware– 共有model.Clientミドルウェア(例: 適応型レート制限)features/policy/basic– allow/block リストと retry hint を扱う簡易ポリシーエンジン
モデルクライアントのスループット & レート制限
Goa-AI は features/model/middleware に provider-agnostic な適応型レートリミッターを提供します。これは任意の model.Client をラップし、リクエストごとのトークンを推定し、トークンバケットで呼び出しをキューイングし、プロバイダがスロットリングを返したときに AIMD(additive-increase/multiplicative-decrease)戦略で実効 TPM 予算を調整します。
import (
"goa.design/goa-ai/features/model/bedrock"
mdlmw "goa.design/goa-ai/features/model/middleware"
)
awsClient := bedrockruntime.NewFromConfig(cfg)
bed, _ := bedrock.New(awsClient, bedrock.Options{
DefaultModel: "us.anthropic.claude-4-5-sonnet-20251120-v1:0",
})
rl := mdlmw.NewAdaptiveRateLimiter(
ctx,
throughputMap, // *rmap.Map joined earlier (nil for process-local)
"bedrock:sonnet", // key for this model family
80_000, // initial TPM
1_000_000, // max TPM
)
limited := rl.Middleware()(bed)
rt := runtime.New()
if err := rt.RegisterModel("bedrock", limited); err != nil {
panic(err)
}
LLM 統合
Goa-AI のプランナーは、provider-agnostic なインターフェースを通じて大規模言語モデルと対話します。この設計により、プランナーコードを変えずに、AWS Bedrock、OpenAI、Google Vertex AI(Gemini / Claude-on-Vertex)、カスタムエンドポイントなどのプロバイダーを切り替えられます。
model.Client インターフェース
すべての LLM 呼び出しは model.Client を通ります。
type Client interface {
Complete(ctx context.Context, req *Request) (*Response, error)
Stream(ctx context.Context, req *Request) (Streamer, error)
}
プロバイダーアダプター
Goa-AI は一般的な LLM プロバイダー向けのアダプターを同梱しています。
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,
})
OpenAI
import "goa.design/goa-ai/features/model/openai"
modelClient, err := openai.New(openai.Options{
APIKey: apiKey,
DefaultModel: "gpt-5-mini",
HighModel: "gpt-5",
SmallModel: "gpt-5-nano",
})
Google Vertex AI(Gemini / Claude-on-Vertex)
features/model/vertex パッケージは、いずれも model.Client を満たす 2 つの
コンストラクタを提供します。ネイティブの Gemini アダプタと、Vertex 上でホスト
される Claude モデルに Anthropic アダプタを向ける純粋なコンストラクタヘルパー
です。
import "goa.design/goa-ai/runtime/agent/runtime"
// Vertex 上の Gemini。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,
})
// Vertex 上の Claude。これは純粋な構築です: SDK の Vertex トランスポートに
// 対して Anthropic SDK クライアントを構築し、features/model/anthropic に渡し
// ます。同パッケージが、Anthropic がホストするすべてのアダプタ(直接 API /
// Vertex ホスト共通)の Messages 変換と HTTP ステータスのエラー分類を担います
// — 別個の変換レイヤはありません。
claudeOnVertexClient, err := rt.NewVertexAnthropicModelClient(ctx, runtime.VertexConfig{
ProjectID: "my-gcp-project",
Location: "us-east5",
DefaultModel: "claude-sonnet-4-5@20250929",
})
Gemini 3 世代のモデルは、ツール呼び出しの背後にある推論チェーンを認証する
ために、functionCall パート(thought/thinking パートだけでなく)に不透明な
thought signature を付与します。Vertex アダプタはこのシグネチャを、
ThinkingPart.Signature と同じ base64 規約で
model.ToolCall.ThoughtSignature / model.ToolUsePart.ThoughtSignature を
通じてラウンドトリップします。ランタイムはこのシグネチャをモデルクライアント
境界で(以下のどちらの統合スタイルでも planner.ToolRequest が生成される前
に)捕捉し、プロバイダー向けトランスクリプトを再構築する際にツールコール ID
で再付与します。planner.ToolRequest にシグネチャフィールドはありません。
プランナーコードがシグネチャの存在を意識する必要はありません。
正規メッセージメタデータと引用の再生
model.Message.Meta には、応答を正確に再生するために必要なプロバイダー生成
データが含まれます。メタデータを永続化または転送する境界では、
model.MarshalMetadata と model.UnmarshalMetadata を使用してください。
これらの codec は単一の JSON オブジェクトを要求し、デコードした数値を
json.Number として保持し、後続データを拒否し、nil または空オブジェクトを
nil に正規化します。
引用の再生はプロバイダー固有であり、引用を通常のテキストに平坦化しては
なりません。Bedrock アダプタは、assistant の CitationsPart をネイティブな
引用コンテンツブロックとして再生し、ソース ID、抜粋、および文書内の文字・
チャンク・ページ位置を保持できます。Bedrock の system content union には
引用メンバーがないため、system 引用はサポートされません。Anthropic と
Vertex は、正規パートに各プロバイダーのプロトコルで必須のフィールドが
ない場合、引用の再生を拒否します。
プランナーでモデルクライアントを使う
プランナーはランタイムの PlannerContext 経由でモデルクライアントを取得します。
現在は、統合スタイルが明示的に 2 つあります。
PlannerModelClient(id)は planner ターン専用の streaming と runtime-owned なイベント発行に使うModelClient(id)は生の transport アクセスが必要で、planner.ConsumeStreamと組み合わせるかPlannerEventsを自前で発行したいときに使う
PlannerModelClient(推奨)
PlannerContext.PlannerModelClient(id) は、AssistantChunk、
PlannerThinkingBlock、UsageDelta の発行を担う planner ターン専用の
クライアントを返します。Stream(...) は基盤となる provider stream を
drain し、planner.StreamSummary を返します。
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")
}
req := &model.Request{
Messages: input.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 テキストはすでにストリーム済み
}, nil
}
これは最も安全な統合スタイルです。planner 専用クライアントは生の
model.Streamer を公開しないため、planner.ConsumeStream と誤って
組み合わせることがありません。また、sum.FinalResponse() はその呼び出しで
捕捉したプロバイダーの正確な応答を選択します。テキストだけのメッセージを
再構築すると、thinking、引用、シグネチャ、メタデータ、メッセージ境界が
失われます。
生の Client + ConsumeStream
生の model.Client が必要な場合は PlannerContext.ModelClient から取得し、
planner.ConsumeStream と組み合わせます。
mc, ok := input.Agent.ModelClient("anthropic.claude-3-5-sonnet-20241022-v2:0")
if !ok {
return nil, errors.New("model not configured")
}
req := &model.Request{
Messages: input.Messages,
Tools: input.Agent.AdvertisedToolDefinitions(),
Stream: true,
}
streamer, err := mc.Stream(ctx, req)
if err != nil {
return nil, err
}
sum, err := planner.ConsumeStream(ctx, streamer, req, input.Events)
if err != nil {
return nil, err
}
この helper は stream を drain し、assistant / thinking / usage の
イベントを発行しつつ、集約済みのテキストとツール呼び出しを含む
StreamSummary を返します。
生の client 経路は、stream 消費を完全に制御したい場合、early-stop の
独自挙動が必要な場合、または PlannerEvents を明示的に扱いたい場合に
使います。PlannerModelClient.Stream(...) と
planner.ConsumeStream を混在させず、planner ターンごとに stream owner
を 1 つ選んでください。
Bedrock メッセージ順序の検証
AWS Bedrock で thinking mode を有効にすると、ランタイムはリクエスト送信前にメッセージ順序制約を検証します。Bedrock は次を要求します。
tool_useを含むアシスタントメッセージは、必ず thinking ブロックから始まることtool_resultを含むユーザーメッセージは、対応するtool_useブロックを持つアシスタントメッセージの直後に続くことtool_resultブロック数は、直前のtool_use数を超えないこと
Bedrock クライアントはこれらを早期に検証し、違反時は説明的なエラーを返します。
bedrock: invalid message ordering with thinking enabled (run=xxx, model=yyy):
bedrock: assistant message with tool_use must start with thinking
この検証により、トランスクリプト ledger の再構築がプロバイダー準拠のメッセージ列を生成することを保証します。
次のステップ
- ツール実行モデルを理解するために Toolsets を学ぶ
- agent-as-tool パターンのために Agent Composition を読む
- トランスクリプト永続化のために Memory & Sessions を読む