新增独立 Agent Runtime Rust 工作区

新增 Core、Engine、Runtime、SQLite、Provider、MCP、Skill、Codex、CLI 与 DAG crate

补齐 OpenAI endpoint 配置、Provider 实例/协议路由和统一工具权限边界

加入持久化、lease、checkpoint、reconciliation、审批恢复与消息历史回归

加入独立 workspace CI、依赖边界、能力集和 Fake Agent 测试脚本

同步建设计划、TODO、架构、测试与验收文档
This commit is contained in:
2026-09-06 17:44:54 +08:00
parent 3e8591040d
commit 202279c6d9
90 changed files with 79439 additions and 0 deletions
@@ -0,0 +1,88 @@
name: Agent Runtime CI
on:
push:
pull_request:
workflow_dispatch:
permissions:
contents: read
env:
CARGO_INCREMENTAL: '0'
CARGO_NET_RETRY: '10'
CARGO_TERM_COLOR: always
RUSTUP_AUTO_INSTALL: '0'
jobs:
rust:
name: Rust workspace
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- name: Prepare isolated temporary directory
shell: bash
run: |
mkdir -p "$HOME/data/tmp"
echo "TMPDIR=$HOME/data/tmp" >> "$GITHUB_ENV"
# rust-toolchain.toml 固定编译器和组件;runner 镜像需预装它们,
# 这里先失败得更明确,不让后续 Clippy 步骤才暴露环境缺口。
- name: Verify pinned Rust toolchain
shell: bash
run: |
rustc --version
cargo fmt --version
cargo clippy --version
- name: Check formatting
run: cargo fmt --all -- --check
- name: Check warnings and targets
run: RUSTFLAGS='-D warnings' cargo check --locked --workspace --all-targets --all-features
- name: Check no-default-features build
run: RUSTFLAGS='-D warnings' cargo check --locked --workspace --all-targets --no-default-features
# The workspace compatibility matrix includes the SQLite service through
# Host. Test the portable Runtime package in isolation so this gate
# proves the physical dependency boundary rather than relying on feature
# unification in the full workspace.
- name: Check portable Runtime without SQLite
shell: bash
run: |
if cargo tree --locked --no-default-features -p agent-runtime -e normal \
| grep -Eiq 'agent-storage-sqlite|rusqlite|libsqlite3-sys'; then
echo 'portable agent-runtime unexpectedly contains SQLite' >&2
exit 1
fi
RUSTFLAGS='-D warnings' cargo check --locked -p agent-runtime --all-targets --no-default-features
cargo test --locked -p agent-runtime --all-targets --no-default-features --no-fail-fast
RUSTDOCFLAGS='-D warnings' cargo doc --locked -p agent-runtime --no-default-features --no-deps
- name: Verify kernel dependency boundary
shell: bash
run: |
./scripts/check-dependencies.sh Cargo.toml
- name: Verify package manifests
shell: bash
run: |
./scripts/check-package-manifests.sh Cargo.toml
# cargo-audit 和 RustSec advisory DB 由 runner 预装/挂载;本步骤刻意不
# 执行 cargo install、git fetch 或其它联网更新。可用 CARGO_AUDIT_BIN
# 指向 runner 固定版本的二进制,并必须提供 RUSTSEC_ADVISORY_DB。
- name: Run offline cargo-audit
shell: bash
run: ./scripts/run-cargo-audit.sh
- name: Verify independent workspace copy
shell: bash
run: ./scripts/verify-independent-workspace.sh
- name: Run tests
run: cargo test --locked --workspace --all-features --no-fail-fast
- name: Run no-default-features tests
run: cargo test --locked --workspace --no-default-features --no-fail-fast
- name: Run deterministic agent test set
shell: bash
run: ./scripts/run-agent-test-set.sh --quick
- name: Build documentation without warnings
run: RUSTDOCFLAGS='-D warnings' cargo doc --locked --workspace --all-features --no-deps
- name: Run Clippy policy gate
run: cargo clippy --locked --workspace --all-features --all-targets -- -D warnings
- name: Run no-default-features Clippy policy gate
run: cargo clippy --locked --workspace --no-default-features --all-targets -- -D warnings
- name: Run portable Runtime Clippy policy gate
run: cargo clippy --locked -p agent-runtime --no-default-features --all-targets -- -D warnings
+8
View File
@@ -0,0 +1,8 @@
/target/
*.db
*.db-shm
*.db-wal
.env
agent.toml
.tmp-test/
.tmp-cli.*/
+1857
View File
File diff suppressed because it is too large Load Diff
+46
View File
@@ -0,0 +1,46 @@
[workspace]
resolver = "2"
members = [
"crates/agent-runtime-core",
"crates/agent-runtime-contracts",
"crates/agent-runtime-engine",
"crates/agent-runtime",
"crates/agent-runtime-sqlite",
"crates/agent-provider-openai",
"crates/agent-provider-fake",
"crates/agent-codex",
"crates/agent-runtime-orchestration",
"crates/agent-storage-sqlite",
"crates/agent-mcp",
"crates/agent-skills",
"crates/agent-app",
"crates/agent-host",
"crates/agent-cli",
]
[workspace.package]
edition = "2024"
version = "0.1.0"
rust-version = "1.96"
license = "UNLICENSED"
[workspace.dependencies]
agent-runtime-core = { path = "crates/agent-runtime-core", version = "0.1.0" }
agent-runtime-contracts = { path = "crates/agent-runtime-contracts", version = "0.1.0" }
agent-runtime-engine = { path = "crates/agent-runtime-engine", version = "0.1.0" }
agent-runtime = { path = "crates/agent-runtime", version = "0.1.0" }
agent-runtime-sqlite = { path = "crates/agent-runtime-sqlite", version = "0.1.0" }
agent-provider-openai = { path = "crates/agent-provider-openai", version = "0.1.0" }
agent-provider-fake = { path = "crates/agent-provider-fake", version = "0.1.0" }
agent-codex = { path = "crates/agent-codex", version = "0.1.0" }
agent-runtime-orchestration = { path = "crates/agent-runtime-orchestration", version = "0.1.0" }
agent-storage-sqlite = { path = "crates/agent-storage-sqlite", version = "0.1.0" }
agent-mcp = { path = "crates/agent-mcp", version = "0.1.0" }
agent-skills = { path = "crates/agent-skills", version = "0.1.0" }
agent-app = { path = "crates/agent-app", version = "0.1.0" }
agent-host = { path = "crates/agent-host", version = "0.1.0" }
serde = { version = "1", features = ["derive"] }
serde_json = "1"
sha2 = "0.10"
toml = "0.8"
thiserror = "2"
+634
View File
File diff suppressed because it is too large Load Diff
+51
View File
@@ -0,0 +1,51 @@
# 复制为 agent.toml 后即可运行 `cargo run -p agent-cli -- run "任务"`。
# 认证字段只保存环境变量名;不要把 api_key、token 或 Cookie 写进此文件。
db = "./agent.db"
provider = "fake"
model = "fake"
stream = false
# provider = "openai" 时可使用以下任一 endpoint 配置:
# openai_base_url = "https://gateway.example/v1" # 自动补 /responses
# openai_endpoint = "https://gateway.example/v1/responses" # 完整地址,优先
# openai_api_key_env = "OPENAI_API_KEY" # 只保存环境变量名
# OPENAI_MODEL 或 AGENT_MODEL 可覆盖 model
# 可选的固定提示词 section;system/developer/context 不会被压成同一条 user 消息。
# system_prompt = "你是一个简洁的助手"
# developer_prompt = "输出可审计的步骤"
# context_prompt = "这是不可信的外部背景"
[skills]
# roots = ["./.codex/skills", "./.agents/skills"]
# names = ["review"]
[mcp]
# server = "workspace"
# stdio_command = "npx"
# stdio_args = ["-y", "@modelcontextprotocol/server-filesystem", "."]
# timeout_secs = 30
# allow = ["read_file"]
# 只有显式列出的内容会被读取并作为不可信上下文注入;默认不读取资源或 prompt。
# context_resources = ["file:///workspace/README.md"]
# context_prompts = ["welcome"] # CLI 使用空参数;带参数请用库 API
# HTTP 或 stdio 的秘密均通过环境变量引用。示例:
# [[mcp.auth]]
# variable = "MCP_TOKEN"
# target = "http_bearer"
#
# stdio 服务器也可以把环境变量转发给子进程:
# [[mcp.auth]]
# variable = "MCP_WORKSPACE_TOKEN"
# target = "stdio_environment"
# name = "WORKSPACE_TOKEN"
# Codex 是显式外部 backend,不替换当前 CLI 的 ModelProvider。
# `agent codex validate` 只校验配置和参数白名单,不启动进程。
# [codex.cli]
# program = "codex"
# args = ["--model=gpt-5-codex"]
# allowed_arg_prefixes = ["--model"]
# timeout_ms = 120000
# max_output_bytes = 1048576
+15
View File
@@ -0,0 +1,15 @@
[package]
name = "agent-app"
version.workspace = true
edition.workspace = true
rust-version.workspace = true
license.workspace = true
description = "通用 Agent 程序的配置与装配边界"
[dependencies]
agent-codex.workspace = true
agent-mcp.workspace = true
agent-provider-openai.workspace = true
serde.workspace = true
serde_json.workspace = true
toml.workspace = true
+332
View File
@@ -0,0 +1,332 @@
//! 通用 Agent 程序的无状态配置与装配输入。
//!
//! 这个 crate 只承接 `agent.toml`、环境变量和非秘密路由 metadata 的解析。
//! 它不依赖 Host、Runtime、线程或数据库,因此 CLI 之外的入口也可以复用
//! 同一套配置优先级,而不会复制一份运行状态机。
use std::env;
use std::fmt;
use std::fs;
use std::path::{Path, PathBuf};
use agent_codex::CodexCliConfig;
use agent_mcp::{McpAuthEnv, McpAuthTarget};
use agent_provider_openai::OpenAiProviderConfig;
use serde::{Deserialize, Serialize};
use serde_json::{Value, json};
/// `agent.toml` 只描述可复现的装配参数;密钥字段故意没有对应结构,
/// 未知字段会被拒绝,避免用户误把明文 token 写入配置文件。
#[derive(Clone, Debug, Default, Deserialize, Serialize)]
#[serde(default, deny_unknown_fields)]
pub struct AgentTomlConfig {
pub db: Option<String>,
pub provider: Option<String>,
pub model: Option<String>,
pub stream: Option<bool>,
pub openai_api_key_env: Option<String>,
/// OpenAI-compatible 网关的 base URL;运行时会自动补 /responses。
pub openai_base_url: Option<String>,
/// OpenAI Responses 的完整 endpoint,优先于 base URL。
pub openai_endpoint: Option<String>,
pub system_prompt: Option<String>,
pub developer_prompt: Option<String>,
pub context_prompt: Option<String>,
#[serde(default)]
pub skills: SkillTomlConfig,
#[serde(default)]
pub mcp: McpTomlConfig,
#[serde(default)]
pub codex: CodexTomlConfig,
}
#[derive(Clone, Debug, Default, Deserialize, Serialize)]
#[serde(default, deny_unknown_fields)]
pub struct SkillTomlConfig {
pub roots: Vec<String>,
pub names: Vec<String>,
}
#[derive(Clone, Debug, Default, Deserialize, Serialize)]
#[serde(default, deny_unknown_fields)]
pub struct McpTomlConfig {
pub server: Option<String>,
pub stdio_command: Option<String>,
pub stdio_args: Vec<String>,
pub http_url: Option<String>,
/// 认证只保存环境变量引用;解析后的 token 由 agent-mcp 在连接时读取。
#[serde(default)]
pub auth: Vec<McpAuthToml>,
pub timeout_secs: Option<u64>,
pub allow: Vec<String>,
/// 只读取这些明确列出的 resource URI 作为不可信上下文。
pub context_resources: Vec<String>,
/// 只展开这些明确列出的 prompt(CLI 使用空参数;需要参数时使用库 API)。
pub context_prompts: Vec<String>,
}
/// Codex 仍是显式外部 backend,不替换 CLI 的 ModelProvider。这里先把一次性
/// CLI 的受限启动配置纳入同一份 agent.toml,并由 `codex validate` 做无副作用
/// 校验;App Server channel 继续由嵌入方注入,避免配置层偷藏第二套运行状态。
#[derive(Clone, Debug, Default, Deserialize, Serialize)]
#[serde(default, deny_unknown_fields)]
pub struct CodexTomlConfig {
pub cli: Option<CodexCliConfig>,
}
/// agent.toml 中的认证引用保持扁平、可读的写法:
/// `target = "http_bearer"`,真正的 Core target 在连接前才构造。
/// 结构里没有 secret 字段,故配置序列化和 Debug 都不会持有凭据原文。
#[derive(Clone, Default, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
pub struct McpAuthToml {
pub variable: String,
pub target: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub name: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub prefix: Option<String>,
}
impl fmt::Debug for McpAuthToml {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("McpAuthToml")
.field("variable", &"<env-ref>")
.field("target", &self.target)
.field("name", &self.name)
.field("prefix", &self.prefix.as_deref().map(|_| "<redacted>"))
.finish()
}
}
impl McpAuthToml {
/// 将非秘密 TOML 引用转换为 MCP 的 Core 认证引用。
pub fn into_core(self) -> Result<McpAuthEnv, Box<dyn std::error::Error>> {
let target = match self.target.as_str() {
"http_bearer" => {
if self.name.is_some() || self.prefix.is_some() {
return Err("http_bearer 认证引用不应带 name/prefix 字段"
.to_owned()
.into());
}
McpAuthTarget::HttpBearer
}
"http_header" => {
let name = self
.name
.filter(|name| !name.trim().is_empty())
.ok_or("http_header 认证引用需要非空 name")?;
McpAuthTarget::HttpHeader {
name,
prefix: self.prefix.unwrap_or_default(),
}
}
"stdio_environment" => {
let name = self
.name
.filter(|name| !name.trim().is_empty())
.ok_or("stdio_environment 认证引用需要非空 name")?;
McpAuthTarget::StdioEnvironment { name }
}
_ => {
return Err(
"MCP 认证 target 无效(支持 http_bearer/http_header/stdio_environment)"
.to_owned()
.into(),
);
}
};
if self.variable.trim().is_empty() {
return Err("MCP 认证 variable 不能为空".to_owned().into());
}
Ok(McpAuthEnv {
variable: self.variable,
target,
})
}
}
impl AgentTomlConfig {
/// 从 `AGENT_CONFIG`(缺省为当前目录 `agent.toml`)读取配置。
pub fn load() -> Result<Self, Box<dyn std::error::Error>> {
let path = env::var_os("AGENT_CONFIG")
.map(PathBuf::from)
.unwrap_or_else(|| PathBuf::from("agent.toml"));
Self::load_from_path(path)
}
/// 从明确路径加载配置;调用方可用它避免在测试中修改全局环境。
pub fn load_from_path(path: impl AsRef<Path>) -> Result<Self, Box<dyn std::error::Error>> {
let path = path.as_ref();
if !path.exists() {
return Ok(Self::default());
}
let text = fs::read_to_string(path)?;
toml::from_str(&text)
.map_err(|error| format!("配置文件 {} 无效: {error}", path.display()).into())
}
/// 解析数据库路径:进程环境覆盖 TOML,最后使用 `agent.db`。
pub fn db_path(&self) -> PathBuf {
env::var_os("AGENT_DB")
.map(PathBuf::from)
.or_else(|| self.db.as_deref().map(PathBuf::from))
.unwrap_or_else(|| PathBuf::from("agent.db"))
}
/// 解析 Provider 名称。没有显式选择时,有可用 OpenAI key 才默认 openai。
pub fn provider(&self) -> String {
non_empty_env("AGENT_PROVIDER")
.or_else(|| self.provider.clone())
.unwrap_or_else(|| {
let key_env = self.openai_api_key_env();
if env::var(key_env)
.map(|key| !key.trim().is_empty())
.unwrap_or(false)
{
"openai".to_owned()
} else {
"fake".to_owned()
}
})
}
/// 解析模型名称,并让环境变量覆盖配置文件。
pub fn model(&self) -> String {
non_empty_env("AGENT_MODEL")
.or_else(|| {
(self.provider() == "openai")
.then(|| non_empty_env("OPENAI_MODEL"))
.flatten()
})
// Treat a blank TOML value like an unset override, matching the
// environment helpers and allowing the normal provider default.
.or_else(|| self.model.clone().filter(|value| !value.trim().is_empty()))
.unwrap_or_else(|| "fake".to_owned())
}
/// API key 只解析环境变量名,从不读取或保存 key 原文。
pub fn openai_api_key_env(&self) -> String {
non_empty_env("AGENT_OPENAI_API_KEY_ENV")
.or_else(|| non_empty_env("OPENAI_API_KEY_ENV"))
.or_else(|| self.openai_api_key_env.clone())
.unwrap_or_else(|| "OPENAI_API_KEY".to_owned())
}
/// 按环境变量优先级构造 OpenAI 非秘密配置。
pub fn openai_provider_config(&self) -> OpenAiProviderConfig {
self.openai_provider_config_with_env(
non_empty_env("OPENAI_ENDPOINT").as_deref(),
non_empty_env("OPENAI_BASE_URL").as_deref(),
)
}
/// 明确传入环境候选,便于嵌入方做纯函数测试而不修改进程环境。
/// 环境完整 endpoint 优先于环境 base URL;任一环境来源都优先于 TOML。
pub fn openai_provider_config_with_env(
&self,
env_endpoint: Option<&str>,
env_base_url: Option<&str>,
) -> OpenAiProviderConfig {
let mut config =
OpenAiProviderConfig::default().with_api_key_env(self.openai_api_key_env());
if let Some(endpoint) = env_endpoint
.filter(|value| !value.trim().is_empty())
.map(str::to_owned)
{
config = config.with_endpoint(endpoint);
} else if let Some(base_url) = env_base_url
.filter(|value| !value.trim().is_empty())
.map(str::to_owned)
{
config = config.with_base_url(base_url);
} else if let Some(endpoint) = self
.openai_endpoint
.clone()
.filter(|value| !value.trim().is_empty())
{
config = config.with_endpoint(endpoint);
} else if let Some(base_url) = self
.openai_base_url
.clone()
.filter(|value| !value.trim().is_empty())
{
config = config.with_base_url(base_url);
}
config
}
/// 解析流式开关;显式环境值优先,OpenAI 缺省开启流式。
pub fn streaming(&self) -> bool {
if let Some(value) = non_empty_env("AGENT_STREAM") {
return matches!(value.as_str(), "1" | "true" | "yes" | "on");
}
self.stream.unwrap_or_else(|| self.provider() == "openai")
}
}
/// 将 provider 与实际模型写成后台 queued run 的非秘密路由观察。
/// worker 仍会在真正执行前重新装配 Provider,不会从 metadata 读取凭据。
pub fn queued_run_metadata(config: &AgentTomlConfig, provider: &str) -> Value {
json!({
"provider": effective_model(config, provider),
"providerKind": provider,
})
}
/// 将配置模型解析为实际执行模型;OpenAI 的 `fake` 只是未指定模型的哨兵。
pub fn effective_model(config: &AgentTomlConfig, provider: &str) -> String {
let model = config.model();
if provider == "openai" && model == "fake" {
"gpt-4.1-mini".to_owned()
} else {
model
}
}
fn non_empty_env(name: &str) -> Option<String> {
env::var(name).ok().filter(|value| !value.trim().is_empty())
}
#[cfg(test)]
mod tests {
use super::{AgentTomlConfig, effective_model, queued_run_metadata};
#[test]
fn blank_model_uses_provider_default_without_secret_values() {
let config: AgentTomlConfig =
toml::from_str("provider = 'openai'\nmodel = ' '").expect("配置应可解析");
assert_eq!(config.model(), "fake");
assert_eq!(effective_model(&config, "openai"), "gpt-4.1-mini");
assert_eq!(
queued_run_metadata(&config, "openai"),
serde_json::json!({"provider": "gpt-4.1-mini", "providerKind": "openai"})
);
}
#[test]
fn load_missing_path_returns_default() {
let config = AgentTomlConfig::load_from_path(
std::env::temp_dir().join("agent-app-missing-config-does-not-exist.toml"),
)
.expect("缺失配置应使用默认值");
assert!(config.db.is_none());
assert_eq!(config.provider, None);
}
#[test]
fn endpoint_prefers_environment_before_toml() {
let config: AgentTomlConfig = toml::from_str(
"provider = 'openai'\nopenai_endpoint = 'https://toml.example/v1/responses'\nopenai_base_url = 'https://toml-base.example/v1'",
)
.expect("配置应可解析");
assert_eq!(
config
.openai_provider_config_with_env(None, Some("https://env.example/v1"))
.resolve_endpoint()
.expect("endpoint 应可解析"),
"https://env.example/v1/responses"
);
}
}
+24
View File
@@ -0,0 +1,24 @@
[package]
name = "agent-cli"
version.workspace = true
edition.workspace = true
rust-version.workspace = true
license.workspace = true
description = "通用 Agent 单智能体命令行程序"
[[bin]]
name = "agent"
path = "src/main.rs"
[dependencies]
agent-app.workspace = true
agent-codex.workspace = true
agent-host.workspace = true
agent-mcp.workspace = true
agent-provider-openai.workspace = true
agent-runtime-core.workspace = true
agent-runtime-engine.workspace = true
agent-skills.workspace = true
serde.workspace = true
serde_json.workspace = true
toml.workspace = true
File diff suppressed because it is too large Load Diff
+13
View File
@@ -0,0 +1,13 @@
[package]
name = "agent-codex"
version.workspace = true
edition.workspace = true
rust-version.workspace = true
license.workspace = true
description = "Codex CLI 与 App Server 的外部 Agent 适配器"
[dependencies]
agent-runtime-core.workspace = true
serde.workspace = true
serde_json.workspace = true
thiserror.workspace = true
@@ -0,0 +1,38 @@
{
"binaryVersion": "0.152.1",
"protocol": "v2",
"appServerArgs": ["app-server", "--stdio"],
"generatedSchemaArgs": [
"app-server",
"generate-json-schema",
"--out",
"<DIR>",
"--experimental"
],
"schemaBundle": "codex_app_server_protocol.v2.schemas.json",
"schemaBundleSha256": "f9e3ca7e56300b4e5a5686419940ef77bbdc42846760d6a9cd53a21f20dc9ebd",
"clientMethods": [
"initialize",
"thread/start",
"turn/start",
"turn/interrupt"
],
"serverRequestMethods": [
"item/commandExecution/requestApproval",
"item/fileChange/requestApproval",
"item/tool/requestUserInput",
"mcpServer/elicitation/request",
"item/permissions/requestApproval",
"item/tool/call",
"applyPatchApproval",
"execCommandApproval"
],
"notificationMethods": [
"thread/started",
"turn/started",
"turn/completed",
"item/agentMessage/delta",
"item/started",
"item/completed"
]
}
File diff suppressed because it is too large Load Diff
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,135 @@
# Codex 0.152.1 协议适配审计
> 本文件按日期追加;测试数字和配置边界以文末“当前配置边界复核”段落为准,前文旧
> 数字保留为历史审计快照。
## 审计范围
本记录针对本机 `codex-cli 0.152.1` 的 App Server stdio 接线,范围只覆盖
`agent-codex`。审计没有使用 API key、登录凭据或网络请求;自动化测试全部使用
内存 JSON-RPC fixture。`0.152.1` 是本适配器明确固定的版本,不代表其它 Codex
版本兼容。
## 本机事实
通过本机 CLI 的帮助和 schema 生成入口核对到:
```text
codex --version
codex-cli 0.152.1
codex app-server --help
--stdio
generate-json-schema --out <DIR> [--experimental]
```
使用 `codex app-server generate-json-schema --out <DIR> --experimental` 得到的
schema 统计如下:
| 文件 | 数量 | SHA-256 |
| --- | ---: | --- |
| `ClientRequest.json` | 154 methods | `7443008decd3f978288accbc22da15e18ca20df4aca179c4750f71fdc0d91587` |
| `ServerRequest.json` | 11 methods | `38bb1c9dbb1dda2a688a7c8712b04319fe6ee2ed28b54acd2f5fd341d25567ff` |
| `ServerNotification.json` | 81 methods | `9adaa7f1d3838cf8328026294cee3194297c4f9ebc53bd23a9290fafc16d33c1` |
| `ClientNotification.json` | 1 method | `706cf248d75027c84a3c63348d0ed507182e8eba40069dd17541793de029145a` |
| `codex_app_server_protocol.v2.schemas.json` | 734 definitions | `f9e3ca7e56300b4e5a5686419940ef77bbdc42846760d6a9cd53a21f20dc9ebd` |
仓库只提交 `fixtures/codex-0.152.1/protocol-manifest.json`,它是上述生成结果的
来源/hash 和本适配器消费集合的 provenance 清单,不冒充完整 generated Rust schema。
清单中的 8 个 server-request method 和 6 个 notification method 是本模块当前有
字段级类型的子集;schema 中另外 3 个 server request(认证 token refresh、attestation、
current time)及其它通知会落到 `Unknown`/中立 JSON 路径。
## 已接线 API
- `ProcessConfig01521::try_new` 要求调用方先核对 `--version` 输出,只接受精确
`0.152.1`(也接受常见的 `v0.152.1` 前缀),并固定无 shell 的
`app-server --stdio` argv。
- `Client01521<R, W>` 可复用任意已连接的 reader/writer,提供严格版本子集的
`initialize`、`thread/start`、`turn/start`、`turn/interrupt` 和显式通知轮询。
- `AppServer01521` 把同一套 DTO 接到已有有界 `CodexAppServerProcess`;进程的
deadline、取消、EOF/异常退出回收和 reader/writer join 仍由通用 process facade
负责。
- `ServerRequest01521` 解码 command/file approval、user input、MCP elicitation、
permissions、dynamic tool、apply-patch 和 exec-command 请求;
`ServerRequestHandler01521` 通过 `HandlerBridge` 回写对应 JSON-RPC result。
handler 由宿主提供,适配器不会自动批准或执行工具;未知 method 必须由宿主显式
处理(或返回 `Raw`)。
- 生成 schema 中的 dynamic tool 有独立的 `namespace` 规格:调用里的 `tool` 是
namespace 内的子工具名,不是可忽略的备注。Host 两条 bridge 现在支持注入
`NamespaceToolResolver`;内置 `StaticNamespaceToolResolver` 只接受调用方显式注册的
`(namespace, tool) -> registered_tool` 映射,命中后仍会重新做工具定义、JSON Schema
和 approval 校验。缺省或 JSON `null` 的 namespace 继续走全局工具名;未注入 resolver、
未知 namespace、未知 namespace 内工具和空 namespace 均在审批/执行前 fail-closed,绝不
猜测分隔符或把 namespace 静默拼到全局工具名上。
- `Notification01521` 提供 agent message delta、thread/turn 生命周期和 item
started/completed 的字段级解析;未知通知保留为中立 `CodexAppServerNotification`。
复杂厂商对象(sandbox、permission profile、dynamic tool arguments、thread item)
在这个窄适配器中保持 `serde_json::Value`,避免把完整 734-definition schema 复制进
通用内核;版本升级时应重新生成 schema、更新清单和对应类型/测试。
## 离线回归
在 `/data/dsk/Genarrative-master` 执行:
```text
cargo fmt --manifest-path rust/Cargo.toml --all -- --check
TMPDIR=/var/tmp cargo test --locked --manifest-path rust/Cargo.toml \
-p agent-codex --all-features --no-fail-fast
RUSTFLAGS='-D warnings' TMPDIR=/var/tmp cargo clippy --locked \
--manifest-path rust/Cargo.toml -p agent-codex --all-targets --all-features -- -D warnings
```
结果:`agent-codex` 64 tests passed(其中 7 个为本版本 adapter 的 manifest、版本
校验、生命周期、审批/工具 handler、response shape、notification 和 UserInput
回归);Clippy 和 rustfmt 通过。测试没有启动真实 Codex,也没有访问网络。
## 明确不做项
- 不声称支持任意 Codex CLI/App Server 版本;不会仅因为都标记为 `v2` 就跳过版本
清单核对。
- 不把完整 generated schema、Codex 登录/认证、网络调用或长期会话持久化塞进
`agent-codex` 的通用内核;外部会话与恢复真相仍由 Host/Runtime 持有。
- 不自动批准 server request,不自动重连或重放已经发出的请求;进程启动后出现
协议/退出/超时错误仍按未知外部副作用交给上层 reconciliation。
## 2026-09-03 继续执行补充
- `NodeRuntimeEventMapper` 现在提供中立 `NodeRequest`、白名单 `NodeEvent` 和
`NodeResult` 到 Core `RuntimeEvent` 的连续 revision 映射。显式调用
`CodexAppServerBackend::invoke_node_with_runtime_events` 时,request、事件和结果会按
顺序交给调用方 sink;sink 自己负责 reducer/RuntimeStore 提交,适配器不会隐式修改
Host 会话或创建第二套状态机。
- Host 的外部会话 cancel 在当前进程没有 active index 时,会先用持久化 request-id 别名
找到原记录;`cancel_persisted`/`AgentHost::cancel_external_request` 可在重开 Host 后
显式发出 cancel,并将记录收束为 `cancel_requested`、`cancelled` 或保守的 `unknown`。
该路径不重新 invoke,也不把未知副作用自动标记为 safe。
- 本地新增 bridge/reopen/带外中断回归后,`agent-codex` 为 64 个测试,Host 为 54 个测试;两种
feature 集合的定向测试和 `-D warnings` Clippy 均通过。证据仍限于本地 fixture;真实
Codex session、完整 generated wire、自动对账/订阅和远端 CI 未验收。
- `CodexAppServerBackend::with_interrupt_hook` 是显式 opt-in 的带外控制入口:调用方必须
提供不重入同一 channel 的独立 control transport;未配置 hook 时仍走 channel 自带的
interrupt,不能把同步 channel 的调用误写成可抢占阻塞 I/O。
## 2026-09-03 当前配置边界复核
- `CodexCliConfig` 与 `CodexAppServerProcessConfig` 的 timeout、output/frame limit 在
serde 和运行时两层拒绝零值;`with_timeout` 拒绝小于一毫秒的 `Duration`,并使用检查式
`u64` 毫秒转换拒绝溢出。`timeout()` 不再用 `max(1)` 静默修正无效配置。
- 新增回归覆盖 camelCase 与 snake_case 配置、公开 struct 直写零值、`u64`/`usize` 解析
溢出、子毫秒以及 `Duration` 转换溢出。当前 `agent-codex` 两套 workspace feature
组合均为 77 个测试通过,Clippy `-D warnings` 与 rustfmt 通过。
- 这仍是本地配置和受控 fixture 的证据;真实 Codex 发行版 session、完整 generated
wire、远端 CI 和自动外部对账不在本地验收范围。
## 2026-09-05 typed 请求边界复核(当前)
- `InitializeParams01521`/`ClientInfo01521` 现在在 `Client01521` 和 `AppServer01521`
的 `initialize` 序列化前复验 `clientInfo.name/version`。公开字段或 serde 构造出的
非法参数会在触碰 transport/handler 前返回 `InvalidConfig`;`thread/start` 字段全为
可选,保持原有空值语义。
- 新增 `typed_client_revalidates_initialize_before_transport` 回归;
`agent-codex` 当前 81 个测试在 all/no-default workspace 均通过,Clippy `-D warnings`
与 rustfmt 通过。该证据仍是本地 fixture,不代表完整 generated wire 或真实发行版
session 已完成。
+20
View File
@@ -0,0 +1,20 @@
[package]
name = "agent-host"
version.workspace = true
edition.workspace = true
rust-version.workspace = true
license.workspace = true
description = "通用 Agent 运行时的组件装配与持久化 Host"
[dependencies]
agent-runtime-core.workspace = true
agent-runtime-engine.workspace = true
agent-runtime-sqlite.workspace = true
agent-codex.workspace = true
agent-provider-openai.workspace = true
agent-provider-fake.workspace = true
agent-mcp.workspace = true
agent-skills.workspace = true
serde.workspace = true
serde_json.workspace = true
thiserror.workspace = true
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,400 @@
use std::collections::BTreeSet;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use agent_host::{AgentHost, HostError};
use agent_provider_fake::{FakeProvider, FakeStep, FakeToolCall};
use agent_runtime_core::{
ApprovalDecision, ApprovalError, ApprovalPolicy, ContentPart, Message, ProviderErrorKind,
RuntimeSnapshot, ToolCall, ToolResult, reduce,
};
use agent_runtime_engine::{CompressionRequest, ContextCompressor, EngineError};
use serde_json::json;
fn assert_snapshot_matches_output(host: &AgentHost, result: &agent_host::HostRunOutput) {
// 结果消息、持久化快照和事件重放必须描述同一条消息历史。
let snapshot = host
.load_runtime_snapshot(&result.runtime_id)
.expect("读取 RuntimeSnapshot")
.expect("RuntimeSnapshot 应存在");
let run = snapshot.run(&result.run_id).expect("run 应存在");
assert_eq!(run.messages(), result.output.messages.as_slice());
// 消息正确还不够,派生工具索引也必须与当前上下文逐项一致。
let mut calls = Vec::new();
let mut results = Vec::new();
for part in run.messages().iter().flat_map(|message| message.content()) {
match part {
ContentPart::ToolCall {
id,
name,
arguments,
} => {
calls.push(ToolCall::try_new(id, name, arguments.clone()).unwrap());
}
ContentPart::ToolResult {
tool_call_id,
output,
is_error,
} => {
results.push(ToolResult::try_new(tool_call_id, output.clone(), *is_error).unwrap());
}
_ => {}
}
}
assert_eq!(run.tool_calls(), calls);
assert_eq!(run.tool_results(), results);
let events = host
.list_runtime_events(&result.runtime_id)
.expect("读取 runtime events");
let mut replayed = RuntimeSnapshot::try_new(&result.runtime_id).expect("创建空快照");
for event in events {
replayed = reduce(&replayed, &event).expect("runtime event 应可重放");
}
assert_eq!(replayed, snapshot);
}
fn assert_tool_rows_are_terminal_and_unique(host: &AgentHost, run_id: &str, expected: usize) {
// 每个工具调用只能有一行最终结果,不能把中间 requested 状态当成完成。
let rows = host.list_tool_calls(run_id).expect("读取工具调用记录");
assert_eq!(rows.len(), expected);
let mut ids = BTreeSet::new();
for row in rows {
assert!(ids.insert(row.id), "工具调用 id 不应重复");
assert_ne!(row.status, "requested", "工具调用不应停留在 requested");
assert!(row.result.is_some(), "已完成工具调用应有结果");
}
}
#[test]
fn normal_complete_persists_exactly_one_message_history_and_replay() {
let host = AgentHost::in_memory()
.expect("创建 Host")
.with_provider(Arc::new(FakeProvider::text("普通完成")), "fake");
let result = host
.run_with_messages("普通完成", vec![Message::user("普通完成").unwrap()])
.expect("普通运行应完成");
assert_eq!(result.output.text, "普通完成");
assert_snapshot_matches_output(&host, &result);
assert_tool_rows_are_terminal_and_unique(&host, &result.run_id, 0);
}
#[test]
fn normal_stream_persists_exactly_one_message_history_and_replay() {
let host = AgentHost::in_memory().expect("创建 Host").with_provider(
Arc::new(FakeProvider::new([FakeStep::stream_text(["流", "式"])])),
"fake",
);
let result = host
.run_with_messages_streaming("流式完成", vec![Message::user("流式完成").unwrap()])
.expect("流式运行应完成");
assert_eq!(result.output.text, "流式");
assert!(!result.output.stream_events.is_empty());
assert_snapshot_matches_output(&host, &result);
assert_tool_rows_are_terminal_and_unique(&host, &result.run_id, 0);
}
#[test]
fn automatic_tool_batches_and_later_round_match_engine_in_both_modes() {
for streaming in [false, true] {
let provider = Arc::new(FakeProvider::new([
FakeStep::ToolCalls {
text: "同一模型响应中的说明".to_owned(),
calls: (1..=3)
.map(|index| FakeToolCall {
id: format!("automatic-{index}"),
name: "echo".to_owned(),
arguments: json!({"index": index}),
})
.collect(),
},
FakeStep::tool_call("automatic-next", "echo", json!({"index": 4})),
FakeStep::text("同一模型响应中的说明"),
]));
let host = AgentHost::in_memory()
.unwrap()
.with_provider(provider.clone(), "fake");
let result = if streaming {
host.run_with_messages_streaming("自动批次", vec![Message::user("自动批次").unwrap()])
} else {
host.run("自动批次")
}
.expect("自动放行的连续工具轮次应完成");
assert_eq!(result.output.messages.len(), 8);
assert_eq!(provider.requests().snapshot().len(), 3);
assert_snapshot_matches_output(&host, &result);
assert_tool_rows_are_terminal_and_unique(&host, &result.run_id, 4);
}
}
struct AllowThenAsk {
asked: AtomicBool,
}
impl AllowThenAsk {
fn new() -> Self {
Self {
asked: AtomicBool::new(false),
}
}
}
impl ApprovalPolicy for AllowThenAsk {
fn decide(
&self,
request: &agent_runtime_core::ApprovalRequest,
) -> Result<ApprovalDecision, ApprovalError> {
if request.call().id() == "batch-call-3" && !self.asked.swap(true, Ordering::AcqRel) {
Ok(ApprovalDecision::Ask)
} else {
Ok(ApprovalDecision::Allow)
}
}
}
#[test]
fn three_tool_batch_and_next_tool_round_keep_one_message_history() {
let provider = Arc::new(FakeProvider::new([
FakeStep::tool_calls([
FakeToolCall {
id: "batch-call-1".to_owned(),
name: "echo".to_owned(),
arguments: json!({"index": 1}),
},
FakeToolCall {
id: "batch-call-2".to_owned(),
name: "echo".to_owned(),
arguments: json!({"index": 2}),
},
FakeToolCall {
id: "batch-call-3".to_owned(),
name: "echo".to_owned(),
arguments: json!({"index": 3}),
},
]),
FakeStep::tool_call("next-round-call", "echo", json!({"index": 4})),
FakeStep::text("三工具批次完成"),
]));
let host = AgentHost::in_memory()
.expect("创建 Host")
.with_provider(provider.clone(), "fake")
.with_approval(Arc::new(AllowThenAsk::new()));
let handle = host.prepare_run("三工具批次").expect("创建 run");
let first = host
.run_existing(&handle.run_id)
.expect_err("第三个调用应 Ask");
assert!(matches!(
first,
HostError::Engine(EngineError::ApprovalRequired { ref call_id, .. })
if call_id == "batch-call-3"
));
let checkpoint = host.read_checkpoint(&handle.run_id).unwrap().unwrap();
let checkpoint_messages: Vec<Message> = serde_json::from_value(checkpoint.messages).unwrap();
let pending_snapshot = host
.load_runtime_snapshot(&handle.runtime_id)
.unwrap()
.unwrap();
assert_eq!(
pending_snapshot.run(&handle.run_id).unwrap().messages(),
checkpoint_messages
);
assert_eq!(checkpoint_messages.len(), 4);
let approval = host
.list_approvals(&handle.run_id)
.expect("读取 approval")
.into_iter()
.find(|item| item.tool_call_id.as_deref() == Some("batch-call-3"))
.expect("第三个调用应有 pending approval");
host.resolve_approval(&approval.id, ApprovalDecision::Allow)
.expect("允许第三个调用");
host.resume_approval(&approval.id).expect("重新排队");
// 恢复时不能重新请求产生首批工具调用的 Provider;只继续未完成的调用和下一轮。
let result = host.run_existing(&handle.run_id).expect("恢复后应完成");
assert_eq!(result.output.text, "三工具批次完成");
assert_snapshot_matches_output(&host, &result);
assert_tool_rows_are_terminal_and_unique(&host, &result.run_id, 4);
assert_eq!(provider.remaining_steps(), 0);
}
#[test]
fn provider_error_keeps_reconciling_snapshot_replay_consistent() {
let host = AgentHost::in_memory().expect("创建 Host").with_provider(
Arc::new(FakeProvider::new([
FakeStep::tool_call("before-error", "echo", json!({"index": 1})),
FakeStep::Error {
kind: ProviderErrorKind::Stream,
message: "provider fixture error".to_owned(),
},
])),
"fake",
);
let handle = host.prepare_run("Provider 错误").expect("创建 run");
let error = host
.run_existing(&handle.run_id)
.expect_err("Provider 错误应返回");
assert!(matches!(error, HostError::Engine(EngineError::Provider(_))));
let record = host
.get_run(&handle.run_id)
.expect("读取 run")
.expect("run 应存在");
assert_eq!(record.status, "reconciling");
let checkpoint = host
.read_checkpoint(&handle.run_id)
.expect("读取 checkpoint")
.expect("Provider 错误应保留 checkpoint");
assert_eq!(checkpoint.phase, "provider_in_flight");
let snapshot = host
.load_runtime_snapshot(&handle.runtime_id)
.expect("读取 RuntimeSnapshot")
.expect("RuntimeSnapshot 应存在");
let events = host
.list_runtime_events(&handle.runtime_id)
.expect("读取 runtime events");
let mut replayed = RuntimeSnapshot::try_new(&handle.runtime_id).expect("创建空快照");
for event in events {
replayed = reduce(&replayed, &event).expect("runtime event 应可重放");
}
assert_eq!(replayed, snapshot);
let checkpoint_messages: Vec<Message> = serde_json::from_value(checkpoint.messages).unwrap();
assert_eq!(
snapshot.run(&handle.run_id).unwrap().messages(),
checkpoint_messages
);
assert_eq!(checkpoint_messages.len(), 3);
assert_tool_rows_are_terminal_and_unique(&host, &handle.run_id, 1);
}
struct ShortCompressor;
impl ContextCompressor for ShortCompressor {
fn compress(&self, _request: &CompressionRequest) -> Result<Vec<Message>, EngineError> {
Ok(vec![Message::user("压缩后的历史").expect("构造摘要消息")])
}
}
#[test]
fn tool_after_context_compression_does_not_restore_old_history() {
let host = AgentHost::in_memory()
.expect("创建 Host")
.with_provider(
Arc::new(FakeProvider::new([
FakeStep::ToolCalls {
text: "旧历史".repeat(10_000),
calls: vec![FakeToolCall {
id: "before-compression".to_owned(),
name: "echo".to_owned(),
arguments: json!({"index": 0}),
}],
},
FakeStep::ToolCalls {
text: "第二轮旧历史".repeat(10_000),
calls: vec![FakeToolCall {
id: "compressed-call".to_owned(),
name: "echo".to_owned(),
arguments: json!({"ok": true}),
}],
},
FakeStep::text("压缩后完成"),
])),
"fake",
)
.with_context_compressor(Arc::new(ShortCompressor));
let result = host
.run_with_messages("压缩测试", vec![Message::user("开始工具后压缩").unwrap()])
.expect("压缩后运行应完成");
assert_eq!(result.output.text, "压缩后完成");
assert_eq!(
result
.output
.context_observations
.iter()
.filter(|observation| observation.compression_attempted)
.count(),
2
);
assert_snapshot_matches_output(&host, &result);
assert_tool_rows_are_terminal_and_unique(&host, &result.run_id, 2);
let runtime = host
.load_runtime_snapshot(&result.runtime_id)
.expect("读取 RuntimeSnapshot")
.expect("RuntimeSnapshot 应存在");
let run = runtime.run(&result.run_id).expect("run");
let messages = run.messages();
assert!(messages.iter().any(|message| {
message
.content()
.iter()
.any(|part| part.as_text() == Some("压缩后的历史"))
}));
assert!(!messages.iter().any(|message| {
message
.content()
.iter()
.any(|part| part.as_text().is_some_and(|text| text.contains("旧历史")))
}));
}
struct FailingCompressor;
impl ContextCompressor for FailingCompressor {
fn compress(&self, _request: &CompressionRequest) -> Result<Vec<Message>, EngineError> {
Err(EngineError::ContextOverflow("压缩失败 fixture".into()))
}
}
#[test]
fn failed_compression_after_tool_keeps_full_history_without_repeating_results() {
let provider = Arc::new(FakeProvider::new([
FakeStep::ToolCalls {
text: "待压缩".repeat(10_000),
calls: vec![FakeToolCall {
id: "before-failed-compression".into(),
name: "echo".into(),
arguments: json!({"ok": true}),
}],
},
FakeStep::text("不能调用这一步"),
]));
let host = AgentHost::in_memory()
.unwrap()
.with_provider(provider.clone(), "fake")
.with_context_compressor(Arc::new(FailingCompressor));
let handle = host.prepare_run("工具后压缩失败").unwrap();
assert!(matches!(
host.run_existing(&handle.run_id),
Err(HostError::Engine(EngineError::ContextOverflow(_)))
));
assert_eq!(provider.requests().snapshot().len(), 1);
assert_eq!(
host.get_run(&handle.run_id).unwrap().unwrap().status,
"reconciling"
);
let checkpoint = host.read_checkpoint(&handle.run_id).unwrap().unwrap();
assert_eq!(checkpoint.phase, "compacting");
let messages: Vec<Message> = serde_json::from_value(checkpoint.messages).unwrap();
assert_eq!(messages.len(), 3);
let snapshot = host
.load_runtime_snapshot(&handle.runtime_id)
.unwrap()
.unwrap();
let run = snapshot.run(&handle.run_id).unwrap();
assert_eq!(run.messages(), messages);
assert_eq!(run.tool_calls().len(), 1);
assert_eq!(run.tool_results().len(), 1);
assert_tool_rows_are_terminal_and_unique(&host, &handle.run_id, 1);
let replayed = host
.list_runtime_events(&handle.runtime_id)
.unwrap()
.iter()
.fold(
RuntimeSnapshot::try_new(&handle.runtime_id).unwrap(),
|snapshot, event| reduce(&snapshot, event).unwrap(),
);
assert_eq!(replayed, snapshot);
}
+13
View File
@@ -0,0 +1,13 @@
[package]
name = "agent-mcp"
version = "0.1.0"
edition = "2024"
rust-version.workspace = true
description = "通用 Agent 的轻量 MCP 协议与传输适配接口"
license = "MIT"
[dependencies]
reqwest = { version = "0.12", features = ["blocking"] }
serde = { version = "1", features = ["derive"] }
serde_json = "1"
tokio = { version = "1", features = ["macros", "net", "rt", "time"] }
+32
View File
@@ -0,0 +1,32 @@
#!/bin/sh
# 一个只用于测试的最小 MCP stdio server。
#
# 它不执行 shell 参数,也不读取环境中的密钥;只按 JSON-RPC method 返回
# 固定的 initialize、tools/list 和 tools/call fixture。测试通过 `sh` 调用
# 本文件,因此不依赖可执行位或额外的 jq/node 运行时。
while IFS= read -r line; do
# 请求 id 是数字时原样回显,保证 client 的 response-id 校验仍然有效。
id=$(printf '%s\n' "$line" | sed -n 's/.*"id"[[:space:]]*:[[:space:]]*\([0-9][0-9]*\).*/\1/p')
[ -n "$id" ] || id=null
case "$line" in
*'"method":"notifications/initialized"'*|*'"method": "notifications/initialized"'*)
# notification 没有 response;继续读取后续请求。
continue
;;
*'"method":"initialize"'*|*'"method": "initialize"'*)
printf '%s\n' "{\"jsonrpc\":\"2.0\",\"id\":$id,\"result\":{\"protocolVersion\":\"2025-06-18\",\"capabilities\":{\"tools\":{}}}}"
;;
*'"method":"tools/list"'*|*'"method": "tools/list"'*)
printf '%s\n' "{\"jsonrpc\":\"2.0\",\"id\":$id,\"result\":{\"tools\":[{\"name\":\"fixture_echo\",\"description\":\"fixture tool\",\"inputSchema\":{\"type\":\"object\",\"properties\":{\"text\":{\"type\":\"string\"}},\"required\":[\"text\"]}}]}}"
;;
*'"method":"tools/call"'*|*'"method": "tools/call"'*)
printf '%s\n' "{\"jsonrpc\":\"2.0\",\"id\":$id,\"result\":{\"content\":[{\"type\":\"text\",\"text\":\"fixture-ok\"}],\"isError\":false}}"
;;
*)
printf '%s\n' "{\"jsonrpc\":\"2.0\",\"id\":$id,\"error\":{\"code\":-32601,\"message\":\"fixture method not found\"}}"
;;
esac
done
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,12 @@
[package]
name = "agent-provider-fake"
version.workspace = true
edition.workspace = true
rust-version.workspace = true
license.workspace = true
description = "用于 Agent 回归测试的确定性 Provider"
[dependencies]
agent-runtime-core.workspace = true
serde.workspace = true
serde_json.workspace = true
+337
View File
@@ -0,0 +1,337 @@
//! 可脚本化的离线 Provider。
//!
//! Fake Provider 只依赖 Core 契约,不读取环境变量、不访问网络,也不把测试
//! 场景塞进 Engine。它既能驱动纯文本/工具循环,也能稳定复现流中断、取消和
//! 压缩响应等边界,供所有适配器和 Host 测试复用。
use std::collections::VecDeque;
use std::sync::{Arc, Mutex};
use agent_runtime_core::{
ContentPart, ModelProvider, ProviderError, ProviderErrorKind, ProviderRequest,
ProviderResponse, ProviderStreamEvent, ProviderStreamSink, ToolCall,
};
use serde::{Deserialize, Serialize};
use serde_json::Value;
/// Fake Provider 的一次确定性响应。
#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
#[serde(tag = "type", rename_all = "kebab-case")]
pub enum FakeStep {
/// 返回一条普通文本响应。
Text { text: String },
/// 返回一个结构化工具调用。
ToolCall {
id: String,
name: String,
arguments: Value,
},
/// 返回文本和一批工具调用,便于覆盖串行审批/执行。
ToolCalls {
#[serde(default)]
text: String,
calls: Vec<FakeToolCall>,
},
/// 把文本按给定片段通过 `ModelProvider::stream` 发出。
StreamText { chunks: Vec<String> },
/// 模拟 Provider 在流中途失败。
StreamError { message: String },
/// 模拟调用方取消或外部调用结果未知。
Cancelled,
/// 作为压缩 Provider 测试脚本使用的摘要响应。
Compression { summary: String },
/// 模拟普通瞬时/上游错误。
Error {
kind: ProviderErrorKind,
message: String,
},
}
#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
pub struct FakeToolCall {
pub id: String,
pub name: String,
pub arguments: Value,
}
impl FakeStep {
pub fn text(text: impl Into<String>) -> Self {
Self::Text { text: text.into() }
}
pub fn tool_call(id: impl Into<String>, name: impl Into<String>, arguments: Value) -> Self {
Self::ToolCall {
id: id.into(),
name: name.into(),
arguments,
}
}
pub fn tool_calls(calls: impl IntoIterator<Item = FakeToolCall>) -> Self {
Self::ToolCalls {
text: String::new(),
calls: calls.into_iter().collect(),
}
}
pub fn stream_text(chunks: impl IntoIterator<Item = impl Into<String>>) -> Self {
Self::StreamText {
chunks: chunks.into_iter().map(Into::into).collect(),
}
}
pub fn compression(summary: impl Into<String>) -> Self {
Self::Compression {
summary: summary.into(),
}
}
}
/// 记录 Provider 实际收到的请求,便于测试消息边界和恢复游标。
#[derive(Clone, Debug, Default)]
pub struct FakeRequestLog(Arc<Mutex<Vec<ProviderRequest>>>);
impl FakeRequestLog {
pub fn snapshot(&self) -> Vec<ProviderRequest> {
self.0
.lock()
.map(|requests| requests.clone())
.unwrap_or_default()
}
}
/// 线程安全的确定性 Provider。每次调用消费脚本中的下一步。
#[derive(Clone, Debug)]
pub struct FakeProvider {
steps: Arc<Mutex<VecDeque<FakeStep>>>,
requests: FakeRequestLog,
}
impl FakeProvider {
pub fn new(steps: impl IntoIterator<Item = FakeStep>) -> Self {
Self {
steps: Arc::new(Mutex::new(steps.into_iter().collect())),
requests: FakeRequestLog::default(),
}
}
pub fn text(text: impl Into<String>) -> Self {
Self::new([FakeStep::text(text)])
}
pub fn tool_then_text(
id: impl Into<String>,
name: impl Into<String>,
arguments: Value,
text: impl Into<String>,
) -> Self {
Self::new([
FakeStep::tool_call(id, name, arguments),
FakeStep::text(text),
])
}
pub fn requests(&self) -> FakeRequestLog {
self.requests.clone()
}
pub fn remaining_steps(&self) -> usize {
self.steps.lock().map(|steps| steps.len()).unwrap_or(0)
}
fn next_step(&self, request: &ProviderRequest) -> Result<FakeStep, ProviderError> {
if let Ok(mut requests) = self.requests.0.lock() {
requests.push(request.clone());
}
self.steps
.lock()
.map_err(|_| ProviderError::new(ProviderErrorKind::Unavailable, "fake 脚本锁已损坏"))?
.pop_front()
.ok_or_else(|| ProviderError::new(ProviderErrorKind::Unavailable, "fake 脚本已耗尽"))
}
fn response_for(
request: &ProviderRequest,
step: FakeStep,
) -> Result<ProviderResponse, ProviderError> {
match step {
FakeStep::Text { text } | FakeStep::Compression { summary: text } => {
ProviderResponse::text(request.request_id(), request.model(), text)
.map_err(Into::into)
}
FakeStep::ToolCall {
id,
name,
arguments,
} => {
let call = ToolCall::try_new(id, name, arguments)?;
ProviderResponse::try_new(request.request_id(), request.model(), [], [call])
.map_err(Into::into)
}
FakeStep::ToolCalls { text, calls } => {
let calls = calls
.into_iter()
.map(|call| ToolCall::try_new(call.id, call.name, call.arguments))
.collect::<Result<Vec<_>, _>>()?;
let content = if text.is_empty() {
Vec::new()
} else {
vec![ContentPart::text(text)?]
};
ProviderResponse::try_new(request.request_id(), request.model(), content, calls)
.map_err(Into::into)
}
FakeStep::StreamText { chunks } => {
ProviderResponse::text(request.request_id(), request.model(), chunks.concat())
.map_err(Into::into)
}
FakeStep::StreamError { message } => {
Err(ProviderError::new(ProviderErrorKind::Stream, message))
}
FakeStep::Cancelled => Err(ProviderError::new(
ProviderErrorKind::Stream,
"fake provider cancelled",
)),
FakeStep::Error { kind, message } => Err(ProviderError::new(kind, message)),
}
}
}
impl ModelProvider for FakeProvider {
fn complete(&self, request: &ProviderRequest) -> Result<ProviderResponse, ProviderError> {
let step = self.next_step(request)?;
match step {
FakeStep::StreamError { message } => {
Err(ProviderError::new(ProviderErrorKind::Stream, message))
}
FakeStep::Cancelled => Err(ProviderError::new(
ProviderErrorKind::Stream,
"fake provider cancelled",
)),
other => Self::response_for(request, other),
}
}
fn stream(
&self,
request: &ProviderRequest,
sink: &mut dyn ProviderStreamSink,
) -> Result<ProviderResponse, ProviderError> {
let step = self.next_step(request)?;
match step {
FakeStep::StreamText { chunks } => {
let mut accumulated = String::new();
for chunk in chunks {
accumulated.push_str(&chunk);
sink.emit(ProviderStreamEvent::TextDelta {
delta: chunk,
accumulated: accumulated.clone(),
})?;
}
sink.emit(ProviderStreamEvent::Completed)?;
ProviderResponse::text(request.request_id(), request.model(), accumulated)
.map_err(Into::into)
}
FakeStep::StreamError { message } => {
Err(ProviderError::new(ProviderErrorKind::Stream, message))
}
FakeStep::Cancelled => Err(ProviderError::new(
ProviderErrorKind::Stream,
"fake provider cancelled",
)),
other => {
let response = Self::response_for(request, other)?;
let mut accumulated = String::new();
for part in response.content() {
if let Some(text) = part.as_text() {
accumulated.push_str(text);
sink.emit(ProviderStreamEvent::TextDelta {
delta: text.to_owned(),
accumulated: accumulated.clone(),
})?;
}
}
for call in response.tool_calls() {
sink.emit(ProviderStreamEvent::ToolCallDelta {
call_id: call.id().to_owned(),
name: Some(call.name().to_owned()),
arguments_delta: call.arguments().to_string(),
})?;
}
sink.emit(ProviderStreamEvent::Completed)?;
Ok(response)
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use agent_runtime_core::{ProviderStreamEvent, ToolChoice, ToolDefinition};
use serde_json::json;
struct Sink(Vec<ProviderStreamEvent>);
impl ProviderStreamSink for Sink {
fn emit(&mut self, event: ProviderStreamEvent) -> Result<(), ProviderError> {
self.0.push(event);
Ok(())
}
}
fn request(_provider: &FakeProvider) -> ProviderRequest {
ProviderRequest::try_new(
"request-1",
"fake",
[agent_runtime_core::Message::user("hi").unwrap()],
)
.unwrap()
.with_tools(
[ToolDefinition::try_new("echo", "echo", json!({"type":"object"})).unwrap()],
ToolChoice::Auto,
)
.unwrap()
}
#[test]
fn scripted_tool_then_text_is_deterministic() {
let provider = FakeProvider::tool_then_text("call-1", "echo", json!({"x": 1}), "done");
let first = provider.complete(&request(&provider)).unwrap();
assert_eq!(first.tool_calls()[0].id(), "call-1");
let second_request = ProviderRequest::try_new(
"request-2",
"fake",
[agent_runtime_core::Message::user("tool result").unwrap()],
)
.unwrap();
let second = provider.complete(&second_request).unwrap();
assert_eq!(second.content()[0].as_text(), Some("done"));
assert_eq!(provider.requests().snapshot().len(), 2);
}
#[test]
fn stream_text_emits_each_chunk_and_completion() {
let provider = FakeProvider::new([FakeStep::stream_text(["你", "好"])]);
let mut sink = Sink(Vec::new());
let response = provider.stream(&request(&provider), &mut sink).unwrap();
assert_eq!(response.content()[0].as_text(), Some("你好"));
assert!(matches!(sink.0[0], ProviderStreamEvent::TextDelta { .. }));
assert!(matches!(sink.0[1], ProviderStreamEvent::TextDelta { .. }));
assert!(matches!(sink.0[2], ProviderStreamEvent::Completed));
}
#[test]
fn error_and_cancel_scenarios_are_explicit() {
let provider = FakeProvider::new([FakeStep::Cancelled]);
let error = provider.complete(&request(&provider)).unwrap_err();
assert_eq!(error.kind(), ProviderErrorKind::Stream);
assert!(error.message().contains("cancelled"));
}
#[test]
fn compression_step_is_a_normal_scripted_response() {
let provider = FakeProvider::new([FakeStep::compression("summary")]);
let response = provider.complete(&request(&provider)).unwrap();
assert_eq!(response.content()[0].as_text(), Some("summary"));
}
}
@@ -0,0 +1,15 @@
[package]
name = "agent-provider-openai"
version.workspace = true
edition.workspace = true
rust-version.workspace = true
license.workspace = true
description = "OpenAI Responses API 的通用 Agent Provider 适配器"
[dependencies]
agent-runtime-core.workspace = true
# 使用系统 TLS,避免把大型 rustls/ring 编译链带进最小 CLI。
reqwest = { version = "0.12", features = ["blocking", "json"] }
serde.workspace = true
serde_json.workspace = true
thiserror.workspace = true
File diff suppressed because it is too large Load Diff

Some files were not shown because too many files have changed in this diff Show More