审批解析移出临界区后再留痕
- respond 的 target 解析抽成 resolve_approval_gate,锁内只做状态判定与 ticket 登记 - 新增 ApprovalGateOutcome(Proceed / Cached / Deny),拒绝原因带出临界区,state 释放后才调用 denied 写日志 - 拒绝留痕不再持有 state 做同步文件 I/O(app_log 写入与日志轮转),不再拉长审批临界区 - 15 条拒绝原因分类逐条等价迁移,行为与返回形状不变
This commit is contained in:
@@ -210,6 +210,19 @@ struct Ticket {
|
||||
response: Option<Value>,
|
||||
}
|
||||
|
||||
/// 锁内解析一次审批 / 交互请求的结果。
|
||||
///
|
||||
/// 解析只做状态判定与登记;拒绝原因的留痕由调用方在**释放 `state` 之后**完成,避免持有审批
|
||||
/// 临界区做同步文件 I/O(见 `ExecutionAdapter::respond`)。
|
||||
enum ApprovalGateOutcome {
|
||||
/// 已登记 ticket,进入租约申请与放行阶段。
|
||||
Proceed(Target),
|
||||
/// 该 request id 已有结论,直接复用,不再申请租约。
|
||||
Cached(Value),
|
||||
/// 拒绝,附留痕用的原因分类。
|
||||
Deny(&'static str),
|
||||
}
|
||||
|
||||
#[derive(Default)]
|
||||
struct ProtocolState {
|
||||
turn_id: Option<String>,
|
||||
@@ -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<Self>, 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();
|
||||
|
||||
Reference in New Issue
Block a user