修复真实游戏链路里 CC 执行器的四处阻塞

过程条目 id 按内容块类型区分,同 id 不同内容时补唯一后缀,避免整轮被判过程历史保存失败
资产下载的 fake-IP 窄例外覆盖 RFC 5180 的 IPv6 段,修复双栈解析下平台图片被拒绝下载
给 CC 子进程下发 MCP_TOOL_TIMEOUT,避免分钟级付费生成被 60 秒默认超时取消
复用 Codex 的进程树归属与退出证明,让 CC 回合能写 executor_stopped 并完成交付封口
This commit is contained in:
kdletters
2026-10-07 08:17:34 +08:00
parent 594ac321bf
commit fba012b42b
4 changed files with 222 additions and 24 deletions
@@ -18,6 +18,17 @@ const CLAUDE_NODE_ENV: &str = "GENARRATIVE_CLAUDE_NODE";
/// DirectProject 的一次 CC 回合可能包含代码生成、试玩、修复和再次复核;8 次会在真实
/// 游戏链路的修复阶段提前收口成 `Reached maximum number of turns (8)`。
const CLAUDE_DIRECT_MAX_TURNS: u32 = 16;
/// cc 侧 loopback MCP 工具调用的硬超时(毫秒)。
///
/// Claude Code 的 MCP 工具调用默认只有 60 秒硬墙,且"进度通知不会延长"它。AGC 的工具里
/// 包含真实付费与长耗时操作:图片生成一次约 56~70 秒,浏览器试玩、构建、验证也都在分钟级。
/// 现场证据(2026-10-06,项目 gameagent-cfd26650):模型在 13:05:30.760 发起两个
/// `agc_generate_image`,13:06:30.767(正好 60.0 秒)MCP 客户端返回
/// `{"content":"The operation timed out.","is_error":true}`,随后同轮的第二个调用被判
/// `direct-execution-interrupted: 原执行未正常结算`,本轮账本落成 `phase=interrupted`。
/// Codex 侧对同一套工具桥已设 `tool_timeout_sec=6600`(110 分钟),这里按同一口径对齐。
const CLAUDE_CODE_MCP_TOOL_TIMEOUT_MS: u64 = 110 * 60 * 1_000;
const CLAUDE_CODE_OUTPUT_MAX_BYTES: usize = 8 * 1024 * 1024;
const CLAUDE_CODE_PROMPT_MAX_BYTES: usize = 4 * 1024 * 1024;
const CLAUDE_CODE_MODEL_MAX_CHARS: usize = 200;
@@ -185,6 +196,24 @@ struct SidecarTurnResult {
result: serde_json::Value,
}
/// 把 cc sidecar 的进程树退出证明写进当前回合账本。
///
/// 交付封口(`finish_sealing`)要求 `executor_stopped`,而这项证明过去只有 Codex app-server
/// 会写:cc 回合因此永远停在 sealing(`deliveryReviews` 烧到上限后落成 interrupted),
/// 下一轮的修改与付费调用再被 `direct-execution-interrupted` 拦住。这里让 cc 执行器按同一
/// 口径上报"后台子树确实已退出",证明缺失时如实上报失败而不是缺席。
fn record_claude_sidecar_exit_proof(root: &Path, proven: Option<bool>) {
let Some(proven) = proven else {
return;
};
app_log!("agent.direct_codex.claude_sidecar_exit_proof proven={proven}");
if let Ok(session) = super::direct_execution::current(root) {
if let Err(error) = session.record_process_exit_proof(proven) {
app_log!("agent.direct_codex.claude_sidecar_exit_proof_unrecorded detail={error}");
}
}
}
/// sidecar 回合的失败分类:静默超时必须先杀进程树再收口,普通失败直接透传。
enum SidecarTurnFailure {
Silent { silent_ms: u64, events: usize },
@@ -215,6 +244,18 @@ async fn run_sidecar_turn(
let pid = child
.id()
.ok_or_else(|| "Claude Agent SDK sidecar 进程缺少 PID".to_string())?;
// 与 Codex app-server 同一套 Job Object 归属:只有把 sidecar 及其子进程树纳入受控
// 归属,回合结束时才拿得到"子树已退出"的可信证明,交付封口才有条件完成。
let owned_tree = match super::codex_app_server::process_tree::OwnedProcessTree::attach_with_role(
&child,
"ClaudeCode",
) {
Ok(tree) => Some(tree),
Err(error) => {
app_log!("agent.direct_codex.claude_sidecar_tree_unowned detail={error}");
None
}
};
let active_key = active_turn.map(|_| claude_project_key(root));
let active_alive = Arc::new(AtomicBool::new(true));
if let (Some(key), Some(client_turn_id)) = (active_key.clone(), active_turn) {
@@ -376,6 +417,15 @@ async fn run_sidecar_turn(
"Claude Agent SDK sidecar 非零退出;exitStatus={exit_status}"
)));
}
let proven = match owned_tree.as_ref() {
Some(tree) => Some(
tree.shutdown(&mut child)
.await
.is_ok_and(|proof| proof.confirmed()),
),
None => None,
};
record_claude_sidecar_exit_proof(root, proven);
Ok(SidecarTurnResult { result })
};
let outcome = match tokio::time::timeout(CLAUDE_CODE_TURN_MAX_DURATION, turn).await {
@@ -401,7 +451,21 @@ async fn run_sidecar_turn(
// 所有失败都走同一条收口:先杀进程树,再以有界窗口读取 stderr。之前只有
// 静默超时读取 stderr,cc 的非零退出 / stdout JSON 错误 / RPC error 会丢掉
// sidecar 给出的真正原因,最终只剩一条 transport 通用句。
let kill_detail = kill_claude_code_process_tree(pid).err();
// 失败路径同样要收口受控子树并留下退出证明:证明确认时不必再叠加 taskkill,
// 证明缺失时回退到原有的进程树强杀,保证不留后台执行器。
let proven = match owned_tree.as_ref() {
Some(tree) => Some(
tree.shutdown(&mut child)
.await
.is_ok_and(|proof| proof.confirmed()),
),
None => None,
};
record_claude_sidecar_exit_proof(root, proven);
let kill_detail = match proven {
Some(true) => None,
_ => kill_claude_code_process_tree(pid).err(),
};
let stderr_detail =
match tokio::time::timeout(CLAUDE_CODE_STDERR_DRAIN_TIMEOUT, stderr_task).await {
Ok(Ok(Ok(bytes))) => match String::from_utf8(bytes) {
@@ -690,7 +754,13 @@ fn configure_claude_code_environment(
.env("DISABLE_AUTOUPDATER", "1")
.env("CI", "1")
.env("NO_COLOR", "1")
.env("TERM", "dumb");
.env("TERM", "dumb")
// 不设这一项,CC 的 MCP 工具调用会回落到 60 秒硬超时,把分钟级的付费生成/试玩
// 直接取消掉,并让本轮租约未结算、后续工具全部被拒。
.env(
"MCP_TOOL_TIMEOUT",
CLAUDE_CODE_MCP_TOOL_TIMEOUT_MS.to_string(),
);
let session = crate::platform_session::current_platform_session();
let route = claude_code_route(llm, session.as_ref());
app_log!(
@@ -1138,7 +1208,13 @@ struct ClaudeCodeToolState {
#[derive(Default, Debug)]
struct ClaudeCodeStreamState {
current_message_id: String,
block_item_ids: HashMap<(String, usize), String>,
/// 内容块 → 过程条目 id。
///
/// 键必须带内容块类型:Claude Code 会把同一条 message 的内容块拆成多条 `assistant`
/// 事件(每条只含一个块),这些事件在各自数组里的下标都是 0。只按 `(message_id, index)`
/// 记忆,会让同一条消息里的 thinking 与 text 抢同一个 id(现场:`…:thinking:0` 先落盘,
/// 随后 text 复用该 id 且内容不同,历史层判冲突并让整轮变成“过程历史保存失败”)。
block_item_ids: HashMap<(String, String, usize), String>,
block_tool_ids: HashMap<(String, usize), String>,
tool_calls: HashMap<String, ClaudeCodeToolState>,
history_error: Option<String>,
@@ -1198,11 +1274,43 @@ fn claude_stream_observe(
}
}
/// 过程条目 id 撞车时的兜底条目:保留原内容,只把 id 换成带唯一后缀的取值。
///
/// 历史层只接受“同 id 同内容”的重复写;同 id 不同内容一律判冲突。Claude Code 会把同一条
/// message 的内容块拆成多条 `assistant` 事件,因此仍有无法用 `(message_id, kind, index)`
/// 区分的情况(例如同一条消息里的两个 text 块)。撞车时宁可多留一条过程历史,也不能让已经
/// 跑完的整轮被判成失败。缺失 id 时给一个自说明的兜底 id,保证追加一定能拿到唯一键。
fn completed_item_conflict_fallback(item: &serde_json::Value) -> serde_json::Value {
let mut retried = item.clone();
let fallback_id = match item.get("id").and_then(serde_json::Value::as_str) {
Some(id) if !id.trim().is_empty() => format!("{id}:{}", uuid::Uuid::new_v4()),
_ => format!("direct-cc:item:{}", uuid::Uuid::new_v4()),
};
if let Some(object) = retried.as_object_mut() {
object.insert("id".to_string(), serde_json::Value::String(fallback_id));
}
retried
}
impl ClaudeCodeStreamState {
/// 已完成过程与 Codex 共用项目历史出口,收到时即保存;终态失败不能撤掉前面已发生的事实。
fn persist_completed_item(&mut self, root: &Path, item: &serde_json::Value) {
if let Err(error) = append_direct_project_history_item_at(root, item) {
self.history_error.get_or_insert(error);
match append_direct_project_history_item_at(root, item) {
Ok(()) => {}
Err(error) if error.starts_with("DirectProject 历史 item id 冲突:") => {
// 同一个 message id 的内容块会在多条 `assistant` 事件里重复出现,条目 id 一旦
// 与已有条目“同 id 不同内容”,历史层会判冲突。这里按收尾回复的既有做法补一个
// 唯一后缀:把 id 撞车降级成多留一条过程历史,而不是让整轮变成过程历史保存失败。
let retried = completed_item_conflict_fallback(item);
if let Err(retry_error) = append_direct_project_history_item_at(root, &retried) {
self.history_error.get_or_insert(format!(
"过程历史追加失败:{retry_error}(首个条目冲突:{error})"
));
}
}
Err(error) => {
self.history_error.get_or_insert(error);
}
}
}
@@ -1215,7 +1323,7 @@ impl ClaudeCodeStreamState {
fn block_item_id(&mut self, message_id: &str, index: usize, kind: &str) -> String {
self.block_item_ids
.entry((message_id.to_string(), index))
.entry((message_id.to_string(), kind.to_string(), index))
.or_insert_with(|| format!("direct-cc:{message_id}:{kind}:{index}"))
.clone()
}
@@ -1495,13 +1603,13 @@ impl ClaudeCodeStreamState {
}
let item_id = self.block_item_id(message_id, index, "text");
self.partial_history
.complete_item(&serde_json::json!({"id": item_id}));
.complete_item(&serde_json::json!({"id": item_id.clone()}));
let at = crate::agent::now_ms();
crate::agent::append_thread_event(
&crate::agent::thread_id_for_project(root),
ThreadEvent::item_completed(
ThreadItem::Message {
item_id,
item_id: item_id.clone(),
role: "assistant".to_string(),
text: text.clone(),
at,
@@ -1509,12 +1617,15 @@ impl ClaudeCodeStreamState {
at,
),
);
self.persist_completed_item(root, &serde_json::json!({
"type": "message",
"role": "assistant",
"id": self.block_item_ids.get(&(message_id.to_string(), index)).cloned().unwrap_or_default(),
"content": [{"type": "output_text", "text": text}],
}));
self.persist_completed_item(
root,
&serde_json::json!({
"type": "message",
"role": "assistant",
"id": item_id,
"content": [{"type": "output_text", "text": text}],
}),
);
self.assistant_texts.push(text);
claude_stream_observe(
observer,
@@ -1870,6 +1981,54 @@ mod tests {
request
}
#[test]
fn claude_block_item_id_keeps_thinking_and_text_apart() {
let mut state = ClaudeCodeStreamState::default();
let thinking = state.block_item_id("msg_1", 0, "thinking");
let text = state.block_item_id("msg_1", 0, "text");
assert_eq!(thinking, "direct-cc:msg_1:thinking:0");
assert_eq!(text, "direct-cc:msg_1:text:0");
assert_ne!(thinking, text);
// 同一个块重复出现(Claude Code 按块分事件重发同一 message)必须复用同一个 id。
assert_eq!(state.block_item_id("msg_1", 0, "thinking"), thinking);
assert_eq!(state.block_item_id("msg_1", 0, "text"), text);
}
#[test]
fn claude_completed_item_conflict_fallback_keeps_content_and_unique_id() {
let item = serde_json::json!({
"type": "reasoning",
"id": "direct-cc:msg_1:thinking:0",
"text": "已经落盘过的思考内容",
});
let fallback = completed_item_conflict_fallback(&item);
assert_eq!(fallback["type"], item["type"]);
assert_eq!(fallback["text"], item["text"]);
let fallback_id = fallback["id"].as_str().expect("fallback id");
assert_ne!(fallback_id, "direct-cc:msg_1:thinking:0");
assert!(
fallback_id.starts_with("direct-cc:msg_1:thinking:0:"),
"兜底 id 必须保留原 id 便于排查:{fallback_id}"
);
assert_ne!(
fallback["id"],
completed_item_conflict_fallback(&item)["id"],
"兜底 id 每次都必须是新的唯一取值"
);
let without_id = serde_json::json!({ "type": "reasoning", "text": "无 id" });
let fallback_without_id = completed_item_conflict_fallback(&without_id);
assert!(
fallback_without_id["id"]
.as_str()
.expect("fallback id")
.starts_with("direct-cc:item:"),
"缺少 id 时也要给出自说明的兜底 id"
);
}
#[test]
fn parses_claude_code_text_result() {
let response = parse_claude_code_result(
@@ -2324,5 +2483,11 @@ mod tests {
] {
assert!(!envs.contains_key(name), "{name} 不应下发到 sidecar");
}
// MCP 工具调用必须显式抬高超时:默认 60 秒会把分钟级的付费生成/试玩取消掉,
// 让本轮租约未结算并拒绝后续工具调用。
assert_eq!(
envs.get("MCP_TOOL_TIMEOUT"),
Some(&Some(CLAUDE_CODE_MCP_TOOL_TIMEOUT_MS.to_string()))
);
}
}
@@ -11,7 +11,7 @@ use tokio::sync::{mpsc, oneshot, Mutex, Notify};
mod direct_project_history_wire;
use direct_project_history_wire::build_direct_project_history_injection_params;
mod process_tree;
pub(crate) mod process_tree;
use process_tree::{OwnedProcessTree, ProcessTreeExitProof};
mod direct_project_identity;
mod execution;
@@ -30,7 +30,7 @@ impl ProcessTreeExitProof {
}
}
pub(super) struct OwnedProcessTree {
pub(crate) struct OwnedProcessTree {
pid: u32,
#[cfg(unix)]
start_identity: String,
@@ -41,7 +41,13 @@ pub(super) struct OwnedProcessTree {
impl OwnedProcessTree {
/// 调用方只能在此成功后发送 initialize,保证所有受控模型执行都在归属内。
pub(super) fn attach(child: &Child) -> Result<Self, String> {
pub(crate) fn attach(child: &Child) -> Result<Self, String> {
Self::attach_with_role(child, "Codex")
}
/// 具名角色版本:Claude Code sidecar 与 Codex app-server 都是受控执行器,
/// 共用同一套 Job Object 归属与退出证明,只是 Job 名字要能区分来源便于排障。
pub(crate) fn attach_with_role(child: &Child, role: &str) -> Result<Self, String> {
let pid = child.id().ok_or("app-server-process-owner-missing")?;
#[cfg(unix)]
let start_identity =
@@ -52,7 +58,7 @@ impl OwnedProcessTree {
#[cfg(windows)]
let job = crate::process_session::WindowsProcessJob::assign_tokio_named(
child,
&format!("Local\\AGCCodex-{}", uuid::Uuid::new_v4()),
&format!("Local\\AGC{role}-{}", uuid::Uuid::new_v4()),
)
.map_err(|error| format!("app-server-process-job-unavailable:{error}"))?;
#[cfg(not(any(windows, unix)))]
@@ -68,7 +74,7 @@ impl OwnedProcessTree {
})
}
pub(super) async fn shutdown(&self, child: &mut Child) -> Result<ProcessTreeExitProof, String> {
pub(crate) async fn shutdown(&self, child: &mut Child) -> Result<ProcessTreeExitProof, String> {
let mut stored = self.exit_proof.lock().await;
if let Some(proof) = stored.as_ref() {
return proof.clone();
@@ -78,7 +84,7 @@ impl OwnedProcessTree {
result
}
pub(super) async fn recorded_exit_proof(&self) -> Result<ProcessTreeExitProof, String> {
pub(crate) async fn recorded_exit_proof(&self) -> Result<ProcessTreeExitProof, String> {
self.exit_proof
.lock()
.await
@@ -1208,6 +1208,9 @@ fn external_asset_host_is_private(host: &str) -> bool {
|| address.is_unique_local()
|| address.is_unicast_link_local()
|| address.is_multicast()
// 与 IPv4 的 198.18/15 对齐:RFC 5180 的 2001:2::/48 同样只有“经已鉴权
// 稳定引用换签”这一条窄例外能放行,用户或上游直接给的地址仍然失败关闭。
|| external_asset_ip_is_proxy_benchmark(std::net::IpAddr::V6(address))
|| address
.to_ipv4_mapped()
.is_some_and(|mapped| external_asset_host_is_private(&mapped.to_string()))
@@ -1221,7 +1224,14 @@ fn external_asset_ip_is_proxy_benchmark(address: std::net::IpAddr) -> bool {
let [first, second, ..] = address.octets();
first == 198 && (18..=19).contains(&second)
}
std::net::IpAddr::V6(_) => false,
// Clash / Clash.Meta 的 IPv6 fake-IP 默认段是 RFC 5180 的 2001:2::/48,与 IPv4 的
// RFC 2544 段是同一类“透明代理占位地址”。双栈解析会同时返回两个 fake-IP(现场:
// OSS 域名解析出 198.18.0.162 + 2001:2::91),只承认 IPv4 会让“全部地址都是
// fake-IP”永远不成立,把正常的平台资产下载误判成私网请求并拒绝。
std::net::IpAddr::V6(address) => {
let segments = address.segments();
segments[0] == 0x2001 && segments[1] == 0x0002 && segments[2] == 0x0000
}
}
}
@@ -1235,9 +1245,9 @@ fn external_asset_resolved_addresses_are_safe(
{
return true;
}
// Clash 等透明代理会把公网域名映射到 RFC 2544 的 198.18.0.0/15 fake-IP。
// 只有经已鉴权 objectKey/legacy path 换签得到的 URL 可以使用这项窄例外;
// 用户或上游直接提供的 URL、其它私网地址及公私混合解析仍然失败关闭。
// Clash 等透明代理会把公网域名映射到 RFC 2544 的 198.18.0.0/15(IPv4)和 RFC 5180 的
// 2001:2::/48(IPv6)fake-IP。只有经已鉴权 objectKey/legacy path 换签得到的 URL 可以使
// 用这项窄例外;用户或上游直接提供的 URL、其它私网地址及公私混合解析仍然失败关闭。
came_from_stable_reference
&& !addresses.is_empty()
&& addresses
@@ -3061,11 +3071,28 @@ mod tests {
assert!(external_asset_resolved_addresses_are_safe(&fake_ip, true));
assert!(!external_asset_resolved_addresses_are_safe(&fake_ip, false));
// Clash 双栈 fake-IP:IPv6 侧是 RFC 5180 的 2001:2::/48,必须与 IPv4 一起放行,
// 否则双栈解析下“全部地址都是 fake-IP”永不成立,正常资产下载会被误拒。
let fake_ipv6 = ["[2001:2::91]:443".parse().expect("parse proxy fake IPv6")];
assert!(external_asset_resolved_addresses_are_safe(&fake_ipv6, true));
assert!(!external_asset_resolved_addresses_are_safe(
&fake_ipv6, false
));
let dual_stack_fake_ips = [
"198.18.0.162:443".parse().expect("parse proxy fake IPv4"),
"[2001:2::91]:443".parse().expect("parse proxy fake IPv6"),
];
assert!(external_asset_resolved_addresses_are_safe(
&dual_stack_fake_ips,
true
));
for address in [
"127.0.0.1:443",
"10.0.0.1:443",
"169.254.169.254:80",
"[::1]:443",
"[fd00::1]:443",
] {
let addresses = [address.parse().expect("parse private address")];
assert!(!external_asset_resolved_addresses_are_safe(