修复第三方 Provider 流式工具计划请求
Project CI / Repository checks (pull_request) Successful in 3m38s
Project CI / Frontend tests (pull_request) Successful in 4m4s
Project CI / Backend tests (pull_request) Successful in 5m0s
Project CI / Native shell tests (pull_request) Failing after 10m53s

Provider 的 llm.stream=true 时改用 stream_run 并聚合完整响应

补充 Anthropic 原生流式工具计划与最终回复回归测试

扩展响应流测试夹具并同步第三方 Provider 兼容性说明
This commit is contained in:
2026-08-27 14:10:03 +08:00
parent 9476ff7648
commit a155529861
5 changed files with 337 additions and 50 deletions
@@ -1694,7 +1694,12 @@ pub(in crate::agent) async fn request_game_creator_agent_runtime_llm_with_persis
&config_path_for_request,
)
.map_err(platform_llm::LlmError::InvalidConfig)?;
client.run(request).await
request_game_creator_agent_runtime_provider_llm(
&client,
&llm_for_request,
request,
)
.await
}
_ => unreachable!("agent mode is normalized"),
}
@@ -1704,6 +1709,18 @@ pub(in crate::agent) async fn request_game_creator_agent_runtime_llm_with_persis
.await
}
async fn request_game_creator_agent_runtime_provider_llm(
client: &platform_llm::LlmClient,
llm: &GameCreatorLlmConfig,
request: platform_llm::LlmRunRequest,
) -> Result<platform_llm::LlmRunResponse, platform_llm::LlmError> {
if llm.stream {
client.stream_run(request, |_| {}).await
} else {
client.run(request).await
}
}
pub(in crate::agent) async fn request_game_creator_agent_runtime_llm_with_transient_retries(
root: &Path,
provider_snapshot: &AgentRuntimeProviderRequestSnapshot,
@@ -1748,9 +1765,8 @@ pub(in crate::agent) async fn request_game_creator_agent_runtime_llm_with_transi
request_game_creator_agent_codex_cli(request.clone()).await
}
GAME_CREATOR_AGENT_MODE_PROVIDER => {
client
.expect("provider mode constructs an HTTP client")
.run(request.clone())
let client = client.expect("provider mode constructs an HTTP client");
request_game_creator_agent_runtime_provider_llm(&client, llm, request.clone())
.await
}
_ => unreachable!("agent mode is normalized"),
@@ -1607,6 +1607,74 @@ pub(crate) fn final_tool_plan_response(response: impl Into<String>) -> String {
.to_string()
}
pub(crate) fn native_anthropic_tool_plan_response(
call_id: &str,
function_name: &str,
arguments: &str,
) -> String {
let events = vec![
serde_json::json!({
"type": "message_start",
"message": { "usage": { "input_tokens": 11, "output_tokens": 22 } }
}),
serde_json::json!({
"type": "content_block_start",
"index": 0,
"content_block": {
"type": "tool_use",
"id": call_id,
"name": function_name,
"input": {}
}
}),
serde_json::json!({
"type": "content_block_delta",
"index": 0,
"delta": { "type": "input_json_delta", "partial_json": arguments }
}),
serde_json::json!({ "type": "content_block_stop", "index": 0 }),
serde_json::json!({ "type": "message_delta", "delta": { "stop_reason": "tool_use" } }),
serde_json::json!({ "type": "message_stop" }),
];
events
.iter()
.map(|event| format!("data: {event}\n\n"))
.collect()
}
pub(crate) fn native_anthropic_text_stream_response(text: &str) -> String {
let events = vec![
serde_json::json!({
"type": "message_start",
"message": { "usage": { "input_tokens": 11, "output_tokens": 0 } }
}),
serde_json::json!({
"type": "content_block_start",
"index": 0,
"content_block": { "type": "text", "text": "" }
}),
serde_json::json!({
"type": "content_block_delta",
"index": 0,
"delta": { "type": "text_delta", "text": text }
}),
serde_json::json!({
"type": "content_block_stop",
"index": 0
}),
serde_json::json!({
"type": "message_delta",
"delta": { "stop_reason": "end_turn" },
"usage": { "output_tokens": 22 }
}),
serde_json::json!({ "type": "message_stop" }),
];
events
.iter()
.map(|event| format!("data: {event}\n\n"))
.collect()
}
fn user_input_tool_plan_response(question: &str) -> String {
serde_json::json!({
"thinkingSummary": "实现路径取决于用户选择,需要先暂停并澄清",
@@ -2290,9 +2358,36 @@ fn spawn_mock_llm_tool_plan_then_transient_final_compaction(
fn spawn_mock_llm_raw_responses_with_capture(
response_bodies: Vec<serde_json::Value>,
request_sender: Option<mpsc::Sender<String>>,
) -> String {
spawn_mock_llm_raw_responses_with_content_type_with_capture(
response_bodies
.into_iter()
.map(|body| body.to_string())
.collect(),
request_sender,
"application/json",
)
}
fn spawn_mock_llm_stream_responses_with_capture(
response_bodies: Vec<String>,
request_sender: Option<mpsc::Sender<String>>,
) -> String {
spawn_mock_llm_raw_responses_with_content_type_with_capture(
response_bodies,
request_sender,
"text/event-stream; charset=utf-8",
)
}
fn spawn_mock_llm_raw_responses_with_content_type_with_capture(
response_bodies: Vec<String>,
request_sender: Option<mpsc::Sender<String>>,
content_type: &str,
) -> String {
let listener = bind_test_tcp_listener("mock raw llm bind");
let base_url = format!("http://{}", listener.local_addr().expect("mock llm addr"));
let content_type = content_type.to_string();
std::thread::spawn(move || {
for response_body in response_bodies {
let (mut stream, _) = listener.accept().expect("mock raw llm accept");
@@ -2300,11 +2395,11 @@ fn spawn_mock_llm_raw_responses_with_capture(
if let Some(sender) = request_sender.as_ref() {
let _ = sender.send(request_text);
}
let body = response_body.to_string();
let response = format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
body.len(),
body
"HTTP/1.1 200 OK\r\nContent-Type: {}\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
content_type,
response_body.len(),
response_body
);
stream
.write_all(response.as_bytes())
@@ -2468,6 +2563,11 @@ enum ResponseStreamMockFinalResponse {
Disconnect,
}
enum ResponseStreamMockPlanningResponse {
NonStream(String),
Stream(String),
}
struct ResponseStreamMockServer {
base_url: String,
first_delta_written: mpsc::Receiver<()>,
@@ -2487,7 +2587,7 @@ impl ResponseStreamMockServer {
fn spawn_response_stream_mock_llm_server(
api_kind: &str,
planning_response: String,
planning_response: ResponseStreamMockPlanningResponse,
final_response: Option<ResponseStreamMockFinalResponse>,
) -> ResponseStreamMockServer {
let listener = bind_test_tcp_listener("response stream mock bind");
@@ -2501,39 +2601,72 @@ fn spawn_response_stream_mock_llm_server(
let (stop, stop_receiver) = mpsc::channel();
let handle = std::thread::spawn(move || {
let mut requests = Vec::new();
let (mut planning_stream, _) = listener
.accept()
.expect("response stream planning request accept");
let planning_request = read_mock_http_request(&mut planning_stream);
requests.push(planning_request);
let planning_body = match api_kind.as_str() {
"openai_responses" => serde_json::json!({
"id": "resp_response_stream_planning",
"model": "response-stream-model",
"output_text": planning_response,
"status": "completed",
"usage": { "input_tokens": 11, "output_tokens": 22, "total_tokens": 33 }
}),
"openai_chat" => serde_json::json!({
"id": "chatcmpl_response_stream_planning",
"model": "response-stream-model",
"choices": [{
"message": { "content": planning_response },
"finish_reason": "stop"
}],
"usage": { "prompt_tokens": 11, "completion_tokens": 22, "total_tokens": 33 }
}),
other => panic!("unsupported response stream mock api kind: {other}"),
let planning_http_response = match planning_response {
ResponseStreamMockPlanningResponse::NonStream(planning_response) => {
let planning_body = match api_kind.as_str() {
"openai_responses" => serde_json::json!({
"id": "resp_response_stream_planning",
"model": "response-stream-model",
"output_text": planning_response,
"status": "completed",
"usage": { "input_tokens": 11, "output_tokens": 22, "total_tokens": 33 }
}),
"openai_chat" => serde_json::json!({
"id": "chatcmpl_response_stream_planning",
"model": "response-stream-model",
"choices": [{
"message": { "content": planning_response },
"finish_reason": "stop"
}],
"usage": { "prompt_tokens": 11, "completion_tokens": 22, "total_tokens": 33 }
}),
other => panic!("unsupported response stream mock api kind: {other}"),
}
.to_string();
format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
planning_body.len(),
planning_body
)
}
ResponseStreamMockPlanningResponse::Stream(planning_response) => {
let planning_body = match api_kind.as_str() {
"openai_responses" => format!(
"data: {}\n\ndata: {}\n\n",
serde_json::json!({
"type": "response.output_text.delta",
"delta": planning_response
}),
serde_json::json!({ "type": "response.completed" })
),
"openai_chat" => format!(
"data: {}\n\ndata: {}\n\ndata: [DONE]\n\n",
serde_json::json!({
"choices": [{ "delta": { "content": planning_response } }]
}),
serde_json::json!({
"choices": [{ "finish_reason": "stop" }]
})
),
other => panic!("unsupported response stream mock api kind: {other}"),
};
format!(
"HTTP/1.1 200 OK\r\nContent-Type: text/event-stream; charset=utf-8\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
planning_body.len(),
planning_body
)
}
};
{
let (mut stream, _) = listener
.accept()
.expect("response stream planning request accept");
let planning_request = read_mock_http_request(&mut stream);
requests.push(planning_request);
stream
.write_all(planning_http_response.as_bytes())
.expect("response stream planning response");
}
.to_string();
let planning_http_response = format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
planning_body.len(),
planning_body
);
planning_stream
.write_all(planning_http_response.as_bytes())
.expect("response stream planning response");
if let Some(final_response) = final_response {
listener
@@ -3978,7 +4111,7 @@ fn run_response_stream_distinct_final_reply_case(api_kind: &str, case_name: &str
let canonical_response = format!("{first_delta}{second_delta}");
let mock = spawn_response_stream_mock_llm_server(
api_kind,
final_tool_plan_response(&planning_fallback),
ResponseStreamMockPlanningResponse::Stream(final_tool_plan_response(&planning_fallback)),
Some(ResponseStreamMockFinalResponse::Deltas(
first_delta.clone(),
second_delta.clone(),
@@ -4080,7 +4213,7 @@ fn run_response_stream_distinct_final_reply_case(api_kind: &str, case_name: &str
.all(|request| request.contains(expected_route)));
let planning_request = mock_http_request_json(&requests[0]);
let final_request = mock_http_request_json(&requests[1]);
assert_eq!(planning_request["stream"], Value::Bool(false));
assert_eq!(planning_request["stream"], Value::Bool(true));
assert_eq!(final_request["stream"], Value::Bool(true));
assert!(requests[0].contains("respond_to_user"));
assert!(!requests[0].contains("\"name\":\"submit_agent_tool_plan\""));
@@ -2505,6 +2505,120 @@ async fn background_agent_runtime_executes_native_function_tool_plan() {
fs::remove_dir_all(root).ok();
}
#[tokio::test]
async fn background_agent_runtime_executes_streamed_native_function_tool_plan() {
let root = unique_project_path();
init_local_game_project_at(&root, "project-stream-native-tool", "流式原生工具项目")
.expect("project init");
let (sender, receiver) = mpsc::channel();
let arguments = serde_json::json!({
"reason": "读取项目索引",
"input": {}
})
.to_string();
let function_name = native_runtime_function_name("project.index").expect("index function");
let base_url = spawn_mock_llm_stream_responses_with_capture(
vec![
native_anthropic_tool_plan_response(
"call-stream-native-index",
&function_name,
&arguments,
),
native_anthropic_tool_plan_response(
"call-stream-native-final",
AGENT_RUNTIME_RESPOND_FUNCTION_NAME,
&serde_json::json!({
"response": "流式原生工具调用已完成聚合。STREAM_NATIVE_TOOL_OK"
})
.to_string(),
),
native_anthropic_text_stream_response(
"最终回复已通过独立流式收束请求生成。STREAM_NATIVE_FINAL_OK",
),
],
Some(sender),
);
let _config_guard = write_test_local_config(format!(
r#"{{
"agentMode": "provider",
"agentLlm": {{
"design-director": {{
"apiKey": "stream-native-tool-key",
"baseUrl": {base_url:?},
"model": "stream-native-tool-model",
"apiKind": "anthropic",
"stream": true,
"webSearchEnabled": false,
"maxRetries": 0
}}
}}
}}"#
));
let run_id = "design-stream-native-function-tool-run";
start_game_creator_agent_background_task_at(
&root,
"design-director",
"用流式原生工具读取项目索引",
run_id,
)
.expect("start streamed native tool task");
let request = receiver
.recv_timeout(Duration::from_secs(2))
.expect("streamed native tool request");
assert!(request.contains("POST /v1/messages HTTP/1.1"));
assert_eq!(
mock_http_request_json(&request)["stream"],
Value::Bool(true)
);
assert!(request.contains(&function_name));
let respond_request = receiver
.recv_timeout(Duration::from_secs(2))
.expect("streamed native respond_to_user request");
assert!(respond_request.contains(AGENT_RUNTIME_RESPOND_FUNCTION_NAME));
assert_eq!(
mock_http_request_json(&respond_request)["stream"],
Value::Bool(true)
);
let final_request = receiver
.recv_timeout(Duration::from_secs(2))
.expect("streamed native final-reply request");
assert!(!final_request.contains(AGENT_RUNTIME_RESPOND_FUNCTION_NAME));
assert_eq!(
mock_http_request_json(&final_request)["stream"],
Value::Bool(true)
);
assert!(receiver.recv_timeout(Duration::from_millis(200)).is_err());
let runtime = wait_for_agent_runtime_idle(&root, "design-director");
assert_eq!(runtime.status, "idle");
assert_eq!(runtime.phase, "completed");
assert_eq!(runtime.recent_tool_calls.len(), 1);
assert_eq!(runtime.recent_tool_calls[0].tool, "project.index");
assert_eq!(runtime.recent_tool_calls[0].status, "ok");
assert_eq!(
runtime.last_response.as_deref(),
Some("最终回复已通过独立流式收束请求生成。STREAM_NATIVE_FINAL_OK")
);
let protocol_records = read_agent_db_records_for_test(&root)
.into_iter()
.filter(|record| {
record["recordType"] == "agent.runtime.tool_plan.protocol" && record["runId"] == run_id
})
.collect::<Vec<_>>();
assert_eq!(protocol_records.len(), 2);
assert_eq!(protocol_records[0]["protocol"], "native_runtime_tools");
assert_eq!(protocol_records[0]["functionCallCount"], 1);
assert_eq!(protocol_records[0]["functionNames"][0], function_name);
assert_eq!(
protocol_records[1]["functionNames"][0],
AGENT_RUNTIME_RESPOND_FUNCTION_NAME
);
fs::remove_dir_all(root).ok();
}
#[tokio::test]
async fn manual_context_compaction_is_private_and_hydrates_runtime_usage() {
let root = unique_project_path();
@@ -335,7 +335,7 @@ async fn response_stream_disabled_keeps_direct_planning_reply_to_one_request() {
let direct_response = "非流配置直接采用 planning response,且只发起一次请求。";
let mock = spawn_response_stream_mock_llm_server(
"openai_responses",
final_tool_plan_response(direct_response),
ResponseStreamMockPlanningResponse::NonStream(final_tool_plan_response(direct_response)),
None,
);
let base_url = mock.base_url.clone();
@@ -426,7 +426,7 @@ async fn response_stream_private_process_output_is_never_published_or_committed_
let raw_provider_response = format!("{first_delta}{second_delta}");
let mock = spawn_response_stream_mock_llm_server(
"openai_responses",
final_tool_plan_response(&planning_fallback),
ResponseStreamMockPlanningResponse::Stream(final_tool_plan_response(&planning_fallback)),
Some(ResponseStreamMockFinalResponse::Deltas(
first_delta.clone(),
second_delta.clone(),
@@ -593,7 +593,7 @@ async fn response_stream_private_process_output_is_never_published_or_committed_
assert_eq!(requests.len(), 2);
assert_eq!(
mock_http_request_json(&requests[0])["stream"],
Value::Bool(false)
Value::Bool(true)
);
assert_eq!(
mock_http_request_json(&requests[1])["stream"],
@@ -629,7 +629,7 @@ async fn response_stream_final_disconnect_with_retry_disabled_fails_without_comm
let planning_fallback = "final stream 失败后只提交这条 planning fallback。";
let mock = spawn_response_stream_mock_llm_server(
"openai_responses",
final_tool_plan_response(planning_fallback),
ResponseStreamMockPlanningResponse::Stream(final_tool_plan_response(planning_fallback)),
Some(ResponseStreamMockFinalResponse::Disconnect),
);
let base_url = mock.base_url.clone();
@@ -718,7 +718,7 @@ async fn response_stream_final_disconnect_with_retry_disabled_fails_without_comm
);
assert_eq!(
mock_http_request_json(&requests[0])["stream"],
Value::Bool(false)
Value::Bool(true)
);
assert_eq!(
mock_http_request_json(&requests[1])["stream"],
@@ -1,8 +1,8 @@
# 【技术说明】AGC 接第三方 Provider 的兼容性缺陷
- 首次记录:2026-08-19
- 最新核对:2026-08-25,当前实现仍保留本文所述 Provider 分发约束
- 结论:**这不是单一策划链路的问题**。各创作流程共用同一套 Provider 分发;第三方端点必须满足当前 `agentMode``apiKind` 和工具调用协议约束。缺陷 4 已修复其余限制仍按本文处理。
- 最新核对:2026-08-27,当前实现仍保留本文所述 Provider 分发约束
- 结论:**这不是单一策划链路的问题**。各创作流程共用同一套 Provider 分发;第三方端点必须满足当前 `agentMode``apiKind` 和工具调用协议约束。缺陷 4 已修复`llm.stream=true` 时 Provider tool-plan 现在按配置发送流式请求并在后端聚合完整响应,前端展示合同不变。其余限制仍按本文处理。
---
@@ -14,6 +14,7 @@
| 2 | `codex_app_server` 模式把第三方端点喂给 codex | apiKind≠openai_responses 时秒挂;否则 413 + 工具误用,180 秒超时后留下待核对的孤儿请求 | 模式前提未被约束 |
| 3 | `provider` 模式下 `tool_choice=required` 与 DeepSeek 思考模式互斥 | 首个 tool-plan 请求 400,整个 runtime 起不来 | 参数空间缺一个值 |
| 4 | 普通 action 批次带 plan update 时,两条预检规则互斥 | 「更新计划 + 委派专业 Agent」同一轮返回就报「批次成员身份或顺序不匹配」 | **本分支回归**(已修) |
| 5 | `llm.stream` 只记录配置,不驱动 Provider tool-plan 传输 | 要求 `stream=true` 的网关第一发 tool-plan 得到 HTTP 400,整轮不可用 | 传输配置失效(已修) |
缺陷 1~3 叠加的结果:**当前代码里没有任何一组配置能让 DeepSeek 跑起来**。缺陷 4 与 provider 无关,换成 `gpt-5.6-terra` 打通 LLM 链路后才暴露出来。
@@ -291,3 +292,26 @@ let expected_member_plan_update = batch
- DeepSeek 网关 413 的具体阈值,以及 `provider` 模式下 AGC 自组的请求体是否也会触顶。
---
## 9. 缺陷 5`llm.stream` 未作用于 Provider tool-plan(已修)
### 现象
`agentMode=provider``llm.stream=true` 时,审计与重试指纹记录 `stream=true`,但首个 tool-plan 仍调用 `LlmClient::run()`,请求体实际为 `stream=false`。只接受流式请求的 OpenAI 兼容网关返回 HTTP 400 `Stream must be set to true`;由于这是本地请求构造错误,重试同一请求无法恢复。
### 修复边界
Provider 的持久化重试分发与常规重试分发统一按 `llm.stream` 选择 `stream_run()` / `run()``stream_run()` 负责聚合文本、工具调用与终态,tool-plan 仍在响应完整后按现有协议解析、校验和交接;不把半截 tool-call 参数发布给前端,也不改变最终回复的 response-stream 合同。
### 回归
- `response_stream_uses_distinct_streamed_final_reply_for_responses_and_chat`:覆盖 Responses / Chat 两种 wire 的 tool-plan 与 final-reply 请求均发送 `stream=true`
- `response_stream_disabled_keeps_direct_planning_reply_to_one_request`:覆盖 `llm.stream=false` 时 tool-plan 仍发送 `stream=false` 且保持单请求直接收束。
- `background_agent_runtime_executes_streamed_native_function_tool_plan`:覆盖 Anthropic tool-use 分片在 Shell Runtime 中聚合为原生工具动作,并完成 tool-plan 协议审计与动作执行。
- `platform-llm` 既有 Chat / Responses 流式工具调用聚合用例继续覆盖分片工具参数装配。
### 升级边界
升级前遗留的 durable retry sidecar 若是在旧实现(审计记录 `stream=true`、实际发送 `stream=false`)期间创建,升级恢复后会按当前配置真实发送流式请求。该行为修正了配置与 wire 行为的一致性,但不保证与升级前已发出的失败请求字节一致;排查跨版本恢复时以 raw failure log 的请求快照为准。
---