From fbfc0b62b9f0a1996cf52c42142c9b9298c97620 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=8E=8B=E5=BE=B7=E5=AE=87?= Date: Thu, 1 Oct 2026 16:04:02 +0800 Subject: [PATCH] =?UTF-8?q?=E5=AE=A1=E6=89=B9=E8=A7=A3=E6=9E=90=E7=A7=BB?= =?UTF-8?q?=E5=87=BA=E4=B8=B4=E7=95=8C=E5=8C=BA=E5=90=8E=E5=86=8D=E7=95=99?= =?UTF-8?q?=E7=97=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - respond 的 target 解析抽成 resolve_approval_gate,锁内只做状态判定与 ticket 登记 - 新增 ApprovalGateOutcome(Proceed / Cached / Deny),拒绝原因带出临界区,state 释放后才调用 denied 写日志 - 拒绝留痕不再持有 state 做同步文件 I/O(app_log 写入与日志轮转),不再拉长审批临界区 - 15 条拒绝原因分类逐条等价迁移,行为与返回形状不变 --- .../src/agent/codex_app_server/execution.rs | 234 ++++++++++-------- 1 file changed, 134 insertions(+), 100 deletions(-) diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/codex_app_server/execution.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/codex_app_server/execution.rs index 10ab21172..2c31b17c9 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/codex_app_server/execution.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/codex_app_server/execution.rs @@ -210,6 +210,19 @@ struct Ticket { response: Option, } +/// 锁内解析一次审批 / 交互请求的结果。 +/// +/// 解析只做状态判定与登记;拒绝原因的留痕由调用方在**释放 `state` 之后**完成,避免持有审批 +/// 临界区做同步文件 I/O(见 `ExecutionAdapter::respond`)。 +enum ApprovalGateOutcome { + /// 已登记 ticket,进入租约申请与放行阶段。 + Proceed(Target), + /// 该 request id 已有结论,直接复用,不再申请租约。 + Cached(Value), + /// 拒绝,附留痕用的原因分类。 + Deny(&'static str), +} + #[derive(Default)] struct ProtocolState { turn_id: Option, @@ -558,114 +571,135 @@ impl ExecutionAdapter { self.spawn_settlement(entries); } + /// 锁内解析一次审批 / 交互请求:登记待放行的 ticket,或给出拒绝原因。 + /// + /// **只做状态判定与登记,不做 I/O、不写日志**:调用方必须在释放 `state` 之后再留痕。 + fn resolve_approval_gate( + &self, + state: &mut ProtocolState, + id: u64, + method: &str, + params: &Value, + fingerprint: &str, + ) -> ApprovalGateOutcome { + if !self.matches_scope(state, params) { + return ApprovalGateOutcome::Deny("request-out-of-turn-scope"); + } + if let Some(ticket) = state.tickets.get(&id) { + return if ticket.fingerprint == fingerprint { + match ticket.response.clone() { + Some(response) => ApprovalGateOutcome::Cached(response), + None => ApprovalGateOutcome::Deny("ticket-without-response"), + } + } else { + ApprovalGateOutcome::Deny("request-id-fingerprint-mismatch") + }; + } + // Approval tasks run concurrently; numeric request IDs need not enter + // this mutex in order. Retain seen IDs rather than rejecting by high-water. + if state.seen_request_ids.len() >= 8192 || !state.seen_request_ids.insert(id) { + return ApprovalGateOutcome::Deny("request-id-replayed-or-capacity"); + } + let target = match method { + "item/commandExecution/requestApproval" | "item/fileChange/requestApproval" => { + let Some(item_id) = identity(params.get("itemId")) else { + return ApprovalGateOutcome::Deny("approval-without-item-id"); + }; + let Some(item) = state.items.get(item_id) else { + return ApprovalGateOutcome::Deny("approval-for-unknown-item"); + }; + let expected = if method == "item/fileChange/requestApproval" { + ItemKind::Patch + } else { + ItemKind::Command + }; + if item.kind != expected || item.terminal.is_some() || item.lease.is_some() { + return ApprovalGateOutcome::Deny("item-not-pending-for-approval"); + } + if state.tickets.values().any(|ticket| matches!(&ticket.target, Target::Item(existing, _) if existing == item_id)) { + return ApprovalGateOutcome::Deny("item-approval-already-pending"); + } + Target::Item( + item_id.to_string(), + if expected == ItemKind::Patch { + EffectKind::Write + } else { + EffectKind::Execute + }, + ) + } + "mcpServer/elicitation/request" => { + if params + .pointer("/_meta/codex_approval_kind") + .and_then(Value::as_str) + != Some("mcp_tool_call") + || params.get("mode").and_then(Value::as_str) != Some("form") + { + return ApprovalGateOutcome::Deny("unsupported-mcp-elicitation-shape"); + } + let Some(server) = params.get("serverName").and_then(Value::as_str) else { + return ApprovalGateOutcome::Deny("mcp-elicitation-without-server-name"); + }; + if server == "agc_tools" || !self.third_party_servers.contains(server) { + return ApprovalGateOutcome::Deny("mcp-server-not-eligible-for-elicitation"); + } + let Some(arguments) = params.pointer("/_meta/tool_params") else { + return ApprovalGateOutcome::Deny("mcp-elicitation-without-tool-params"); + }; + let key = cohort_key(server, arguments); + let Some(cohort) = state.cohorts.get(&key) else { + return ApprovalGateOutcome::Deny("mcp-cohort-not-registered"); + }; + let pending_members = cohort + .members + .iter() + .filter(|member| { + state + .items + .get(*member) + .is_some_and(|entry| entry.terminal.is_none()) + }) + .count(); + let seen_tickets = state.tickets.values().filter(|ticket| matches!(&ticket.target, Target::Cohort(existing) if existing == &key)).count(); + if pending_members == 0 || cohort.members.len() <= seen_tickets { + return ApprovalGateOutcome::Deny("mcp-cohort-already-covered"); + } + state.cohorts.get_mut(&key).unwrap().admissions += 1; + Target::Cohort(key) + } + _ => return ApprovalGateOutcome::Deny("unsupported-approval-method"), + }; + if state.tickets.len() >= MAX_REQUEST_CACHE { + state.tickets.retain(|_, ticket| ticket.response.is_none()); + } + state.tickets.insert( + id, + Ticket { + fingerprint: fingerprint.to_string(), + target: target.clone(), + response: None, + }, + ); + ApprovalGateOutcome::Proceed(target) + } + pub(super) async fn respond(self: &Arc, id: u64, method: &str, params: &Value) -> Value { if self.closed.load(Ordering::Acquire) || self.is_host_ending() { return self.denied(id, method, "session-closed-or-host-ending"); } let fingerprint = cohort_key(method, params); - let target = { + // 解析与登记在锁内一次完成;`app_log!` 是同步文件 I/O(必要时还会轮转日志),所以拒绝 + // 留痕必须等 `state` 释放之后再做,否则会把审批临界区拉长到磁盘延迟。 + let outcome = { let Ok(mut state) = self.state.lock() else { return self.denied(id, method, "state-lock-poisoned"); }; - if !self.matches_scope(&state, params) { - return self.denied(id, method, "request-out-of-turn-scope"); - } - if let Some(ticket) = state.tickets.get(&id) { - return if ticket.fingerprint == fingerprint { - ticket - .response - .clone() - .unwrap_or_else(|| self.denied(id, method, "ticket-without-response")) - } else { - self.denied(id, method, "request-id-fingerprint-mismatch") - }; - } - // Approval tasks run concurrently; numeric request IDs need not enter - // this mutex in order. Retain seen IDs rather than rejecting by high-water. - if state.seen_request_ids.len() >= 8192 || !state.seen_request_ids.insert(id) { - return self.denied(id, method, "request-id-replayed-or-capacity"); - } - let target = match method { - "item/commandExecution/requestApproval" | "item/fileChange/requestApproval" => { - let Some(item_id) = identity(params.get("itemId")) else { - return self.denied(id, method, "approval-without-item-id"); - }; - let Some(item) = state.items.get(item_id) else { - return self.denied(id, method, "approval-for-unknown-item"); - }; - let expected = if method == "item/fileChange/requestApproval" { - ItemKind::Patch - } else { - ItemKind::Command - }; - if item.kind != expected || item.terminal.is_some() || item.lease.is_some() { - return self.denied(id, method, "item-not-pending-for-approval"); - } - if state.tickets.values().any(|ticket| matches!(&ticket.target, Target::Item(existing, _) if existing == item_id)) { - return self.denied(id, method, "item-approval-already-pending"); - } - Target::Item( - item_id.to_string(), - if expected == ItemKind::Patch { - EffectKind::Write - } else { - EffectKind::Execute - }, - ) - } - "mcpServer/elicitation/request" => { - if params - .pointer("/_meta/codex_approval_kind") - .and_then(Value::as_str) - != Some("mcp_tool_call") - || params.get("mode").and_then(Value::as_str) != Some("form") - { - return self.denied(id, method, "unsupported-mcp-elicitation-shape"); - } - let Some(server) = params.get("serverName").and_then(Value::as_str) else { - return self.denied(id, method, "mcp-elicitation-without-server-name"); - }; - if server == "agc_tools" || !self.third_party_servers.contains(server) { - return self.denied(id, method, "mcp-server-not-eligible-for-elicitation"); - } - let Some(arguments) = params.pointer("/_meta/tool_params") else { - return self.denied(id, method, "mcp-elicitation-without-tool-params"); - }; - let key = cohort_key(server, arguments); - let Some(cohort) = state.cohorts.get(&key) else { - return self.denied(id, method, "mcp-cohort-not-registered"); - }; - let pending_members = cohort - .members - .iter() - .filter(|member| { - state - .items - .get(*member) - .is_some_and(|entry| entry.terminal.is_none()) - }) - .count(); - let seen_tickets = state.tickets.values().filter(|ticket| matches!(&ticket.target, Target::Cohort(existing) if existing == &key)).count(); - if pending_members == 0 || cohort.members.len() <= seen_tickets { - return self.denied(id, method, "mcp-cohort-already-covered"); - } - state.cohorts.get_mut(&key).unwrap().admissions += 1; - Target::Cohort(key) - } - _ => return self.denied(id, method, "unsupported-approval-method"), - }; - if state.tickets.len() >= MAX_REQUEST_CACHE { - state.tickets.retain(|_, ticket| ticket.response.is_none()); - } - state.tickets.insert( - id, - Ticket { - fingerprint, - target: target.clone(), - response: None, - }, - ); - target + self.resolve_approval_gate(&mut state, id, method, params, &fingerprint) + }; + let target = match outcome { + ApprovalGateOutcome::Proceed(target) => target, + ApprovalGateOutcome::Cached(response) => return response, + ApprovalGateOutcome::Deny(reason) => return self.denied(id, method, reason), }; self.settling.fetch_add(1, Ordering::AcqRel); let root = self.root.clone();