15c227cebf
完善Codex CLI、App Server、ConPTY与进程会话在Windows下的发现、启动、恢复和退出行为 修复Agent Runtime、Provider重试、项目写锁及工具交接账本的并发与跨测试串线问题 补齐配置目录、路径脱敏、原子写入、浏览器探测和本地Provider smoke的跨平台兼容 增强Goal Contract、自动策略、资源生成及运行态恢复的契约和回归测试 更新AI游戏创作智能体App技术文档中的Windows稳定性说明 验证AGC开发态、Release打包、打包后GUI运行及完整agc:check门禁
222 lines
7.9 KiB
Rust
222 lines
7.9 KiB
Rust
use super::*;
|
||
|
||
pub(crate) fn active_process_session_records_at(
|
||
root: &Path,
|
||
agent_id: Option<&str>,
|
||
run_id: Option<&str>,
|
||
) -> Result<Vec<ProcessSessionRecord>, String> {
|
||
let live_sessions = process_session_registry()
|
||
.lock()
|
||
.map_err(|_| "process session registry 锁已损坏".to_string())?
|
||
.sessions
|
||
.values()
|
||
.filter(|live| live.root == root)
|
||
.cloned()
|
||
.collect::<Vec<_>>();
|
||
let mut records = Vec::with_capacity(live_sessions.len());
|
||
for live in live_sessions {
|
||
if agent_id.is_some_and(|value| value != live.identity.agent_id.as_str())
|
||
|| run_id.is_some_and(|value| value != live.identity.run_id.as_str())
|
||
{
|
||
continue;
|
||
}
|
||
let output = live
|
||
.output
|
||
.lock()
|
||
.map_err(|_| "process session output 锁已损坏".to_string())?;
|
||
records.push(process_session_record_from_live(&live, &output));
|
||
}
|
||
let directory = root.join(".agent/runtime/process-sessions");
|
||
let entries = match fs::read_dir(&directory) {
|
||
Ok(entries) => entries,
|
||
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
|
||
records.sort_by(|left, right| left.started_at.cmp(&right.started_at));
|
||
return Ok(records);
|
||
}
|
||
Err(error) => return Err(format!("读取 process session 目录失败:{error}")),
|
||
};
|
||
for entry in entries {
|
||
let entry = entry.map_err(|error| format!("读取 process session 目录项失败:{error}"))?;
|
||
let name = entry.file_name();
|
||
let Some(name) = name.to_str() else {
|
||
continue;
|
||
};
|
||
let Some(process_id) = name
|
||
.strip_suffix(".json")
|
||
.filter(|value| !value.ends_with(".output"))
|
||
else {
|
||
continue;
|
||
};
|
||
if validate_process_id(process_id).is_err() {
|
||
continue;
|
||
}
|
||
if records.iter().any(|record| record.process_id == process_id) {
|
||
continue;
|
||
}
|
||
let Some(mut record) = read_process_session_record(root, process_id)? else {
|
||
continue;
|
||
};
|
||
if matches!(
|
||
record.status.as_str(),
|
||
"prepared" | "launching" | "running" | "terminating"
|
||
) && record.owner_boot_id != process_session_boot_id()
|
||
{
|
||
reconcile_stale_active_process_session(&mut record);
|
||
write_process_session_record(root, &record)?;
|
||
}
|
||
if (record.needs_reconciliation
|
||
|| matches!(
|
||
record.status.as_str(),
|
||
"prepared" | "launching" | "running" | "terminating" | "needs-reconciliation"
|
||
))
|
||
&& agent_id.is_none_or(|value| value == record.agent_id)
|
||
&& run_id.is_none_or(|value| value == record.run_id)
|
||
{
|
||
records.push(record);
|
||
}
|
||
}
|
||
records.sort_by(|left, right| left.started_at.cmp(&right.started_at));
|
||
Ok(records)
|
||
}
|
||
|
||
pub(crate) fn has_active_process_sessions_at(root: &Path) -> Result<bool, String> {
|
||
if !active_process_session_records_at(root, None, None)?.is_empty() {
|
||
return Ok(true);
|
||
}
|
||
#[cfg(target_os = "linux")]
|
||
{
|
||
return Ok(pending_process_launch_registry()
|
||
.lock()
|
||
.map_err(|_| "pending process launch registry 锁已损坏".to_string())?
|
||
.launches
|
||
.values()
|
||
.any(|launch| launch.root == root));
|
||
}
|
||
#[cfg(not(target_os = "linux"))]
|
||
Ok(false)
|
||
}
|
||
|
||
pub(crate) fn terminate_process_sessions_for_run_at(
|
||
root: &Path,
|
||
agent_id: &str,
|
||
run_id: &str,
|
||
) -> Result<(), String> {
|
||
let records = active_process_session_records_at(root, Some(agent_id), Some(run_id))?;
|
||
for record in &records {
|
||
if record.status == "needs-reconciliation" || record.needs_reconciliation {
|
||
return Err(format!(
|
||
"进程会话 {} 需要人工核对,不能把 run 标记为已取消",
|
||
record.process_id
|
||
));
|
||
}
|
||
let identity = ProcessSessionIdentity {
|
||
project_id: record.project_id.clone(),
|
||
agent_id: record.agent_id.clone(),
|
||
task_id: record.task_id.clone(),
|
||
conversation_session_id: record.conversation_session_id.clone(),
|
||
run_id: record.run_id.clone(),
|
||
start_action_id: record.start_action_id.clone(),
|
||
start_action_fingerprint: record.start_action_fingerprint.clone(),
|
||
};
|
||
let terminal = terminate_process_session_at(root, &identity, &record.process_id, None)?;
|
||
if terminal.status == "running" || terminal.needs_reconciliation {
|
||
return Err(format!(
|
||
"进程会话 {} 尚未形成可信终态(status={},needsReconciliation={}),不能把 run 标记为已取消",
|
||
record.process_id, terminal.status, terminal.needs_reconciliation
|
||
));
|
||
}
|
||
}
|
||
if active_process_session_records_at(root, Some(agent_id), Some(run_id))?.is_empty() {
|
||
Ok(())
|
||
} else {
|
||
Err("仍有未收束的 process session,不能把 run 标记为已取消".to_string())
|
||
}
|
||
}
|
||
|
||
pub(crate) fn shutdown_all_process_sessions() {
|
||
#[cfg(target_os = "linux")]
|
||
{
|
||
let pending = pending_process_launch_registry()
|
||
.lock()
|
||
.ok()
|
||
.map(|mut registry| {
|
||
registry
|
||
.launches
|
||
.values_mut()
|
||
.filter_map(|launch| {
|
||
launch.shutdown_requested = true;
|
||
launch.process_group_leader
|
||
})
|
||
.collect::<Vec<_>>()
|
||
})
|
||
.unwrap_or_default();
|
||
for process_group_leader in pending {
|
||
unsafe {
|
||
libc::kill(-process_group_leader, libc::SIGKILL);
|
||
}
|
||
}
|
||
}
|
||
let sessions = process_session_registry()
|
||
.lock()
|
||
.ok()
|
||
.map(|registry| registry.sessions.values().cloned().collect::<Vec<_>>())
|
||
.unwrap_or_default();
|
||
for live in sessions {
|
||
let _ = live.control.send(ProcessControl::Shutdown);
|
||
}
|
||
}
|
||
|
||
pub(crate) fn shutdown_all_process_sessions_and_wait(timeout: Duration) -> Result<(), String> {
|
||
shutdown_all_process_sessions();
|
||
let deadline = std::time::Instant::now() + timeout;
|
||
loop {
|
||
let sessions = process_session_registry()
|
||
.lock()
|
||
.map_err(|_| "process session registry 锁已损坏".to_string())?
|
||
.sessions
|
||
.values()
|
||
.cloned()
|
||
.collect::<Vec<_>>();
|
||
let running = sessions.iter().filter(|live| {
|
||
live.output
|
||
.lock()
|
||
.map(|output| {
|
||
matches!(
|
||
output.status.as_str(),
|
||
"prepared" | "launching" | "running" | "terminating"
|
||
)
|
||
})
|
||
.unwrap_or(true)
|
||
});
|
||
if running.count() == 0 {
|
||
#[cfg(target_os = "linux")]
|
||
let pending_empty = pending_process_launch_registry()
|
||
.lock()
|
||
.map_err(|_| "pending process launch registry 锁已损坏".to_string())?
|
||
.launches
|
||
.is_empty();
|
||
#[cfg(not(target_os = "linux"))]
|
||
let pending_empty = true;
|
||
if pending_empty {
|
||
return Ok(());
|
||
}
|
||
}
|
||
if std::time::Instant::now() >= deadline {
|
||
return Err("Runner 退出前未能回收全部 process session".to_string());
|
||
}
|
||
thread::sleep(Duration::from_millis(25));
|
||
}
|
||
}
|
||
|
||
#[cfg(test)]
|
||
pub(crate) fn clear_process_session_registry_for_tests() {
|
||
shutdown_all_process_sessions();
|
||
if let Ok(mut registry) = process_session_registry().lock() {
|
||
registry.sessions.clear();
|
||
}
|
||
#[cfg(target_os = "linux")]
|
||
if let Ok(mut registry) = pending_process_launch_registry().lock() {
|
||
registry.launches.clear();
|
||
}
|
||
}
|