管理エージェントランタイム
安定性: 実験的 — この機能は追加機能であり、
managed-runtimeによって機能ゲートされています。 機能が無効な場合、既存のRunner/LlmAgentAPIs には影響しません。 API のインターフェースは、将来のリリースで変更される可能性があります。
概要
Managed Agent Runtime(adk-managed)は、プロバイダーに依存せず、永続化および再開が可能な
エージェント実行エンジンです。宣言的な ManagedAgentDef を受け取り、実行可能な
エージェントを構築し、チェックポイントから再開でき、イベントをストリーミングするバックグラウンドセッションとして実行します。
ランタイムはサービスではなく、ライブラリです。プラットフォームがこれをホストします。つまり、次のようになります。
- 単独でテスト可能: HTTP/認証/課金への依存関係がありません
- 組み込み可能: セルフホスト型のデプロイでは、同じランタイムトレイトを直接使用します
- プラットフォームを交換可能: 異なるプラットフォームで同じランタイムをホストできます
- プロバイダーに依存しない: モデルプロバイダーにかかわらず、イベントシーケンスは同一です
クイックスタート
Cargo.toml に機能を追加します。
[dependencies]
adk-rust = { version = "2.1.0", features = ["managed-runtime"] }
または、adk-managed クレートを直接使用します。
[dependencies]
adk-managed = "2.1.0"
adk-session = "2.1.0"
最小限の例(ScriptedLlm — API キーなし)
use std::sync::Arc;
use adk_managed::{
DefaultManagedAgentRuntime, ManagedAgentRuntime, ModelResolver,
ScriptedLlm, ScriptedTurn,
resolver::ResolverResult,
types::{ContentBlock, ManagedAgentDef, ModelRef, UserEvent},
};
use adk_session::InMemorySessionService;
use async_trait::async_trait;
use futures::StreamExt;
// A resolver that returns our scripted LLM
struct MockResolver { llm: Arc<dyn adk_core::Llm> }
#[async_trait]
impl ModelResolver for MockResolver {
async fn resolve(&self, _: &ModelRef) -> ResolverResult<Arc<dyn adk_core::Llm>> {
Ok(self.llm.clone())
}
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
// 1. Create a scripted LLM (deterministic, offline, $0)
let llm = Arc::new(ScriptedLlm::new("test-model", vec![
ScriptedTurn { text: Some("Hello!".into()), tool_calls: vec![] },
]));
// 2. Build the runtime
let runtime = DefaultManagedAgentRuntime::new(
Arc::new(MockResolver { llm }),
Arc::new(InMemorySessionService::new()),
);
// 3. Create an agent
let def = ManagedAgentDef::new("my-agent", ModelRef::Shorthand("test-model".into()))
.with_system("You are helpful.");
let agent = runtime.create(def).await?;
// 4. Start a session
let session = runtime.start_session(&agent, None).await?;
// 5. Subscribe to events and send a message
let mut stream = runtime.stream_events(&session, None).await?;
runtime.send_event(&session, UserEvent::Message {
content: vec![ContentBlock::Text { text: "Hi!".into() }],
}).await?;
// 6. Collect events
while let Some(event) = stream.next().await {
println!("{event:?}");
}
Ok(())
}
アーキテクチャ
┌─────────────────────────────────────────────────────────────┐
│ Platform Layer (ep-* crates) │
│ HTTP Routes │ Auth │ Billing │ Multi-tenancy │
└──────────────────────────┬──────────────────────────────────┘
│ Rust trait calls (in-process)
▼
┌─────────────────────────────────────────────────────────────┐
│ Runtime Layer (adk-managed) │
│ │
│ ManagedAgentRuntime trait + DefaultManagedAgentRuntime │
│ ─────────────────────────────────────────────────── │
│ • Builds runnable agents from ManagedAgentDef │
│ • Runs supervised session loop (durable, resumable) │
│ • Emits provider-neutral SessionEvent stream │
│ • Manages custom tool parking, checkpoints, interrupts │
│ • Resolves ModelRef → Arc<dyn Llm> │
│ │
│ Composes existing crates: │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌────────────┐ │
│ │adk-runner│ │adk-session│ │adk-model │ │adk-tool │ │
│ └──────────┘ └──────────┘ └──────────┘ └────────────┘ │
└─────────────────────────────────────────────────────────────┘
コア型
ManagedAgentRuntime トレイト
エージェントのライフサイクル全体を定義する中心的な非同期トレイトです。
| メソッド | 説明 |
|---|---|
create(def) | エージェント定義を登録し、AgentHandle を返します |
start_session(agent, env?) | 新しいセッションを開始します。初期ステータスは Queued です |
send_event(session, event) | セッションにUserEventを送信 |
stream_events(session, from_seq?) | SessionEventストリームを購読 |
interrupt(session) | 次の境界で停止し、status.idleを出力 |
pause(session) | チェックポイントを作成して処理を一時停止 |
resume(session) | 一時停止または再起動から再開 |
status(session) | 現在の SessionStatus を照会 |
archive(session) | 終端状態、データ保持 |
delete_session(session) | セッションデータを削除 |
ManagedAgentDef
builder APIを使用した宣言型エージェント定義:
let def = ManagedAgentDef::new("my-agent", ModelRef::Shorthand("gemini-3.7-flash".into()))
.with_system("You are a helpful assistant.")
.with_description("Research agent with web search")
.with_tools(vec![ToolConfig::BuiltIn(ManagedBuiltinTool::WebSearch)]);
SessionEvent
単調増加するシーケンス番号を備えた、プロバイダーに依存しないイベントストリーム:
agent.message— アシスタントのテキストコンテンツagent.tool_use— 組み込みツールの呼び出しagent.custom_tool_use— クライアントで実行されるカスタムツール(ループは一時停止)agent.mcp_tool_use— MCPツールの呼び出しstatus.running— ターン開始status.idle— ターン完了(stop_reasonを伴う)error— 実行エラー
UserEvent
クライアントからエージェントへのイベント:
user.message— エージェントにコンテンツを送信user.interrupt— 現在のターンを停止user.tool_confirmation— ツール実行を許可または拒否user.custom_tool_result— カスタムツールの結果を返すuser.tool_result— 組み込みツールの結果(セルフホストの場合のみ)user.define_outcome— 成功条件を設定
ModelRef
すべてのプロバイダーをサポートする、プロバイダーに依存しないモデル参照:
// Shorthand (provider inferred from name)
ModelRef::Shorthand("gemini-3.7-flash".into())
ModelRef::Shorthand("gpt-5.6-terra".into())
ModelRef::Shorthand("claude-sonnet-5".into())
// Structured (explicit provider)
ModelRef::Structured {
provider: Provider::OpenaiCompatible,
model: ModelConfig::Compatible {
model: "my-model".into(),
base_url: "https://my-endpoint.com/v1".into(),
api_key: "sk-...".into(),
},
speed: None,
}
主な機能
永続的なセッション
すべてのイベントはアトミックにチェックポイント化されます。プロセスがクラッシュした場合、resume()は
最後の一貫したチェックポイントから再構築され、イベントが失われることはありません:
// Before crash: events 0..5 committed
// After restart:
runtime.resume(&session).await?;
// Continues from seq=5, no gap, no duplicate
カスタムツールの一時停止
エージェントがagent.custom_tool_useを出力すると、クライアントが結果を返すか、設定可能なタイムアウトが
経過するまでループは一時停止します:
// Agent emits: agent.custom_tool_use { custom_tool_use_id: "ct_1", name: "deploy" }
// Client executes the tool, then:
runtime.send_event(&session, UserEvent::CustomToolResult {
custom_tool_use_id: "ct_1".into(),
content: vec![ContentBlock::Text { text: "Deployed successfully".into() }],
}).await?;
イベントの再生
シーケンスに基づく再生により、SSE Last-Event-IDの再接続をサポートします:
// Reconnect from seq 42 — replays events 43, 44, ... then live tail
let stream = runtime.stream_events(&session, Some(42)).await?;
プロバイダー間の同等性
同一のManagedAgentDefにより、Gemini、OpenAI、Anthropic、Ollama、およびOpenAI互換プロバイダー間で、バイト単位で同一のイベント型シーケンスが生成されます(fixture F-8)。
ScriptedLlmによるテスト
ScriptedLlmは、ランタイムパイプライン全体を実行する決定論的なLLMダブルです。置き換えられるのは、プロバイダーのAPI呼び出しのみです:
use adk_managed::testing::{ScriptedLlm, ScriptedTurn, ScriptedToolCall};
use serde_json::json;
let llm = ScriptedLlm::new("test", vec![
ScriptedTurn {
text: Some("I'll search for that.".into()),
tool_calls: vec![ScriptedToolCall {
name: "web_search".into(),
input: json!({"query": "rust agents"}),
id: Some("tc_1".into()),
}],
},
ScriptedTurn {
text: Some("Here are the results...".into()),
tool_calls: vec![],
},
]);
APIリファレンス
完全なAPIドキュメントは docs.rs で利用できます:
スモークテストの例
プラットフォームチーム向けに、スタンドアロンのサンプルクレートが用意されています。
cargo run --manifest-path examples/managed_runtime_hello/Cargo.toml
これは、ScriptedLlm を使用してフィクスチャ F-1 をエンドツーエンドで実行します(API キーは不要です)。
管理対象状態の永続性
管理対象のセッション状態(イベントログ、シーケンス位置、保留中のツール呼び出し、ライフサイクルのステータス)は、ManagedStateStore に格納されます。ストアは独自の保証を報告します。
| 耐久性 | 意味 |
|---|---|
ProcessLocal | プロセスの実行中は再生と再開が機能します。クラッシュすると状態が失われ、別のプロセスでセッションを再開することはできません。 |
CrashDurable | 書き込みが確認応答される前に状態がバッキングストアに書き込まれるため、別のプロセスでセッションを再構築できます。 |
InMemoryManagedStateStore のみが同梱され、それは ProcessLocal です。 チェックポイント機能の存在から推測するのではなく、保証内容を確認してください。
use adk_managed::{Durability, InMemoryManagedStateStore, ManagedStateStore};
let store = InMemoryManagedStateStore::new();
assert_eq!(store.durability(), Durability::ProcessLocal);
assert!(!store.durability().survives_process_loss());
チェックポイント処理とフラッシュ処理の違い
CheckpointManager::checkpoint はイベントと新しい実行状態を一緒に記録するため、リプレイ時にどちらか一方だけが存在することはありません。これはマネージャー自身のフィールドへの書き込みです。flush はスナップショットを設定されたストアに書き込み、restore はそこからマネージャーを再構築します。
use adk_managed::{CheckpointManager, InMemoryManagedStateStore, ManagedStateStore};
use std::sync::Arc;
# async fn example() -> Result<(), adk_managed::types::RuntimeError> {
let store: Arc<dyn ManagedStateStore> = Arc::new(InMemoryManagedStateStore::new());
let manager = CheckpointManager::new("session-1".to_string()).with_store(Arc::clone(&store));
manager.flush().await?;
let restored = CheckpointManager::restore("session-1".to_string(), store).await?;
assert_eq!(restored.session_id(), "session-1");
# Ok(())
# }
ステータスの報告
ManagedAgentRuntime::status はセッションループが書き込むものと同じハンドルを読み取るため、コントロールプレーンの遷移だけでなく、通常の遷移も確認できます。
| 遷移 | 原因 |
|---|---|
Queued → Running | ターンが開始する |
Running → Idle | ターンが完了し、使用量が記録される |
任意 → Paused | pause |
任意 → Archived | archive または delete_session |
注: 以前はこれが1つの共有ハンドルでした。
statusは、ターンの実行中も含め、セッションの存続期間全体にわたってQueuedを報告していました。制御プレーンの遷移(pause、resume、archive)は、ハンドルに直接書き込まれていたため確認できました。
削除のセマンティクス
delete_session は両方のプレーンを削除します。
- セッションを終了状態にし、そのループをキャンセルします。
- ランタイムハンドルを削除します。
- 注入された
SessionServiceを通じて、永続化された会話を、作成時に使用されたものと同じアイデンティティstart_sessionのもとで削除します。
ステップ3が失敗した場合、delete_session は、データを保持しているアプリ、ユーザー、セッションを明示したエラーを返します。その時点ではハンドルはすでに消失しているため、呼び出し元には成功したと思い込ませるのではなく、手動で何をクリーンアップする必要があるかを伝える必要があります。
runtime.delete_session(&session).await?;
// The handle is gone and the conversation is no longer in the session backend.
重要: 削除すると、所有者のもとにあるセッションの会話が削除されます。 セッションの所有権を参照してください。
セッションの所有権
start_session には ManagedOwner が必要です。セッションはそのアイデンティティのもとで永続化され、セッションループが行うすべての Runner 呼び出しでそのアイデンティティが使用されます。
use adk_managed::{ManagedAgentRuntime, ManagedOwner};
# async fn start(runtime: &dyn ManagedAgentRuntime, agent: &adk_managed::AgentHandle)
# -> Result<(), adk_managed::RuntimeError> {
let owner = ManagedOwner::new("support-console", "user-42")?;
let session = runtime.start_session(agent, &owner, None).await?;
# Ok(())
# }
重要:
checkpointは「アトミックに永続化する」と記述され、「どのようなクラッシュの後でも、リプレイは一貫したビューを参照できる」という保証があり、読み込みは「再起動後にセッションを再構築するために必要なすべて」を返すと説明されていました。しかし、どちらも実際には成立していませんでした。両方とも、永続ストアに対するトランザクションなしで、インメモリフィールドを操作していました。提供されたストアでは、新しいプロセスでrestoreを実行しても何も見つかりません。 両方のコンポーネントが必要であり、空白以外の値でなければなりません。異なる所有者に属するセッションは個別にアドレス指定されるため、検索と削除は1人の所有者にスコープされ、別の所有者のデータに到達することはできません。
注記: すべての管理対象セッションは以前、定数
managed/managed_userの下に永続化されていたため、すべてが1つの論理名前空間を共有していました。呼び出し元ごとにスコープを設定することはできず、どのセッションがどの呼び出し元に属するかを特定することもできませんでした。
環境設定
EnvironmentConfig は env_vars と working_dir を保持します。このランタイムは、次のいずれかを要求する設定を拒否します。
invalid request: EnvironmentConfig cannot be honoured by this runtime: sessions run
in-process, so per-session environment variables and working directories would have to mutate
process-global state shared with other sessions. Pass `None`, or configure a sandboxed runtime.
セッションはプロセス内で実行されるため、セッションごとの環境変数や作業ディレクトリを適用すると、他のすべてのセッションと共有される状態が変更されます。拒否することが正直な対応です。要求を満たすには、サンドボックス化された実行境界が必要です。
注記: この引数は以前
_envという名前で、破棄されていました。そのため、環境設定を指定した呼び出し元には、それを黙って無視するセッションが返されていました。