修复 Direct 回合收尾与进程退出判断
Project CI / AI game creator shell Rust crates (pull_request) Successful in 5m43s
Project CI / AI game creator shell Rust lane 2/2 (pull_request) Successful in 6m21s
Project CI / AI game creator shell Rust lane 1/2 (pull_request) Successful in 7m20s
Project CI / Backend tests (pull_request) Successful in 8m32s
Project CI / Frontend tests (pull_request) Successful in 3m13s
Project CI / AI game creator shell web tests (pull_request) Successful in 2m59s
Project CI / Repository checks (pull_request) Successful in 5m12s
Project CI / Native shell tests (pull_request) Successful in 6m57s

使用受控归属退出证明完成 Unix 回合,保留完整子树证明边界
保守处理 Linux 进程扫描错误,避免将未确认的活跃进程判为空组
优先返回收尾时已持久化的预算耗尽报告,保留其他关闭错误
阻止正常关闭的迟到通知创建或覆盖成功回合报告
补充进程扫描、完整响应和预算回归测试并同步规范与验收记录
This commit is contained in:
2026-10-05 07:04:51 +01:00
parent 063905ee5f
commit 691e51e85c
10 changed files with 396 additions and 28 deletions
@@ -892,7 +892,7 @@ impl ExecutionAdapter {
/// 中断会话,`closed` 与阶段都要等那个任务跑到才变,所以"标志已置、阶段未变"的窗口里到达的
/// 通道断开 / 中断都是宿主自己收尾的结果,不能记成 `transport-failed`。
///
/// 不记失败事实不等于不收束:原因照样写进报告(`interrupt` 会把它追加进去),便于核对。
/// 非宿主关闭仍记录失败原因;宿主关闭的通知只在失败终态追加,正常回复不受影响。
///
/// **事实要落在适配器上,不能落在调用点的局部变量里。** 回合还开着的时候,看门狗会在同一个
/// `inner.closed` 标志上把本轮收束掉(见 [`Self::start_watchdog`]),谁先谁后取决于调度,而终态
@@ -904,8 +904,7 @@ impl ExecutionAdapter {
let reason = failure.diagnostic_detail();
if self.is_closed() {
// 宿主自己收尾:连接是我们先关的,紧随其后的 `TransportClosed` 只是收尾的副产物。
// 只把原因留给报告,不改阶段——否则正常的宿主收尾会被改写成 `interrupted`
// 并把回执文案换成「执行通道已断开」(见 pitfalls 2026-09-28)。
// 不改阶段,且仅在已有失败终态追加说明,避免正常回合被通道关闭报告截断。
let session = Arc::clone(&self.session);
let note = reason.clone();
let _ = tokio::task::spawn_blocking(move || session.append_terminal_note(note)).await;
@@ -1115,7 +1114,7 @@ impl ExecutionAdapter {
self.closed.store(true, Ordering::Release);
let proven = shutdown_game_creator_codex_app_server_inner(inner, "宿主执行预算或交付收尾")
.await
.map(|proof| proof.confirmed())
.map(|proof| proof.owned_scope_retired())
.unwrap_or(false);
self.drain().await;
let session = Arc::clone(&self.session);
@@ -1187,7 +1186,7 @@ impl ExecutionAdapter {
"模型本次执行结束,回收原生后台子树",
)
.await
.map(|proof| proof.confirmed())
.map(|proof| proof.owned_scope_retired())
.unwrap_or(false);
self.drain().await;
let session = Arc::clone(&self.session);
@@ -1485,8 +1484,13 @@ mod tests {
assert!(adapter.turn_failure().is_none());
assert!(adapter.host_stop_requested());
// 原因照样进报告:不算失败不等于不用记。
assert!(adapter.report().contains("模型本次执行结束"));
// 正常关闭通知不能为尚未完成的回合创建终态报告。
assert!(adapter
.session
.snapshot()
.unwrap()
.terminal_report
.is_none());
// 阶段不能被改成 Interrupted:宿主自己收尾时,账本必须保留它自己的结论。
// 这一条是 2026-09-28「CLI 每轮回执都变成『执行通道已断开』」缺陷的回归判据——
// 缺陷版本里 `fail_turn` 会无条件 `session.interrupt(reason)`,把阶段打成 Interrupted。
@@ -8242,17 +8242,23 @@ done
#[cfg(unix)]
#[tokio::test]
async fn direct_project_turn_does_not_forward_codex_user_echo_as_chat_items() {
run_direct_response_fixture(false).await;
run_direct_response_fixture(false, false).await;
}
#[cfg(unix)]
#[tokio::test]
async fn no_contract_operation_completes_with_original_response_after_group_shutdown() {
run_direct_response_fixture(false, true).await;
}
#[cfg(unix)]
#[tokio::test]
async fn ready_contract_allows_operation_then_complete_response_before_host_shutdown() {
run_direct_response_fixture(true).await;
run_direct_response_fixture(true, true).await;
}
#[cfg(unix)]
async fn run_direct_response_fixture(ready_contract: bool) {
async fn run_direct_response_fixture(ready_contract: bool, has_operation: bool) {
use std::os::unix::fs::PermissionsExt;
let temp = tempfile::tempdir().expect("temp dir");
@@ -8283,7 +8289,7 @@ while IFS= read -r line; do
esac
done
"#;
let script = if ready_contract {
let script = if has_operation {
script.replace(" printf '%s\\n' '{\"method\":\"item/agentMessage/delta\"", r#" printf '%s\n' '{"method":"item/started","params":{"threadId":"thread-echo","turnId":"turn-echo","item":{"id":"cmd-ready","type":"commandExecution","status":"inProgress"}}}'
printf '%s\n' '{"id":999,"method":"item/commandExecution/requestApproval","params":{"threadId":"thread-echo","turnId":"turn-echo","itemId":"cmd-ready"}}'
IFS= read -r approval
@@ -8360,10 +8366,11 @@ done
snapshot.project_id = direct_codex_canonical_project_identity(&project)
.expect("canonical Provider snapshot identity")
.1;
let _execution_guard = super::super::direct_execution::register_for_test(execution)
.expect("register host execution");
let _execution_guard =
super::super::direct_execution::register_for_test(Arc::clone(&execution))
.expect("register host execution");
let mut observer = |_observation| {};
connection
let response = connection
.run_turn_with_direct_observer_and_history(
&snapshot,
&llm,
@@ -8379,6 +8386,42 @@ done
.expect("run direct-project turn");
drop(observer);
if has_operation {
let proof = connection
.inner
.process_tree
.recorded_exit_proof()
.await
.unwrap();
assert!(proof.owned_scope_retired());
assert!(!proof.confirmed(), "Unix 退出仍只证明受控进程组范围");
assert!(execution.snapshot().unwrap().executor_stopped);
}
if !ready_contract {
assert_eq!(
execution.snapshot().unwrap().phase,
super::super::direct_execution::ExecutionPhase::Working
);
}
let report = super::super::direct_delivery::review_reply(&project, &execution)
.await
.unwrap();
let state = execution.snapshot().unwrap();
assert_eq!(
state.phase,
super::super::direct_execution::ExecutionPhase::Completed
);
if ready_contract {
let report = report.expect("已就绪合同必须返回宿主验收报告");
assert!(report.contains("本轮已完成宿主验收"));
assert_eq!(response.text, report);
assert_eq!(state.terminal_report.as_deref(), Some(report.as_str()));
} else {
assert_eq!(response.text, "好的");
assert!(report.is_none());
assert!(state.terminal_report.is_none());
}
let consumed =
crate::agent::consume_thread(&bootstrap.subscription_id).expect("consume events");
// 回合起止必须与开口用户条目同源:前端在「只有锚点 + 历史、运行态为空」的回合里靠这个
@@ -20,7 +20,7 @@ impl ProcessTreeExitProof {
self.main_process_exited && self.observed_tree_empty && self.full_tree_verified
}
/// 生命周期退役仅证明已控制的归属为空;不能替代验收使用的完整子树证明。
/// 回合收尾及交付只要求受控归属退役;进程组不能冒充完整子树证明。
pub(super) fn owned_scope_retired(&self) -> bool {
match self.scope {
"windows-job" => self.confirmed(),
@@ -191,28 +191,80 @@ impl OwnedProcessTree {
/// 僵尸对退役语义等价于已退出,扫描 /proc 时只把非僵尸成员视为存活。
#[cfg(all(unix, target_os = "linux"))]
fn linux_process_group_has_live_member(pgid: i32) -> Option<bool> {
for entry in std::fs::read_dir("/proc").ok()? {
let entry = entry.ok()?;
let Ok(stat) = std::fs::read_to_string(entry.path().join("stat")) else {
linux_process_group_members(
pgid,
std::fs::read_dir("/proc")
.ok()?
.map(|entry| entry.map(|entry| entry.path())),
)
}
#[cfg(target_os = "linux")]
fn linux_process_group_members(
pgid: i32,
entries: impl IntoIterator<Item = std::io::Result<std::path::PathBuf>>,
) -> Option<bool> {
let mut uncertain = false;
for entry in entries {
let Ok(path) = entry else {
uncertain = true;
continue;
};
if !path
.file_name()
.and_then(|name| name.to_str())
.is_some_and(|name| !name.is_empty() && name.bytes().all(|byte| byte.is_ascii_digit()))
{
continue;
}
let stat = match std::fs::read_to_string(path.join("stat")) {
Ok(stat) => stat,
// /proc 枚举与读取之间 PID 可消失;仅在确认目录已消失时忽略。
Err(error)
if error.kind() == std::io::ErrorKind::NotFound
&& path.try_exists().is_ok_and(|exists| !exists) =>
{
continue
}
Err(_) => {
uncertain = true;
continue;
}
};
let Some(close) = stat.rfind(')') else {
uncertain = true;
continue;
};
let mut fields = stat[close + 1..].split_whitespace();
let (Some(state), Some(_ppid), Some(member_pgid)) =
(fields.next(), fields.next(), fields.next())
else {
uncertain = true;
continue;
};
if member_pgid.parse::<i32>().ok() != Some(pgid) {
let Ok(member_pgid) = member_pgid.parse::<i32>() else {
uncertain = true;
continue;
};
if member_pgid != pgid {
continue;
}
if !matches!(
state,
"R" | "S" | "D" | "Z" | "T" | "t" | "X" | "x" | "K" | "W" | "P" | "I"
) {
uncertain = true;
continue;
}
if state != "Z" {
return Some(true);
}
}
Some(false)
if uncertain {
None
} else {
Some(false)
}
}
impl Drop for OwnedProcessTree {
@@ -231,6 +283,115 @@ mod tests {
const FIXTURE: &str =
"agent::codex_app_server::process_tree::tests::owned_tree_process_fixture";
#[cfg(target_os = "linux")]
#[test]
fn process_scan_distinguishes_live_zombie_foreign_and_disappeared_pids() {
let root = tempfile::tempdir().unwrap();
let zombie = root.path().join("101");
let foreign = root.path().join("102");
let live = root.path().join("103");
for (path, stat) in [
(&zombie, "101 (name with ) brackets) Z 1 100"),
(&foreign, "102 (other) R 1 200"),
(&live, "103 (worker) S 1 100"),
] {
std::fs::create_dir(path).unwrap();
std::fs::write(path.join("stat"), stat).unwrap();
}
let gone = root.path().join("104");
let non_pid = root.path().join("self");
assert_eq!(
linux_process_group_members(
100,
[Ok(zombie.clone()), Ok(foreign), Ok(gone), Ok(non_pid)]
),
Some(false)
);
assert_eq!(
linux_process_group_members(100, [Ok(zombie), Ok(live)]),
Some(true)
);
}
#[cfg(target_os = "linux")]
#[test]
fn incomplete_process_scan_never_proves_empty() {
let root = tempfile::tempdir().unwrap();
let pid = root.path().join("101");
std::fs::create_dir(&pid).unwrap();
// PID 目录仍在:缺失 stat 与进程确实消失不同。
assert_eq!(linux_process_group_members(100, [Ok(pid.clone())]), None);
std::fs::create_dir(pid.join("stat")).unwrap();
assert_eq!(linux_process_group_members(100, [Ok(pid.clone())]), None);
std::fs::remove_dir(pid.join("stat")).unwrap();
for malformed in [
"invalid",
"101 (worker) S",
"101 (worker) S 1 bad",
"101 (worker) unknown 1 100",
] {
std::fs::write(pid.join("stat"), malformed).unwrap();
assert_eq!(linux_process_group_members(100, [Ok(pid.clone())]), None);
}
std::fs::write(pid.join("stat"), "101 (worker) Z 1 100").unwrap();
assert_eq!(
linux_process_group_members(
100,
[
Err(std::io::Error::from(std::io::ErrorKind::PermissionDenied)),
Ok(pid.clone())
]
),
None
);
// 枚举错误不提前停止:其后有确定的活跃成员时仍直接报告存活。
std::fs::write(pid.join("stat"), "101 (worker) R 1 100").unwrap();
assert_eq!(
linux_process_group_members(
100,
[
Err(std::io::Error::from(std::io::ErrorKind::PermissionDenied)),
Ok(pid)
]
),
Some(true)
);
}
#[cfg(target_os = "linux")]
#[test]
fn real_zombie_is_not_a_live_group_member_before_reaping() {
use std::os::unix::process::CommandExt;
let mut child = std::process::Command::new("sleep")
.arg("30")
.process_group(0)
.spawn()
.unwrap();
let pgid = child.id() as i32;
let live = linux_process_group_has_live_member(pgid);
child.kill().unwrap();
let stat_path = format!("/proc/{pgid}/stat");
let deadline = std::time::Instant::now() + Duration::from_secs(2);
let mut zombie = false;
while std::time::Instant::now() < deadline {
zombie = std::fs::read_to_string(&stat_path).is_ok_and(|stat| {
stat.rsplit_once(')')
.is_some_and(|(_, fields)| fields.split_whitespace().next() == Some("Z"))
});
if zombie {
break;
}
std::thread::sleep(Duration::from_millis(10));
}
let group_exists = unsafe { libc::kill(-pgid, 0) } == 0;
let after_kill = linux_process_group_has_live_member(pgid);
// 断言前回收直属子进程,测试失败也不留下僵尸。
child.wait().unwrap();
assert_eq!(live, Some(true));
assert!(zombie && group_exists);
assert_eq!(after_kill, Some(false));
}
#[test]
fn group_retirement_allows_a_new_turn_without_claiming_full_tree_acceptance() {
let mut proof = ProcessTreeExitProof {
@@ -519,7 +519,12 @@ pub(super) async fn review_reply(
}
let ledger = session.snapshot()?;
if ledger.contract.is_none() {
session.finish_without_contract()?;
let finished = session.finish_without_contract();
// 收尾的预算检查可能刚写入 Exhausted;优先呈现已持久化的终态报告。
if let Some(report) = terminal_report(session) {
return Ok(Some(report));
}
finished?;
return Ok(None);
}
if try_seal(root, session).await? || session.snapshot()?.phase == ExecutionPhase::Sealing {
@@ -1167,7 +1167,7 @@ impl ExecutionSession {
next.phase = ExecutionPhase::Interrupted;
let prior = next.terminal_report.take().unwrap_or_default();
next.terminal_report = Some(format!(
"{prior}\n执行器缺少完整子进程退出证明。本轮仍未验收,需核对后台操作结果。"
"{prior}\n执行器缺少受控归属退出证明。本轮仍未验收,需核对后台操作结果。"
));
}
self.commit(&mut data, next)
@@ -1263,13 +1263,19 @@ impl ExecutionSession {
self.commit(&mut data, next)
}
/// 只把原因追加进终态说明,**不改阶段**。
/// 仅追加到失败终态说明,不为正常收尾生成终态报告,也不改写完成后的回复。
///
/// 宿主自己收尾(正常终态、预算与交付收尾)时会先关掉 app-server,连接随之关闭;
/// 这类「关闭原因」要留痕给排障看,但不能把已经/正在正常收口的回合改写成 `Interrupted`
/// —— CLI、单回合宿主每轮都会命中这个窗口(见 pitfalls 2026-09-28)。
pub(super) fn append_terminal_note(&self, reason: String) -> Result<(), String> {
let mut data = self.lock()?;
if !matches!(
data.ledger.phase,
ExecutionPhase::Interrupted | ExecutionPhase::Exhausted
) {
return Ok(());
}
let mut next = data.ledger.clone();
let prior = next.terminal_report.take().unwrap_or_default();
next.terminal_report = Some(
@@ -1,5 +1,36 @@
use super::*;
#[tokio::test]
async fn no_contract_review_returns_budget_report_created_during_close() {
let (_temp, session) = fixture_without_contract(DirectValidationConfig {
max_turn_seconds: 1,
..Default::default()
});
// 不先 tick:让 review_reply 进入无合同分支后,收尾本身触发预算终态。
session.lock().unwrap().elapsed_offset_ms = 1001;
assert!(session.snapshot().unwrap().terminal_report.is_none());
let report = super::super::direct_delivery::review_reply(&session.root, &session)
.await
.unwrap()
.expect("预算报告应优先于关闭错误");
let state = session.snapshot().unwrap();
assert_eq!(state.phase, ExecutionPhase::Exhausted);
assert!(report.contains("时间预算已耗尽"));
assert_eq!(state.terminal_report.as_deref(), Some(report.as_str()));
}
#[tokio::test]
async fn no_contract_review_preserves_close_error_without_terminal_report() {
let (_temp, session) = fixture_without_contract(Default::default());
let operation = session.admit(EffectKind::Execute, None).unwrap();
let result = super::super::direct_delivery::review_reply(&session.root, &session).await;
assert!(result.is_err());
let state = session.snapshot().unwrap();
assert_eq!(state.phase, ExecutionPhase::Working);
assert!(state.terminal_report.is_none());
operation.finish(true, false, None).unwrap();
}
fn analytics_metadata(user: &str) -> crate::analytics::run::Metadata {
use crate::analytics::contract::{Context, Route, RunSource, Source};
crate::analytics::run::Metadata::new(
@@ -182,6 +213,16 @@ fn legacy_run_without_metadata_is_not_assigned_current_users_identity() {
}
fn fixture(config: DirectValidationConfig) -> (tempfile::TempDir, Arc<ExecutionSession>) {
let (temp, session) = fixture_without_contract(config);
session
.freeze_contract(json!({"requirements":[{"id":"test"}]}))
.unwrap();
(temp, session)
}
fn fixture_without_contract(
config: DirectValidationConfig,
) -> (tempfile::TempDir, Arc<ExecutionSession>) {
let temp = tempfile::tempdir().unwrap();
let root = temp.path().join("project");
crate::init_local_game_project_at(&root, "execution-test", "执行测试").unwrap();
@@ -194,12 +235,60 @@ fn fixture(config: DirectValidationConfig) -> (tempfile::TempDir, Arc<ExecutionS
)
.unwrap();
assert!(session.newly_accepted);
session
.freeze_contract(json!({"requirements":[{"id":"test"}]}))
.unwrap();
(temp, session)
}
#[test]
fn shutdown_notes_do_not_create_or_overwrite_success_reports() {
let (_temp, ordinary) = fixture_without_contract(Default::default());
for closing in [false, true] {
if closing {
ordinary.begin_closing().unwrap();
}
ordinary.append_terminal_note("连接已关闭".into()).unwrap();
assert!(ordinary.snapshot().unwrap().terminal_report.is_none());
}
ordinary.record_process_exit_proof(true).unwrap();
ordinary.finish_attempt().unwrap();
ordinary.finish_without_contract().unwrap();
ordinary
.append_terminal_note("迟到关闭通知".into())
.unwrap();
assert_eq!(
ordinary.snapshot().unwrap().phase,
ExecutionPhase::Completed
);
assert!(ordinary.snapshot().unwrap().terminal_report.is_none());
let (_temp, delivery) = fixture(Default::default());
assert!(delivery
.begin_sealing(delivery.snapshot().unwrap().revision)
.unwrap());
delivery.append_terminal_note("连接已关闭".into()).unwrap();
assert!(delivery.snapshot().unwrap().terminal_report.is_none());
delivery.record_process_exit_proof(true).unwrap();
delivery.complete("验收完成".into()).unwrap();
delivery
.append_terminal_note("迟到关闭通知".into())
.unwrap();
assert_eq!(
delivery.snapshot().unwrap().terminal_report.as_deref(),
Some("验收完成")
);
let (_temp, failed) = fixture(Default::default());
failed.interrupt("用户取消".into()).unwrap();
failed.append_terminal_note("连接已关闭".into()).unwrap();
assert_eq!(
failed.snapshot().unwrap().phase,
ExecutionPhase::Interrupted
);
assert_eq!(
failed.snapshot().unwrap().terminal_report.as_deref(),
Some("用户取消\n连接已关闭")
);
}
#[test]
fn repeated_failures_do_not_consume_delivery_reviews_or_block_other_operations() {
let (_temp, session) = fixture(DirectValidationConfig {