托管智能体运行时
稳定性:实验性 — 此功能为增量功能,并受
managed-runtime特性门控。 在禁用该功能时,不会影响现有的Runner/LlmAgentAPIs。 API 接口未来版本可能会发生变化。
概述
托管智能体运行时(adk-managed)是一个与提供商无关、持久化且可恢复的
智能体执行引擎。它接收声明式的 ManagedAgentDef,构建可运行的智能体,并将其作为
可通过检查点恢复、流式传输事件的后台会话运行。
该运行时是一个库,而不是服务。平台负责托管它。这意味着:
- 可独立测试:不依赖任何 HTTP/身份验证/计费
- 可嵌入:自托管部署直接使用相同的运行时 trait
- 可替换的平台:不同平台可以托管相同的运行时
- 与提供商无关:无论模型提供商为何,事件序列都完全一致
快速开始
将该特性添加到你的 Cargo.toml 中:
[dependencies]
adk-rust = { version = "2.1.0", features = ["managed-runtime"] }
或者直接使用 adk-managed crate:
[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 Trait
定义完整智能体生命周期的核心异步 trait:
| 方法 | 描述 |
|---|---|
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
使用构建器 API 的声明式 agent 定义:
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
与 provider 无关、带有单调递增序列号的事件流:
agent.message— Assistant 文本内容agent.tool_use— 内置工具调用agent.custom_tool_use— 由客户端执行的自定义工具(循环暂停)agent.mcp_tool_use— MCP 工具调用status.running— Turn 开始status.idle— Turn 完成(包含stop_reason)error— 执行错误
UserEvent
客户端到 agent 的事件:
user.message— 向 agent 发送内容user.interrupt— 停止当前 turnuser.tool_confirmation— 允许/拒绝工具执行user.custom_tool_result— 返回自定义工具结果user.tool_result— 内置工具结果(仅限自托管)user.define_outcome— 设置成功标准
ModelRef
支持所有 provider、与 provider 无关的模型引用:
// 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 发出 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?;
Provider 对等性
相同的 ManagedAgentDef 在 Gemini、OpenAI、Anthropic、Ollama 以及兼容 OpenAI 的 provider 中,会生成字节级一致的事件类型序列(fixture F-8)。
使用 ScriptedLlm 进行测试
ScriptedLlm 是一个确定性的 LLM 测试替身,用于执行完整的运行时流水线。仅替换 provider 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 上查看:
烟雾测试示例
为平台团队提供了一个独立的示例 crate:
cargo run --manifest-path examples/managed_runtime_hello/Cargo.toml
该示例使用 ScriptedLlm 端到端运行 fixture 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 |
注意: 在此之前,这是一个共享句柄,
status会在会话的整个生命周期内报告Queued,包括会话正在执行各轮操作时。控制平面转换(暂停、恢复、归档)之所以可见,是因为它们会直接写入句柄。
删除语义
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在新进程中找不到任何内容。
这两个组件都是必需的,并且不能留空。属于不同所有者的会话会分别寻址,因此查找和删除都限定在单个所有者范围内,无法访问其他所有者的数据。
注意: 每个受管理的会话此前都持久化在常量
managed/managed_user下,因此它们共享同一个逻辑命名空间:无法将任何内容限定到某个调用方,也无法将任何会话归属于某个调用方。
环境配置
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,但会被丢弃,因此提供环境配置的调用方获得的会话会静默忽略该配置。