合并 lane B:清理 AGC 客户端退役运行时(session record / 内容 diff / ledger 写入层)遗留告警

- 并入本地分支 chore/agc-warnings-session @ 4e8662e87(18 文件):process_session 退役模块删除、project/* 与 runner/tests.rs 未用项收敛。
This commit is contained in:
2026-10-06 17:02:11 +08:00
18 changed files with 56 additions and 6302 deletions
@@ -1,22 +1,13 @@
use super::*;
use portable_pty::{native_pty_system, Child, CommandBuilder, MasterPty, PtySize};
use sha2::{Digest, Sha256};
use std::collections::HashMap;
use std::sync::Condvar;
#[cfg(target_os = "linux")]
use crate::process_session_bridge::*;
mod io;
mod lifecycle;
mod model;
mod persistence;
mod recovery;
pub(crate) use io::*;
pub(crate) use lifecycle::*;
pub(crate) use model::*;
pub(crate) use persistence::*;
pub(crate) use recovery::*;
#[cfg(test)]
@@ -1,370 +0,0 @@
use super::*;
pub(super) fn live_process_session(
process_id: &str,
) -> Result<Option<Arc<LiveProcessSession>>, String> {
validate_process_id(process_id)?;
Ok(process_session_registry()
.lock()
.map_err(|_| "process session registry 锁已损坏".to_string())?
.sessions
.get(process_id)
.cloned())
}
pub(super) fn poll_result_from_output(
process_id: &str,
output: &str,
state: &ProcessOutputState,
sandbox_backend: &str,
sandbox_mode: &str,
network_access: &str,
sandbox_profile_version: &str,
sandbox_establishment: &str,
target_exec: &str,
launch_failure_kind: Option<&str>,
cursor: Option<&str>,
max_chars: usize,
) -> Result<ProcessSessionPollResult, String> {
let offset = parse_process_session_cursor(process_id, cursor, output)?;
let end = output[offset..]
.char_indices()
.nth(max_chars)
.map(|(index, _)| offset + index)
.unwrap_or(output.len());
let next_cursor = process_session_cursor(process_id, end);
Ok(ProcessSessionPollResult {
process_id: process_id.to_string(),
status: state.status.clone(),
output: output[offset..end].to_string(),
cursor: process_session_cursor(process_id, offset),
next_cursor,
has_more: end < output.len(),
stdin_open: state.stdin_open,
exit_code: state.exit_code,
signal: state.signal.clone(),
output_bytes: output.len(),
output_sha256: format!("{:x}", Sha256::digest(output.as_bytes())),
source_changed: state.source_changed,
needs_reconciliation: state.needs_reconciliation,
sandbox_backend: sandbox_backend.to_string(),
sandbox_mode: sandbox_mode.to_string(),
network_access: network_access.to_string(),
sandbox_profile_version: sandbox_profile_version.to_string(),
sandbox_establishment: sandbox_establishment.to_string(),
target_exec: target_exec.to_string(),
launch_failure_kind: launch_failure_kind.map(str::to_string),
})
}
pub(crate) fn poll_process_session_at(
root: &Path,
identity: &ProcessSessionIdentity,
process_id: &str,
cursor: Option<&str>,
max_chars: Option<usize>,
wait_ms: Option<u64>,
) -> Result<ProcessSessionPollResult, String> {
validate_process_session_identity(identity)?;
let max_chars = max_chars
.unwrap_or(PROCESS_SESSION_DEFAULT_POLL_CHARS)
.min(PROCESS_SESSION_MAX_POLL_CHARS);
let wait_ms = wait_ms.unwrap_or(0).min(PROCESS_SESSION_MAX_POLL_WAIT_MS);
if let Some(live) = live_process_session(process_id)? {
if live.root != root || live.identity != *identity {
return Err("process session 不属于当前 Agent run".to_string());
}
let mut output = live
.output
.lock()
.map_err(|_| "process session output 锁已损坏".to_string())?;
let initial_offset = parse_process_session_cursor(process_id, cursor, &output.text)?;
if wait_ms > 0 && initial_offset == output.text.len() && output.status == "running" {
let waited = live
.output_changed
.wait_timeout(output, Duration::from_millis(wait_ms))
.map_err(|_| "process session output 锁已损坏".to_string())?;
output = waited.0;
}
return poll_result_from_output(
process_id,
&output.text,
&output,
&live.sandbox_backend,
&live.sandbox_mode,
&live.network_access,
&live.sandbox_profile_version,
&live.sandbox_establishment,
&live.target_exec,
output.launch_failure_kind.as_deref(),
cursor,
max_chars,
);
}
let mut record = read_process_session_record(root, process_id)?
.ok_or_else(|| "process session 不存在".to_string())?;
validate_process_session_access(&record, identity)?;
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)?;
}
let transcript = if let Some(output_ref) = record.output_ref.as_deref() {
read_agent_runtime_json_sidecar_with_max_bytes::<ProcessSessionTranscript>(
root,
output_ref,
"Agent Runtime process transcript",
PROCESS_SESSION_TRANSCRIPT_MAX_BYTES,
)?
} else {
None
};
if let Some(transcript) = &transcript {
if let Err(error) = validate_process_session_transcript(transcript, &record) {
record.status = "needs-reconciliation".to_string();
record.stdin_open = false;
record.needs_reconciliation = true;
record.terminal_at = Some(unix_timestamp());
record.updated_at = unix_timestamp();
let _ = write_process_session_record(root, &record);
return Err(error);
}
}
let output = transcript
.as_ref()
.map(|value| value.output.as_str())
.unwrap_or_default();
let state = ProcessOutputState {
text: output.to_string(),
status: record.status,
exit_code: record.exit_code,
signal: record.signal,
stdin_open: record.stdin_open,
reader_finished: true,
output_limit_exceeded: false,
source_fingerprint_after: record.source_fingerprint_after,
source_changed: record.source_changed,
needs_reconciliation: record.needs_reconciliation,
launch_failure_kind: record.launch_failure_kind.clone(),
};
poll_result_from_output(
process_id,
output,
&state,
&record.sandbox_backend,
&record.sandbox_mode,
&record.network_access,
&record.sandbox_profile_version,
&record.sandbox_establishment,
&record.target_exec,
record.launch_failure_kind.as_deref(),
cursor,
max_chars,
)
}
pub(crate) fn write_process_session_stdin_at(
root: &Path,
identity: &ProcessSessionIdentity,
process_id: &str,
data: &str,
append_newline: bool,
eof: bool,
) -> Result<ProcessSessionStdinResult, String> {
write_process_session_stdin_at_with_after_write(
root,
identity,
process_id,
data,
append_newline,
eof,
|_| {},
)
}
pub(super) fn write_process_session_stdin_at_with_after_write<F>(
root: &Path,
identity: &ProcessSessionIdentity,
process_id: &str,
data: &str,
append_newline: bool,
eof: bool,
after_write: F,
) -> Result<ProcessSessionStdinResult, String>
where
F: FnOnce(&LiveProcessSession),
{
validate_process_session_identity(identity)?;
let live = live_process_session(process_id)?
.ok_or_else(|| "process session 不在当前 Runner 中运行".to_string())?;
if live.root != root || live.identity != *identity {
return Err("process session 不属于当前 Agent run".to_string());
}
let mut bytes = data.as_bytes().to_vec();
if append_newline {
#[cfg(windows)]
bytes.extend_from_slice(b"\r\n");
#[cfg(not(windows))]
bytes.push(b'\n');
}
if bytes.len() > PROCESS_SESSION_MAX_STDIN_BYTES {
return Err(format!(
"command.stdin 单次最多写入 {PROCESS_SESSION_MAX_STDIN_BYTES} 字节"
));
}
if bytes.iter().any(|byte| *byte == 0) {
return Err("command.stdin 不接受 NUL 或二进制正文".to_string());
}
let content_sha256 = format!("{:x}", Sha256::digest(&bytes));
let mut writer = live
.writer
.lock()
.map_err(|_| "process session stdin 锁已损坏".to_string())?;
if live
.output
.lock()
.map_err(|_| "process session output 锁已损坏".to_string())?
.status
!= "running"
{
return Err("process session 已进入终态".to_string());
}
if !bytes.is_empty() {
let stream = writer
.as_mut()
.ok_or_else(|| "process session stdin 已关闭".to_string())?;
stream
.write_all(&bytes)
.and_then(|()| stream.flush())
.map_err(|error| format!("写入 process session stdin 失败:{error}"))?;
}
if eof {
writer.take();
}
drop(writer);
after_write(&live);
let mut output = live
.output
.lock()
.map_err(|_| "process session output 锁已损坏".to_string())?;
if eof || output.status != "running" {
output.stdin_open = false;
}
let record = process_session_record_from_live(&live, &output);
if let Err(error) = write_process_session_record(root, &record) {
output.status = "needs-reconciliation".to_string();
output.needs_reconciliation = true;
output.stdin_open = false;
let reconciliation = process_session_record_from_live(&live, &output);
let _ = write_process_session_record(root, &reconciliation);
let _ = live.control.send(ProcessControl::Terminate);
return Err(format!(
"command.stdin 已写入但状态无法落盘,需要人工核对:{error}"
));
}
Ok(ProcessSessionStdinResult {
process_id: process_id.to_string(),
bytes_written: data.len() + usize::from(append_newline),
content_sha256,
stdin_open: output.stdin_open,
eof,
sandbox_backend: live.sandbox_backend.clone(),
sandbox_mode: live.sandbox_mode.clone(),
network_access: live.network_access.clone(),
sandbox_profile_version: live.sandbox_profile_version.clone(),
})
}
pub(crate) fn terminate_process_session_at(
root: &Path,
identity: &ProcessSessionIdentity,
process_id: &str,
cursor: Option<&str>,
) -> Result<ProcessSessionPollResult, String> {
validate_process_session_identity(identity)?;
if let Some(live) = live_process_session(process_id)? {
if live.root != root || live.identity != *identity {
return Err("process session 不属于当前 Agent run".to_string());
}
let running = live
.output
.lock()
.map_err(|_| "process session output 锁已损坏".to_string())?
.status
== "running";
if running {
live.control
.send(ProcessControl::Terminate)
.map_err(|_| "process session 监督线程已结束".to_string())?;
let mut output = live
.output
.lock()
.map_err(|_| "process session output 锁已损坏".to_string())?;
let deadline = std::time::Instant::now()
+ Duration::from_millis(PROCESS_SESSION_TERMINATE_GRACE_MS + 1_500);
while output.status == "running" && std::time::Instant::now() < deadline {
let remaining = deadline.saturating_duration_since(std::time::Instant::now());
let waited = live
.output_changed
.wait_timeout(output, remaining.min(Duration::from_millis(100)))
.map_err(|_| "process session output 锁已损坏".to_string())?;
output = waited.0;
}
}
let mut result =
poll_process_session_at(root, identity, process_id, cursor, Some(1), Some(0))?;
let cursor_offset = result
.cursor
.rsplit_once(':')
.and_then(|(_, offset)| offset.parse::<usize>().ok())
.ok_or_else(|| "command.terminate 返回了无效 cursor".to_string())?;
result.output.clear();
result.next_cursor = result.cursor.clone();
result.has_more = cursor_offset < result.output_bytes;
return Ok(result);
}
let mut result = poll_process_session_at(root, identity, process_id, cursor, Some(1), Some(0))?;
let cursor_offset = result
.cursor
.rsplit_once(':')
.and_then(|(_, offset)| offset.parse::<usize>().ok())
.ok_or_else(|| "command.terminate 返回了无效 cursor".to_string())?;
result.output.clear();
result.next_cursor = result.cursor.clone();
result.has_more = cursor_offset < result.output_bytes;
Ok(result)
}
pub(crate) fn mark_process_session_start_audit_failure_at(
root: &Path,
process_id: &str,
error: &str,
) -> Result<(), String> {
if let Some(live) = live_process_session(process_id)? {
if live.root != root {
return Err("process session 不属于当前项目".to_string());
}
mark_process_session_reconciliation(
&live,
&format!("command.start audit persistence failed: {error}"),
);
if let Ok(mut output) = live.output.lock() {
output.launch_failure_kind = Some("start-audit-failed".to_string());
}
let _ = live.control.send(ProcessControl::Terminate);
}
let mut record = read_process_session_record(root, process_id)?
.ok_or_else(|| "command.start audit 失败后 process record 缺失".to_string())?;
record.status = "needs-reconciliation".to_string();
record.stdin_open = false;
record.needs_reconciliation = true;
record.launch_failure_kind = Some("start-audit-failed".to_string());
record.signal = Some(redact_agent_runtime_project_paths(root, error, 240));
record.terminal_at = Some(unix_timestamp());
record.updated_at = unix_timestamp();
write_process_session_record(root, &record)
}
File diff suppressed because it is too large Load Diff
@@ -1,205 +1,24 @@
use super::*;
pub(super) const PROCESS_SESSION_SCHEMA_VERSION: &str = "3";
pub(super) const PROCESS_SESSION_TRANSCRIPT_SCHEMA_VERSION: &str = "2";
pub(super) const PROCESS_SESSION_CURSOR_VERSION: &str = "v1";
pub(super) const PROCESS_SESSION_MAX_PER_PROJECT: usize = 4;
pub(super) const PROCESS_SESSION_MAX_PER_AGENT: usize = 2;
pub(super) const PROCESS_SESSION_MAX_OUTPUT_BYTES: usize = 256 * 1024;
pub(super) const PROCESS_SESSION_MAX_PENDING_LINE_BYTES: usize = 16 * 1024;
pub(super) const PROCESS_SESSION_MAX_STDIN_BYTES: usize = 8 * 1024;
pub(super) const PROCESS_SESSION_DEFAULT_POLL_CHARS: usize = 8_000;
pub(super) const PROCESS_SESSION_MAX_POLL_CHARS: usize = 16_000;
pub(super) const PROCESS_SESSION_MAX_POLL_WAIT_MS: u64 = 30_000;
pub(super) const PROCESS_SESSION_RECORD_MAX_BYTES: usize = 32 * 1024;
pub(super) const PROCESS_SESSION_TRANSCRIPT_MAX_BYTES: usize = 320 * 1024;
pub(super) const PROCESS_SESSION_TERMINATE_GRACE_MS: u64 = 800;
#[cfg(target_os = "linux")]
pub(super) const PROCESS_SESSION_OWNER_PID_ENV: &str = "GENARRATIVE_PROCESS_SESSION_OWNER_PID";
#[cfg(target_os = "linux")]
pub(super) const PROCESS_SESSION_CHILD_MODE: &str = "--process-session-child";
#[derive(Clone, Debug, Eq, PartialEq)]
pub(crate) struct ProcessSessionIdentity {
pub(crate) project_id: String,
pub(crate) agent_id: String,
pub(crate) task_id: String,
pub(crate) conversation_session_id: String,
pub(crate) run_id: String,
pub(crate) start_action_id: String,
pub(crate) start_action_fingerprint: String,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(deny_unknown_fields, rename_all = "camelCase")]
pub(crate) struct ProcessSessionRecord {
pub(crate) schema_version: String,
pub(crate) project_id: String,
pub(crate) agent_id: String,
pub(crate) task_id: String,
pub(crate) conversation_session_id: String,
pub(crate) run_id: String,
pub(crate) start_action_id: String,
pub(crate) start_action_fingerprint: String,
pub(crate) process_id: String,
pub(crate) owner_boot_id: String,
pub(crate) command_id: String,
pub(crate) program: String,
pub(crate) cwd: String,
#[serde(default)]
pub(crate) sandbox_backend: String,
#[serde(default)]
pub(crate) sandbox_mode: String,
#[serde(default)]
pub(crate) network_access: String,
#[serde(default)]
pub(crate) sandbox_profile_version: String,
#[serde(default)]
pub(crate) sandbox_establishment: String,
#[serde(default)]
pub(crate) target_exec: String,
#[serde(default)]
pub(crate) launch_failure_kind: Option<String>,
#[serde(default)]
pub(crate) sandbox_ready_at: Option<u64>,
#[serde(default)]
pub(crate) exec_established_at: Option<u64>,
pub(crate) status: String,
pub(crate) exit_code: Option<i32>,
pub(crate) signal: Option<String>,
pub(crate) stdin_open: bool,
pub(crate) output_bytes: usize,
pub(crate) output_sha256: String,
pub(crate) output_ref: Option<String>,
pub(crate) source_fingerprint_before: String,
pub(crate) source_fingerprint_after: Option<String>,
pub(crate) source_changed: Option<bool>,
pub(crate) needs_reconciliation: bool,
pub(crate) started_at: u64,
pub(crate) terminal_at: Option<u64>,
pub(crate) updated_at: u64,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(deny_unknown_fields, rename_all = "camelCase")]
pub(super) struct ProcessSessionTranscript {
pub(super) schema_version: String,
pub(super) project_id: String,
pub(super) agent_id: String,
pub(super) task_id: String,
pub(super) conversation_session_id: String,
pub(super) run_id: String,
pub(super) start_action_id: String,
pub(super) start_action_fingerprint: String,
pub(super) process_id: String,
pub(super) output: String,
pub(super) output_sha256: String,
pub(super) output_bytes: usize,
pub(super) updated_at: u64,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize)]
#[serde(rename_all = "camelCase")]
pub(crate) struct ProcessSessionPollResult {
pub(crate) process_id: String,
pub(crate) status: String,
pub(crate) output: String,
pub(crate) cursor: String,
pub(crate) next_cursor: String,
pub(crate) has_more: bool,
pub(crate) stdin_open: bool,
pub(crate) exit_code: Option<i32>,
pub(crate) signal: Option<String>,
pub(crate) output_bytes: usize,
pub(crate) output_sha256: String,
pub(crate) source_changed: Option<bool>,
pub(crate) needs_reconciliation: bool,
pub(crate) sandbox_backend: String,
pub(crate) sandbox_mode: String,
pub(crate) network_access: String,
pub(crate) sandbox_profile_version: String,
pub(crate) sandbox_establishment: String,
pub(crate) target_exec: String,
pub(crate) launch_failure_kind: Option<String>,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize)]
#[serde(rename_all = "camelCase")]
pub(crate) struct ProcessSessionStdinResult {
pub(crate) process_id: String,
pub(crate) bytes_written: usize,
pub(crate) content_sha256: String,
pub(crate) stdin_open: bool,
pub(crate) eof: bool,
pub(crate) sandbox_backend: String,
pub(crate) sandbox_mode: String,
pub(crate) network_access: String,
pub(crate) sandbox_profile_version: String,
}
#[cfg(not(test))]
#[derive(Debug)]
pub(super) struct ProcessOutputState {
pub(super) text: String,
pub(super) status: String,
pub(super) exit_code: Option<i32>,
pub(super) signal: Option<String>,
pub(super) stdin_open: bool,
pub(super) reader_finished: bool,
pub(super) output_limit_exceeded: bool,
pub(super) source_fingerprint_after: Option<String>,
pub(super) source_changed: Option<bool>,
pub(super) needs_reconciliation: bool,
pub(super) launch_failure_kind: Option<String>,
}
impl ProcessOutputState {
pub(super) fn running() -> Self {
Self {
text: String::new(),
status: "running".to_string(),
exit_code: None,
signal: None,
stdin_open: true,
reader_finished: false,
output_limit_exceeded: false,
source_fingerprint_after: None,
source_changed: None,
needs_reconciliation: false,
launch_failure_kind: None,
}
}
}
#[derive(Debug)]
pub(super) enum ProcessControl {
Terminate,
OutputLimit,
Shutdown,
}
pub(super) struct LiveProcessSession {
pub(super) root: PathBuf,
pub(super) identity: ProcessSessionIdentity,
pub(super) process_id: String,
pub(super) command_id: String,
pub(super) program: String,
pub(super) cwd: String,
pub(super) sandbox_backend: String,
pub(super) sandbox_mode: String,
pub(super) network_access: String,
pub(super) sandbox_profile_version: String,
pub(super) sandbox_establishment: String,
pub(super) target_exec: String,
pub(super) sandbox_ready_at: Option<u64>,
pub(super) exec_established_at: Option<u64>,
pub(super) source_fingerprint_before: String,
pub(super) started_at: u64,
#[cfg(not(test))]
pub(super) output: Mutex<ProcessOutputState>,
pub(super) output_changed: Condvar,
pub(super) writer: Mutex<Option<Box<dyn std::io::Write + Send>>>,
pub(super) master: Mutex<Option<Box<dyn MasterPty + Send>>>,
#[cfg(windows)]
pub(super) job: Mutex<Option<WindowsProcessJob>>,
pub(super) control: std::sync::mpsc::Sender<ProcessControl>,
}
@@ -215,14 +34,6 @@ unsafe impl Sync for WindowsProcessJob {}
#[cfg(windows)]
impl WindowsProcessJob {
pub(super) fn assign(child: &dyn Child) -> Result<Self, String> {
let process = child
.as_raw_handle()
.ok_or_else(|| "command.start Windows child 缺少 process handle".to_string())?
as windows_sys::Win32::Foundation::HANDLE;
Self::assign_handle(process)
}
pub(crate) fn assign_std(child: &std::process::Child) -> Result<Self, String> {
use std::os::windows::io::AsRawHandle;
Self::assign_handle(
@@ -491,17 +302,6 @@ impl Drop for WindowsProcessJob {
}
}
impl std::fmt::Debug for LiveProcessSession {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("LiveProcessSession")
.field("process_id", &self.process_id)
.field("agent_id", &self.identity.agent_id)
.field("run_id", &self.identity.run_id)
.finish_non_exhaustive()
}
}
#[derive(Default)]
pub(super) struct ProcessSessionRegistry {
pub(super) sessions: HashMap<String, Arc<LiveProcessSession>>,
@@ -510,8 +310,6 @@ pub(super) struct ProcessSessionRegistry {
#[cfg(target_os = "linux")]
#[derive(Clone, Debug)]
pub(super) struct PendingProcessLaunch {
pub(super) root: PathBuf,
pub(super) agent_id: String,
pub(super) process_group_leader: Option<i32>,
pub(super) shutdown_requested: bool,
}
@@ -522,21 +320,8 @@ pub(super) struct PendingProcessLaunchRegistry {
pub(super) launches: HashMap<String, PendingProcessLaunch>,
}
#[cfg(target_os = "linux")]
pub(super) struct PendingProcessLaunchGuard {
process_id: String,
}
#[cfg(target_os = "linux")]
impl Drop for PendingProcessLaunchGuard {
fn drop(&mut self) {
if let Ok(mut registry) = pending_process_launch_registry().lock() {
registry.launches.remove(&self.process_id);
}
}
}
static PROCESS_SESSION_REGISTRY: OnceLock<Mutex<ProcessSessionRegistry>> = OnceLock::new();
#[cfg(not(test))]
static PROCESS_SESSION_BOOT_ID: OnceLock<String> = OnceLock::new();
#[cfg(target_os = "linux")]
static PENDING_PROCESS_LAUNCH_REGISTRY: OnceLock<Mutex<PendingProcessLaunchRegistry>> =
@@ -552,77 +337,7 @@ pub(super) fn pending_process_launch_registry() -> &'static Mutex<PendingProcess
.get_or_init(|| Mutex::new(PendingProcessLaunchRegistry::default()))
}
#[cfg(target_os = "linux")]
pub(super) fn reserve_pending_process_launch(
root: &Path,
agent_id: &str,
process_id: &str,
) -> Result<PendingProcessLaunchGuard, String> {
let mut registry = pending_process_launch_registry()
.lock()
.map_err(|_| "pending process launch registry 锁已损坏".to_string())?;
if registry.launches.contains_key(process_id) {
return Err("command.start pending launch 身份冲突".to_string());
}
registry.launches.insert(
process_id.to_string(),
PendingProcessLaunch {
root: root.to_path_buf(),
agent_id: agent_id.to_string(),
process_group_leader: None,
shutdown_requested: false,
},
);
Ok(PendingProcessLaunchGuard {
process_id: process_id.to_string(),
})
}
#[cfg(target_os = "linux")]
pub(super) fn activate_pending_process_launch(
process_id: &str,
process_group_leader: i32,
) -> Result<(), String> {
if process_group_leader <= 1 {
return Err("command.start wrapper 进程组身份无效".to_string());
}
let mut registry = pending_process_launch_registry()
.lock()
.map_err(|_| "pending process launch registry 锁已损坏".to_string())?;
let launch = registry
.launches
.get_mut(process_id)
.ok_or_else(|| "command.start pending launch reservation 缺失".to_string())?;
launch.process_group_leader = Some(process_group_leader);
if launch.shutdown_requested {
unsafe {
libc::kill(-process_group_leader, libc::SIGKILL);
}
return Err("Runner shutdown 已取消 pending process launch".to_string());
}
Ok(())
}
pub(crate) fn process_session_boot_id() -> &'static str {
PROCESS_SESSION_BOOT_ID
.get_or_init(|| {
let mut digest = Sha256::new();
digest.update(std::process::id().to_le_bytes());
digest.update(
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_nanos()
.to_le_bytes(),
);
digest.update(unix_timestamp().to_le_bytes());
let value = format!("{:x}", digest.finalize());
format!("boot-{}", &value[..32])
})
.as_str()
}
// Runner 启动时绑定 owner;共享的会话归属与清理逻辑仍参与测试。
// Runner 启动时绑定 owner;该入口只由 Runner 生产接入(runner/server/desktop.rs)调用。
#[cfg(not(test))]
pub(crate) fn initialize_process_session_boot_id(boot_id: &str) -> Result<(), String> {
let boot_id = boot_id.trim();
File diff suppressed because it is too large Load Diff
@@ -1,138 +1,5 @@
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")]
{
@@ -166,6 +33,8 @@ pub(crate) fn shutdown_all_process_sessions() {
}
}
/// Runner 退出等待只由生产接入(runner/server/desktop.rs)调用。
#[cfg(not(test))]
pub(crate) fn shutdown_all_process_sessions_and_wait(timeout: Duration) -> Result<(), String> {
shutdown_all_process_sessions();
let deadline = std::time::Instant::now() + timeout;
@@ -207,15 +76,3 @@ pub(crate) fn shutdown_all_process_sessions_and_wait(timeout: Duration) -> Resul
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();
}
}
File diff suppressed because it is too large Load Diff
@@ -1,127 +0,0 @@
use std::fs;
use std::path::Path;
use std::process::Child;
use std::thread;
use std::time::{Duration, Instant};
pub(super) fn project_processes(root: &Path) -> Vec<i32> {
let Ok(canonical_root) = fs::canonicalize(root) else {
return Vec::new();
};
fs::read_dir("/proc")
.into_iter()
.flatten()
.flatten()
.filter_map(|entry| {
let process_id = entry.file_name().to_string_lossy().parse::<i32>().ok()?;
if process_id <= 1 || process_id == std::process::id() as i32 {
return None;
}
let cwd = fs::read_link(entry.path().join("cwd")).ok()?;
(cwd == canonical_root).then_some(process_id)
})
.collect()
}
pub(super) struct OwnerFixtureCleanup<'a> {
pub(super) owner: &'a mut Child,
pub(super) root: &'a Path,
}
impl Drop for OwnerFixtureCleanup<'_> {
fn drop(&mut self) {
let _ = self.owner.kill();
let _ = self.owner.wait();
// 正常路径先验证子进程自行退出;这里只兜底作用域退出(包括 panic)后的残留。
// 项目目录由每条用例独占,不能按进程名清理其他用例或开发进程。
let deadline = Instant::now() + Duration::from_secs(5);
loop {
let remaining = project_processes(self.root);
if remaining.is_empty() {
return;
}
for process_id in &remaining {
unsafe {
libc::kill(*process_id, libc::SIGKILL);
}
}
if Instant::now() >= deadline {
// Drop 可能在 panic 展开期间执行,不能再次 panic。
use std::io::Write;
let _ = writeln!(
std::io::stderr(),
"owner fixture cleanup timed out: pids={remaining:?}"
);
return;
}
thread::sleep(Duration::from_millis(25));
}
}
}
#[test]
fn owner_fixture_cleanup_reaps_processes_on_panic_without_touching_other_projects() {
use std::panic::{catch_unwind, AssertUnwindSafe};
use std::process::{Command, Stdio};
// 回归夹具自己的回收不能依赖被测 guard,否则 guard 回归时测试也会泄漏。
struct Sleeper(Child);
impl Drop for Sleeper {
fn drop(&mut self) {
let _ = self.0.kill();
let _ = self.0.wait();
}
}
fn sleeper(root: &Path) -> Sleeper {
Sleeper(
Command::new("sleep")
.arg("60")
.current_dir(root)
.stdin(Stdio::null())
.stdout(Stdio::null())
.stderr(Stdio::null())
.spawn()
.expect("spawn cleanup fixture"),
)
}
// 同时覆盖 owner 刚启动就失败,以及已有残留进程时失败。
for has_residual in [false, true] {
let project = tempfile::tempdir().expect("cleanup project");
let other_project = tempfile::tempdir().expect("unrelated project");
let mut owner = sleeper(project.path());
let cleanup = OwnerFixtureCleanup {
owner: &mut owner.0,
root: project.path(),
};
// 故意不依赖 owner 退出监测,验证兜底能清理仍留在项目目录的进程。
let mut residual = has_residual.then(|| sleeper(project.path()));
let mut other = sleeper(other_project.path());
let result = catch_unwind(AssertUnwindSafe(move || {
let _cleanup = cleanup;
panic!("simulate an assertion failure before owner shutdown");
}));
let owner_status = owner.0.try_wait();
let residual_status = residual.as_mut().map(|child| child.0.try_wait());
let other_status = other.0.try_wait();
// 即使 guard 回归,先收口本测试持有的进程再断言,避免回归用例自身泄漏。
drop(owner);
drop(residual);
drop(other);
assert!(result.is_err());
assert!(matches!(owner_status, Ok(Some(_))), "owner must exit");
if has_residual {
assert!(
matches!(residual_status, Some(Ok(Some(_)))),
"residual process must exit"
);
}
assert!(
matches!(other_status, Ok(None)),
"other project must survive"
);
}
}
@@ -1,10 +1,10 @@
use super::*;
use sha2::{Digest, Sha256};
use similar::TextDiff;
use std::io::{Seek, SeekFrom};
mod agent_db;
mod asset_export;
#[cfg(not(test))]
mod asset_rename;
mod bootstrap;
mod checkpoint;
@@ -19,12 +19,14 @@ mod resource_dependency_graph;
mod resource_editor;
mod resource_layout;
mod verification;
#[cfg(not(test))]
mod version_resource_replacement;
mod write_lock;
pub(crate) use agent_db::*;
#[cfg(not(test))]
pub(crate) use asset_export::*;
#[cfg(not(test))]
pub(crate) use asset_rename::*;
pub(crate) use bootstrap::*;
pub(crate) use checkpoint::*;
@@ -40,5 +42,6 @@ pub(crate) use resource_dependency_graph::*;
pub(crate) use resource_editor::*;
pub(crate) use resource_layout::*;
pub(crate) use verification::*;
#[cfg(not(test))]
pub(crate) use version_resource_replacement::*;
pub(crate) use write_lock::*;
File diff suppressed because it is too large Load Diff
@@ -28,26 +28,6 @@ pub(crate) struct RenameLocalProjectAssetResult {
pub(crate) committed_project_revision: u64,
}
/// 测试注入:`.agent/runtime/test-fail-next-asset-rename-manifest-write` 存在时,下一步
/// manifest 写入按失败返回,用来验证"文件已改名必须改回原名"的回滚。
///
/// 与 `agent_db` 的 `test-fail-next-*` 同一套约定,且整段是 `#[cfg(test)]`:生产签名
/// ([`rename_local_project_asset_at`])不接受任何故障注入参数,也没有第二个入口能把函数
/// 推上"只回滚、绝不写 manifest"的那条路。
#[cfg(test)]
fn take_rename_manifest_write_failure_injection(root: &Path) -> Result<(), String> {
let failure_path = root.join(".agent/runtime/test-fail-next-asset-rename-manifest-write");
match std::fs::read_to_string(&failure_path) {
Ok(_) => {
std::fs::remove_file(&failure_path)
.map_err(|error| format!("清理素材改名 manifest 写失败注入标记失败:{error}"))?;
Err("fault-injected:rename-asset-manifest-write".to_string())
}
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
Err(error) => Err(format!("读取素材改名 manifest 写失败注入标记失败:{error}")),
}
}
/// 校验调用方给出的新文件名,返回 trim 后的名字。
///
/// 判据(按顺序):
@@ -291,13 +271,7 @@ pub(crate) fn rename_local_project_asset_at(
let asset = manifest.assets[index].clone();
// 文件已改名:从这里开始的任何失败都必须把文件改回原名。
// 注入结果必须并进 `write_error`、不能就地 `?` 返回,否则会绕过下面的回滚。
#[cfg(test)]
let injected_write_error = take_rename_manifest_write_failure_injection(root).err();
#[cfg(not(test))]
let injected_write_error: Option<String> = None;
let write_error =
injected_write_error.or_else(|| write_manifest(&manifest_path, &manifest).err());
let write_error = write_manifest(&manifest_path, &manifest).err();
if let Some(error) = write_error {
return Err(rollback_asset_file_rename(
&current_absolute,
@@ -123,46 +123,6 @@ pub(crate) fn create_local_project_checkpoint_at(
})
}
const PROJECT_CONTENT_DIFF_CONTEXT_LINES: usize = 3;
const PROJECT_CONTENT_DIFF_MAX_FILE_BYTES: u64 = 2 * 1024 * 1024;
#[derive(Debug, Eq, PartialEq, Serialize)]
#[serde(rename_all = "camelCase")]
pub(crate) struct LocalProjectContentDiffResult {
pub(crate) checkpoint_id: String,
pub(crate) content: String,
pub(crate) file_count: usize,
pub(crate) truncated: bool,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum LocalProjectContentDiffStatus {
Added,
Changed,
Deleted,
}
impl LocalProjectContentDiffStatus {
fn as_str(self) -> &'static str {
match self {
Self::Added => "added",
Self::Changed => "changed",
Self::Deleted => "deleted",
}
}
}
struct LocalProjectContentDiffFile {
path: String,
status: LocalProjectContentDiffStatus,
}
struct LocalProjectContentDiffSource {
bytes: Option<Vec<u8>>,
size: u64,
sha256: String,
}
pub(crate) fn open_project_snapshot_regular_file(
path: &Path,
label: &str,
@@ -255,272 +215,6 @@ pub(crate) fn open_project_private_regular_file(
open_project_snapshot_regular_file(path, label)
}
fn read_local_project_content_diff_source(
path: &Path,
) -> Result<LocalProjectContentDiffSource, String> {
let (mut file, metadata) = open_project_private_regular_file(path, "内容 diff 文件")?;
let mut hasher = Sha256::new();
let mut bytes = (metadata.len() <= PROJECT_CONTENT_DIFF_MAX_FILE_BYTES)
.then(|| Vec::with_capacity(metadata.len() as usize));
let mut size = 0_u64;
let mut buffer = [0_u8; 64 * 1024];
loop {
let read = file
.read(&mut buffer)
.map_err(|error| format!("读取内容 diff 文件失败:{}: {error}", path.display()))?;
if read == 0 {
break;
}
hasher.update(&buffer[..read]);
size = size.saturating_add(read as u64);
if let Some(content) = bytes.as_mut() {
if size <= PROJECT_CONTENT_DIFF_MAX_FILE_BYTES {
content.extend_from_slice(&buffer[..read]);
} else {
bytes = None;
}
}
}
let final_metadata = file
.metadata()
.map_err(|error| format!("复核内容 diff 文件失败:{}: {error}", path.display()))?;
if final_metadata.len() != size {
return Err(format!(
"内容 diff 读取期间文件发生漂移:{}",
path.display()
));
}
#[cfg(unix)]
{
use std::os::unix::fs::MetadataExt;
if final_metadata.nlink() != 1 {
return Err(format!(
"内容 diff 文件不能是硬链接文件:{}",
path.display()
));
}
}
#[cfg(windows)]
validate_windows_regular_file_handle(&file, "内容 diff 文件")?;
Ok(LocalProjectContentDiffSource {
bytes,
size,
sha256: format!("{:x}", hasher.finalize()),
})
}
fn local_project_content_diff_source_at(
root: &Path,
checkpoint_id: &str,
file: &LocalProjectContentDiffFile,
checkpoint: bool,
) -> Result<Option<LocalProjectContentDiffSource>, String> {
if (checkpoint && file.status == LocalProjectContentDiffStatus::Added)
|| (!checkpoint && file.status == LocalProjectContentDiffStatus::Deleted)
{
return Ok(None);
}
let relative_path = if checkpoint {
checkpoint_file_relative_path(checkpoint_id, &file.path)
} else {
file.path.clone()
};
let path = resolve_local_project_path(root, &relative_path)?;
read_local_project_content_diff_source(&path).map(Some)
}
fn render_local_project_content_diff_section(
file: &LocalProjectContentDiffFile,
checkpoint: Option<&LocalProjectContentDiffSource>,
current: Option<&LocalProjectContentDiffSource>,
) -> (String, bool) {
let checkpoint_sha256 = checkpoint
.map(|source| source.sha256.as_str())
.unwrap_or("-");
let current_sha256 = current.map(|source| source.sha256.as_str()).unwrap_or("-");
let mut output = format!(
"diff --git a/{0} b/{0}\nstatus: {1}\ncheckpoint-sha256: {2}\ncurrent-sha256: {3}\n",
file.path,
file.status.as_str(),
checkpoint_sha256,
current_sha256,
);
if checkpoint.is_some_and(|source| source.bytes.is_none())
|| current.is_some_and(|source| source.bytes.is_none())
{
let checkpoint_size = checkpoint
.map(|source| source.size.to_string())
.unwrap_or_else(|| "-".to_string());
let current_size = current
.map(|source| source.size.to_string())
.unwrap_or_else(|| "-".to_string());
output.push_str(&format!(
"[content omitted: per-file limit exceeded; limit={PROJECT_CONTENT_DIFF_MAX_FILE_BYTES} bytes; checkpoint={checkpoint_size} bytes; current={current_size} bytes; truncated]\n\n"
));
return (output, true);
}
let checkpoint_bytes = checkpoint
.and_then(|source| source.bytes.as_deref())
.unwrap_or_default();
let current_bytes = current
.and_then(|source| source.bytes.as_deref())
.unwrap_or_default();
if checkpoint_bytes.contains(&0) || current_bytes.contains(&0) {
output.push_str("[content omitted: binary data detected; truncated]\n\n");
return (output, true);
}
let (Ok(checkpoint_text), Ok(current_text)) = (
std::str::from_utf8(checkpoint_bytes),
std::str::from_utf8(current_bytes),
) else {
output.push_str("[content omitted: non-UTF-8 data detected; truncated]\n\n");
return (output, true);
};
let diff = TextDiff::from_lines(checkpoint_text, current_text);
let checkpoint_header = if file.status == LocalProjectContentDiffStatus::Added {
"/dev/null".to_string()
} else {
format!("a/{}", file.path)
};
let current_header = if file.status == LocalProjectContentDiffStatus::Deleted {
"/dev/null".to_string()
} else {
format!("b/{}", file.path)
};
let mut unified = diff.unified_diff();
unified.context_radius(PROJECT_CONTENT_DIFF_CONTEXT_LINES);
unified.header(&checkpoint_header, &current_header);
output.push_str(&unified.to_string());
if !output.ends_with('\n') {
output.push('\n');
}
output.push('\n');
(output, false)
}
fn append_local_project_content_diff_truncation(
output: &mut String,
output_chars: &mut usize,
max_chars: usize,
detail: &str,
) {
let marker = format!("[content diff truncated: {detail}]\n");
let marker_chars = marker.chars().count();
if output_chars.saturating_add(marker_chars) <= max_chars {
output.push_str(&marker);
*output_chars += marker_chars;
} else if *output_chars < max_chars {
output.push('…');
*output_chars += 1;
}
}
pub(crate) fn diff_local_project_checkpoint_content_at(
root: &Path,
checkpoint_id: &str,
max_files: usize,
max_chars: usize,
) -> Result<LocalProjectContentDiffResult, String> {
let initial_diff = diff_local_project_checkpoint_at(root, checkpoint_id)?;
let checkpoint_id = initial_diff.checkpoint_id.clone();
let mut files = Vec::with_capacity(
initial_diff.added.len() + initial_diff.changed.len() + initial_diff.deleted.len(),
);
files.extend(
initial_diff
.added
.iter()
.map(|entry| LocalProjectContentDiffFile {
path: entry.path.clone(),
status: LocalProjectContentDiffStatus::Added,
}),
);
files.extend(
initial_diff
.changed
.iter()
.map(|entry| LocalProjectContentDiffFile {
path: entry.path.clone(),
status: LocalProjectContentDiffStatus::Changed,
}),
);
files.extend(
initial_diff
.deleted
.iter()
.map(|entry| LocalProjectContentDiffFile {
path: entry.path.clone(),
status: LocalProjectContentDiffStatus::Deleted,
}),
);
files.sort_by(|left, right| left.path.cmp(&right.path));
let mut content = String::new();
let mut content_chars = 0_usize;
let mut file_count = 0_usize;
let mut truncated = false;
let mut character_budget_exhausted = false;
let file_limit = max_files.min(files.len());
for (index, file) in files.iter().take(file_limit).enumerate() {
let checkpoint = local_project_content_diff_source_at(root, &checkpoint_id, file, true)?;
let current = local_project_content_diff_source_at(root, &checkpoint_id, file, false)?;
let (section, section_truncated) =
render_local_project_content_diff_section(file, checkpoint.as_ref(), current.as_ref());
let section_chars = section.chars().count();
let has_unprocessed_files = index + 1 < files.len();
let truncation_reserve = usize::from(has_unprocessed_files && max_chars > 0);
if content_chars
.saturating_add(section_chars)
.saturating_add(truncation_reserve)
> max_chars
{
truncated = true;
character_budget_exhausted = true;
append_local_project_content_diff_truncation(
&mut content,
&mut content_chars,
max_chars,
&format!(
"character budget reached; {} file diff(s) omitted",
files.len() - file_count
),
);
break;
}
content.push_str(&section);
content_chars += section_chars;
file_count += 1;
truncated |= section_truncated;
}
if !character_budget_exhausted && file_limit < files.len() {
truncated = true;
append_local_project_content_diff_truncation(
&mut content,
&mut content_chars,
max_chars,
&format!(
"file budget reached; {} file diff(s) omitted",
files.len() - file_count
),
);
}
let final_diff = diff_local_project_checkpoint_at(root, &checkpoint_id)?;
if final_diff != initial_diff {
return Err("内容 diff 读取期间项目文件发生变化,请重新执行 project.diff".to_string());
}
Ok(LocalProjectContentDiffResult {
checkpoint_id,
content,
file_count,
truncated,
})
}
pub(crate) fn diff_local_project_checkpoint_at(
root: &Path,
checkpoint_id: &str,
@@ -50,130 +50,9 @@ fn write_checkpoint_file(root: &Path, checkpoint_id: &str, path: &str, content:
fs::write(target, content).expect("write checkpoint file");
}
#[test]
fn checkpoint_content_diff_renders_added_changed_and_deleted_unified_hunks() {
let root = unique_checkpoint_test_root("content-diff-hunks");
fs::create_dir_all(root.join("game")).expect("create game directory");
fs::write(
root.join("game/changed.txt"),
"line-1\nline-2\nline-3\nline-4\nold-line\nline-6\nline-7\nline-8\n",
)
.expect("write changed fixture");
fs::write(root.join("game/deleted.txt"), "deleted\n").expect("write deleted fixture");
let checkpoint = create_local_project_checkpoint_at(&root).expect("create checkpoint");
fs::write(
root.join("game/changed.txt"),
"line-1\nline-2\nline-3\nline-4\nnew-line\nline-6\nline-7\nline-8\n",
)
.expect("update changed fixture");
fs::remove_file(root.join("game/deleted.txt")).expect("delete fixture");
fs::write(root.join("game/added.txt"), "added\n").expect("write added fixture");
let result =
diff_local_project_checkpoint_content_at(&root, &checkpoint.checkpoint_id, 10, 20_000)
.expect("render checkpoint content diff");
assert_eq!(result.checkpoint_id, checkpoint.checkpoint_id);
assert_eq!(result.file_count, 3);
assert!(!result.truncated);
assert!(result.content.contains("status: added"));
assert!(result.content.contains("status: changed"));
assert!(result.content.contains("status: deleted"));
assert!(result
.content
.contains("--- /dev/null\n+++ b/game/added.txt"));
assert!(result
.content
.contains("--- a/game/deleted.txt\n+++ /dev/null"));
assert!(result
.content
.contains(" line-2\n line-3\n line-4\n-old-line\n+new-line\n line-6\n line-7\n line-8"));
assert!(result.content.contains("checkpoint-sha256:"));
assert!(result.content.contains("current-sha256:"));
fs::remove_dir_all(root).ok();
}
#[test]
fn checkpoint_content_diff_marks_file_and_character_budget_truncation() {
let root = unique_checkpoint_test_root("content-diff-budgets");
fs::create_dir_all(root.join("game")).expect("create game directory");
fs::write(root.join("game/a.txt"), "before a\n").expect("write a fixture");
fs::write(root.join("game/b.txt"), "before b\n").expect("write b fixture");
let checkpoint = create_local_project_checkpoint_at(&root).expect("create checkpoint");
fs::write(root.join("game/a.txt"), "after a\n").expect("update a fixture");
fs::write(root.join("game/b.txt"), "after b\n").expect("update b fixture");
let file_limited =
diff_local_project_checkpoint_content_at(&root, &checkpoint.checkpoint_id, 1, 20_000)
.expect("render file-limited content diff");
assert_eq!(file_limited.file_count, 1);
assert!(file_limited.truncated);
assert!(file_limited.content.contains("game/a.txt"));
assert!(!file_limited.content.contains("game/b.txt"));
assert!(file_limited.content.contains("file budget reached"));
let character_limited =
diff_local_project_checkpoint_content_at(&root, &checkpoint.checkpoint_id, 10, 120)
.expect("render character-limited content diff");
assert_eq!(character_limited.file_count, 0);
assert!(character_limited.truncated);
assert!(character_limited
.content
.contains("character budget reached"));
assert!(character_limited.content.chars().count() <= 120);
let tiny_budget =
diff_local_project_checkpoint_content_at(&root, &checkpoint.checkpoint_id, 10, 1)
.expect("render tiny-budget content diff");
assert_eq!(tiny_budget.content, "…");
assert_eq!(tiny_budget.file_count, 0);
assert!(tiny_budget.truncated);
fs::remove_dir_all(root).ok();
}
#[test]
fn checkpoint_content_diff_marks_binary_non_utf8_and_oversized_files() {
let root = unique_checkpoint_test_root("content-diff-omissions");
fs::create_dir_all(root.join("game")).expect("create game directory");
fs::write(root.join("game/binary.bin"), [0_u8, 1, 2]).expect("write binary fixture");
fs::write(root.join("game/non-utf8.txt"), [0xff_u8, b'a']).expect("write non-UTF-8 fixture");
fs::write(
root.join("game/oversized.txt"),
vec![b'a'; PROJECT_CONTENT_DIFF_MAX_FILE_BYTES as usize + 1],
)
.expect("write oversized fixture");
let checkpoint = create_local_project_checkpoint_at(&root).expect("create checkpoint");
fs::write(root.join("game/binary.bin"), [0_u8, 1, 3]).expect("update binary fixture");
fs::write(root.join("game/non-utf8.txt"), [0xfe_u8, b'a']).expect("update non-UTF-8 fixture");
fs::write(
root.join("game/oversized.txt"),
vec![b'b'; PROJECT_CONTENT_DIFF_MAX_FILE_BYTES as usize + 1],
)
.expect("update oversized fixture");
let result =
diff_local_project_checkpoint_content_at(&root, &checkpoint.checkpoint_id, 10, 8_000)
.expect("render omitted content diff markers");
assert_eq!(result.file_count, 3);
assert!(result.truncated);
assert!(result.content.contains("binary data detected; truncated"));
assert!(result
.content
.contains("non-UTF-8 data detected; truncated"));
assert!(result.content.contains("per-file limit exceeded"));
assert!(result.content.contains("checkpoint-sha256:"));
assert!(result.content.contains("current-sha256:"));
fs::remove_dir_all(root).ok();
}
#[cfg(unix)]
#[test]
fn checkpoint_and_content_diff_reject_hard_linked_project_files() {
fn checkpoint_rejects_hard_linked_project_files() {
let root = unique_checkpoint_test_root("hard-link-boundary");
let outside = unique_checkpoint_test_root("hard-link-outside");
fs::create_dir_all(root.join("game")).expect("create game directory");
@@ -190,23 +69,6 @@ fn checkpoint_and_content_diff_reject_hard_linked_project_files() {
);
assert!(!root.join(".agent/checkpoints").exists());
fs::remove_file(root.join("game/linked.txt")).expect("remove first hard link");
fs::write(root.join("game/linked.txt"), "checkpoint version\n")
.expect("write safe checkpoint fixture");
let checkpoint = create_local_project_checkpoint_at(&root).expect("create safe checkpoint");
fs::remove_file(root.join("game/linked.txt")).expect("remove safe fixture");
fs::hard_link(&outside_file, root.join("game/linked.txt"))
.expect("replace current file with hard link");
let diff_error =
diff_local_project_checkpoint_content_at(&root, &checkpoint.checkpoint_id, 10, 10_000)
.expect_err("content diff must reject project hard links");
assert!(
diff_error.contains("hard link") || diff_error.contains("硬链接"),
"{diff_error}"
);
assert!(!diff_error.contains("outside secret"));
fs::remove_dir_all(root).ok();
fs::remove_dir_all(outside).ok();
}
@@ -551,18 +413,6 @@ fn snapshot_workflows_exclude_and_preserve_sensitive_paths() {
);
assert!(diff.deleted.is_empty());
let content_diff =
diff_local_project_checkpoint_content_at(&root, &checkpoint.checkpoint_id, 10, 20_000)
.expect("content diff ignores legacy sensitive checkpoint entries");
assert_eq!(content_diff.file_count, 2);
assert!(!content_diff.truncated);
assert!(content_diff.content.contains("game/extra.txt"));
assert!(content_diff.content.contains("game/state.txt"));
assert!(!content_diff.content.contains("checkpoint secret"));
for path in sensitive_paths {
assert!(!content_diff.content.contains(path));
}
let restored = restore_local_project_checkpoint_at(&root, &checkpoint.checkpoint_id)
.expect("restore ignores legacy sensitive checkpoint entries");
assert_eq!(restored.restored_count, 1);
@@ -629,13 +479,6 @@ fn checkpoint_create_and_restore_reject_symbolic_link_boundaries() {
.expect_err("checkpoint restore must reject a linked content source");
assert!(source_error.contains("符号链接"), "{source_error}");
assert!(!root.join("game/notes.txt").exists());
let content_diff_error =
diff_local_project_checkpoint_content_at(&root, source_checkpoint_id, 10, 20_000)
.expect_err("content diff must reject a linked checkpoint source");
assert!(
content_diff_error.contains("符号链接"),
"{content_diff_error}"
);
let target_checkpoint_id = "checkpoint-linked-target";
write_checkpoint_manifest(&root, target_checkpoint_id, &["linked/target.txt"]);
@@ -321,47 +321,7 @@ fn write_agent_conversation_session_catalog_unlocked(
)
}
pub(crate) fn ensure_agent_conversation_session_at(
root: &Path,
agent_id: &str,
session_id: &str,
title: &str,
) -> Result<(), String> {
validate_project_root(root)?;
let agent_id = normalize_conversation_agent_id(agent_id)?;
let session_id = normalize_agent_conversation_session_id(session_id)?;
with_agent_conversation_session_lane_at(root, &agent_id, "Agent Session 确保存在", || {
let catalog_path = agent_conversation_session_catalog_path(root, &agent_id);
let lock = project_append_lock_for(&catalog_path)?;
let _guard = lock.lock("Agent Session 目录")?;
let mut catalog = read_agent_conversation_session_catalog_unlocked(root, &agent_id)?;
if let Some(existing) = catalog
.sessions
.iter()
.find(|session| session.session_id == session_id)
{
if existing.archived_at.is_some() {
return Err(format!("Agent Session 已归档:{session_id}"));
}
} else {
let now = unix_timestamp();
catalog.sessions.push(AgentConversationSessionRecord {
session_id: session_id.clone(),
title: title.trim().chars().take(80).collect::<String>(),
created_at: now,
updated_at: now,
archived_at: None,
message_count: 0,
legacy: false,
forked_from_session_id: None,
forked_message_count: None,
});
}
catalog.active_session_id = session_id;
write_agent_conversation_session_catalog_unlocked(root, &catalog)
})
}
#[cfg(not(test))]
fn agent_conversation_session_list_result(
root: &Path,
catalog: AgentConversationSessionCatalogFile,
@@ -376,6 +336,7 @@ fn agent_conversation_session_list_result(
}
}
#[cfg(not(test))]
pub(crate) fn list_game_creator_agent_sessions_at(
root: &Path,
agent_id: &str,
@@ -388,6 +349,7 @@ pub(crate) fn list_game_creator_agent_sessions_at(
.map(|catalog| agent_conversation_session_list_result(root, catalog))
}
#[cfg(not(test))]
pub(crate) fn create_game_creator_agent_session_at(
root: &Path,
agent_id: &str,
@@ -452,6 +414,7 @@ pub(crate) fn create_game_creator_agent_session_at(
})
}
#[cfg(not(test))]
fn new_agent_conversation_session_id(
catalog: &AgentConversationSessionCatalogFile,
agent_id: &str,
@@ -475,6 +438,7 @@ fn new_agent_conversation_session_id(
.ok_or_else(|| "无法生成唯一 Agent Session ID".to_string())
}
#[cfg(not(test))]
pub(crate) fn fork_game_creator_agent_session_at(
root: &Path,
agent_id: &str,
@@ -490,6 +454,7 @@ pub(crate) fn fork_game_creator_agent_session_at(
)
}
#[cfg(not(test))]
fn fork_game_creator_agent_session_with_catalog_writer_at<F>(
root: &Path,
agent_id: &str,
@@ -608,45 +573,7 @@ where
})
}
#[cfg(test)]
pub(crate) fn fork_game_creator_agent_session_with_catalog_failure_at(
root: &Path,
agent_id: &str,
source_session_id: &str,
title: &str,
) -> Result<AgentConversationSessionListResult, String> {
fork_game_creator_agent_session_with_catalog_writer_at(
root,
agent_id,
source_session_id,
title,
|_, _| Err("测试注入 Agent Session 目录写入失败".to_string()),
)
}
#[cfg(test)]
pub(crate) fn fork_game_creator_agent_session_with_catalog_write_hook_at<F>(
root: &Path,
agent_id: &str,
source_session_id: &str,
title: &str,
before_catalog_write: F,
) -> Result<AgentConversationSessionListResult, String>
where
F: FnOnce(),
{
fork_game_creator_agent_session_with_catalog_writer_at(
root,
agent_id,
source_session_id,
title,
|root, catalog| {
before_catalog_write();
write_agent_conversation_session_catalog_unlocked(root, catalog)
},
)
}
#[cfg(not(test))]
pub(crate) fn set_active_game_creator_agent_session_at(
root: &Path,
agent_id: &str,
@@ -675,6 +602,7 @@ pub(crate) fn set_active_game_creator_agent_session_at(
})
}
#[cfg(not(test))]
pub(crate) fn archive_game_creator_agent_session_at(
root: &Path,
agent_id: &str,
@@ -729,24 +657,6 @@ fn resolve_agent_conversation_session_at(
.ok_or_else(|| format!("Agent Session 不存在:{session_id}"))
}
pub(crate) fn resolve_agent_conversation_session_id_at(
root: &Path,
agent_id: &str,
session_id: Option<&str>,
require_writable: bool,
) -> Result<String, String> {
validate_project_root(root)?;
let agent_id = normalize_conversation_agent_id(agent_id)?;
let session = resolve_agent_conversation_session_at(root, &agent_id, session_id)?;
if require_writable && session.archived_at.is_some() {
return Err(format!(
"Agent Session 已归档,只能读取:{}",
session.session_id
));
}
Ok(session.session_id)
}
fn touch_agent_conversation_session_at(
root: &Path,
agent_id: &str,
@@ -950,17 +860,6 @@ pub(crate) fn read_local_conversation_at(
read_local_conversation_for_session_at(root, agent_id, None)
}
pub(crate) fn read_local_conversation_message_by_id_for_session_at(
root: &Path,
agent_id: Option<&str>,
session_id: Option<&str>,
message_id: &str,
) -> Result<Option<LocalConversationMessageRecord>, String> {
read_local_conversation_message_by_id_for_session_internal_at(
root, agent_id, session_id, message_id, true,
)
}
pub(crate) fn read_local_conversation_message_by_id_for_session_without_touch_at(
root: &Path,
agent_id: Option<&str>,
@@ -1236,23 +1135,7 @@ pub(crate) fn append_local_conversation_message_for_session_idempotent_at(
.map(|(conversation, _)| conversation)
}
pub(crate) fn append_local_conversation_message_for_session_idempotent_with_status_at(
root: &Path,
agent_id: Option<&str>,
session_id: Option<&str>,
message: LocalConversationMessage,
message_id: &str,
) -> Result<(LocalConversationResult, bool), String> {
append_local_conversation_message_for_session_internal_at(
root,
agent_id,
session_id,
message,
Some(message_id),
None,
)
}
#[cfg(test)]
pub(crate) fn append_local_conversation_message_for_session_idempotent_with_finalization_at(
root: &Path,
agent_id: Option<&str>,
@@ -1354,6 +1237,7 @@ pub(crate) fn append_markdown_entry(
append_game_creator_private_file(path, &bytes, error_label)
}
#[cfg(not(test))]
pub(crate) fn append_local_permission_log_at(
root: &Path,
event: &str,
@@ -218,6 +218,8 @@ pub(crate) fn reject_agent_runtime_private_control_path(
Ok(())
}
/// file.delete 入口退役后,只剩测试通过它验证 `.agent` 控制面删除门禁。
#[cfg(test)]
fn reject_agent_control_path_delete(normalized_path: &str) -> Result<(), String> {
if matches!(
normalized_path.split('/').next(),
@@ -267,6 +269,8 @@ pub(crate) fn write_local_project_file_at(
})
}
/// file.delete 入口已退役:仅测试仍用它验证 `.agent`/checkpoint 控制面与硬链接门禁。
#[cfg(test)]
pub(crate) fn delete_local_project_file_at(
root: &Path,
relative_path: &str,
@@ -825,6 +825,7 @@ pub(crate) fn import_local_cocos_project_at(
})
}
#[cfg(not(test))]
pub(crate) fn import_local_unity_project_at(
root: &Path,
project_id: &str,
@@ -980,6 +981,7 @@ pub(crate) fn game_iteration_resource_bindings(
/// revision has produced a durable successful browser-playtest receipt.
/// Replays are idempotent: once any formal version exists, validation never
/// rewrites or appends another initial record.
#[cfg(test)]
pub(crate) fn ensure_initial_game_iteration_version_at(
root: &Path,
project_revision: u64,
@@ -1159,160 +1161,6 @@ pub(crate) fn ensure_manifest_seed_tasks(manifest: &mut GameCreationAppManifest)
}
}
pub(crate) fn validate_manifest_required_visual_asset(
root: &Path,
manifest: &GameCreationAppManifest,
task_id: &str,
) -> Result<(), String> {
let (expected_path, expected_kind) = match task_id {
"art-director" => ("assets/art-spec.png", GameCreationAppAssetKind::IconSpec),
"design-foundation" => (
"assets/ui-prototype.png",
GameCreationAppAssetKind::UiDesign,
),
"art-asset-plan" => (
"assets/art-spritesheet.png",
GameCreationAppAssetKind::IconSpritesheet,
),
_ => return Ok(()),
};
let asset = manifest
.assets
.iter()
.find(|asset| asset.local_path == expected_path && asset.kind == expected_kind)
.ok_or_else(|| format!("缺少规范视觉资产:{expected_path} ({expected_kind})"))?;
if !asset.media_type.starts_with("image/")
|| asset.source.kind != GameCreationAppAssetSourceKind::Canvas
{
return Err(format!("规范视觉资产文件或来源无效:{expected_path}"));
}
let bytes = resolve_local_project_path(root, &asset.local_path)
.ok()
.and_then(|path| fs::read(path).ok())
.filter(|bytes| bytes.starts_with(b"\x89PNG\r\n\x1a\n"))
.ok_or_else(|| format!("规范视觉资产不是有效登记的 PNG 文件:{expected_path}"))?;
let decoded = image::load_from_memory(&bytes)
.map_err(|_| format!("规范视觉资产 PNG 无法完整解码:{expected_path}"))?;
if decoded.width() == 0 || decoded.height() == 0 {
return Err(format!("规范视觉资产 PNG 尺寸无效:{expected_path}"));
}
if task_id == "art-asset-plan" && !decoded.to_rgba8().pixels().any(|pixel| pixel[3] < u8::MAX) {
return Err("首版美术素材图没有真实透明像素".to_string());
}
let canvas_project_id = asset
.source
.canvas_project_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.ok_or_else(|| format!("规范视觉资产缺少 canvasProjectId:{expected_path}"))?;
asset
.source
.resource_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.ok_or_else(|| format!("规范视觉资产缺少 resourceId:{expected_path}"))?;
let expected_route = match task_id {
"art-asset-plan" => "/api/external/v1/editor/icon-spritesheets/generations",
_ => "/api/external/v1/editor/images/generations",
};
let expected_generation_kind = match task_id {
"art-director" => "spec",
"design-foundation" => "ui-design",
"art-asset-plan" => "icon-spritesheet",
_ => unreachable!("non-visual tasks returned above"),
};
if asset.source.generation_route.as_deref() != Some(expected_route)
|| asset.source.generation_kind.as_deref() != Some(expected_generation_kind)
{
return Err(format!(
"规范视觉资产缺少匹配的持久生成来源,按 legacy 资产处理:{expected_path}"
));
}
if task_id == "art-director" {
// 规范图是视觉来源链的根:它自身不派生任何视觉资产,但 icon-spec 生成允许用户参考
// (没有规范前置,最多总上限),这些参考只是风格输入,不构成派生关系。这里改为验证
// 参考集合仍符合 icon-spec 请求合同;route / generation kind / canvasProjectId /
// resourceId / PNG 解码等身份判据全部保持不变。
if !crate::agent::platform_art_runtime_references_match_request_contract(
&asset.source.reference_resource_ids,
expected_kind,
) {
return Err(format!(
"统一视觉规范图的参考集合不符合请求合同:{expected_path}"
));
}
return Ok(());
}
validate_manifest_required_visual_asset(root, manifest, "art-director")?;
let art_spec = manifest
.assets
.iter()
.find(|asset| {
asset.local_path == "assets/art-spec.png"
&& asset.kind == GameCreationAppAssetKind::IconSpec
})
.ok_or_else(|| "缺少当前统一视觉规范图".to_string())?;
let art_spec_project_id = art_spec
.source
.canvas_project_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.ok_or_else(|| "统一视觉规范图缺少 canvasProjectId".to_string())?;
let art_spec_resource_id = art_spec
.source
.resource_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.ok_or_else(|| "统一视觉规范图缺少 resourceId".to_string())?;
// 派生素材的参考合同是「规范图前置在最前,用户参考按顺序追加在后」,图集不接受用户参考:
// 规范身份仍只由首项承担,用户参考不能顶替也不能冒充规范引用。
if !crate::agent::platform_art_runtime_references_match_request_contract(
&asset.source.reference_resource_ids,
expected_kind,
) {
return Err(format!(
"派生视觉资产未精确引用当前统一视觉规范图:{expected_path}"
));
}
let reference_resource_id = asset.source.reference_resource_ids[0].as_str();
let original_provenance_matches =
canvas_project_id == art_spec_project_id && reference_resource_id == art_spec_resource_id;
let rebound_local_source_matches = if original_provenance_matches {
true
} else {
let art_spec_bytes = resolve_local_project_path(root, &art_spec.local_path)
.ok()
.and_then(|path| fs::read(path).ok())
.ok_or_else(|| "读取统一视觉规范图失败".to_string())?;
let source_identity = new_external_editor_source_identity(
&art_spec.id,
&format!("{:x}", Sha256::digest(&art_spec_bytes)),
&art_spec.media_type,
art_spec.kind.as_str(),
)?;
external_editor_remote_reference_matches_local_source_at(
root,
&manifest.project_id,
canvas_project_id,
reference_resource_id,
&source_identity,
)?
};
if !rebound_local_source_matches {
return Err(format!(
"派生视觉资产未绑定当前统一视觉规范图的本地内容身份:{expected_path}"
));
}
Ok(())
}
pub(crate) fn set_task_status(
manifest: &mut GameCreationAppManifest,
task_id: &str,
@@ -1682,154 +1530,6 @@ pub(crate) fn add_manifest_asset_tags_at(
/// 批量标签写入的审计类型:一次批量追加只留一条记录,装的是"谁被追加了什么"。
pub(crate) const ASSET_BATCH_TAG_AUDIT_RECORD_TYPE: &str = "asset.tags.append";
pub(crate) fn create_manifest_task_at(
root: &Path,
task_id: &str,
title: &str,
group: GameCreationAppAgentGroup,
role: &str,
status: GameCreationAppTaskStatus,
dependencies: Vec<String>,
artifacts: Vec<String>,
acceptance_criteria: Vec<String>,
) -> Result<GameCreationAppTaskState, String> {
let (manifest_path, mut manifest) = read_or_create_manifest(root)?;
ensure_manifest_seed_tasks(&mut manifest);
let fallback_id = format!(
"agent-task-{}-{}",
unix_timestamp(),
manifest.tasks.len() + 1
);
let task_id = normalize_manifest_task_id(task_id, &fallback_id)?;
if manifest.tasks.iter().any(|task| task.id == task_id) {
return Err(format!("项目任务已存在:{task_id}"));
}
let title = normalize_manifest_task_text(title, "任务标题", 120)?;
let role = normalize_manifest_task_text(role, "任务角色", 64)?;
let known_task_ids = manifest
.tasks
.iter()
.map(|task| task.id.clone())
.collect::<Vec<_>>();
let dependencies = normalize_manifest_task_id_list(dependencies, "依赖任务", 8)?;
for dependency in &dependencies {
if dependency == &task_id {
return Err("任务不能依赖自己".to_string());
}
if !known_task_ids.iter().any(|known| known == dependency) {
return Err(format!("依赖任务不存在:{dependency}"));
}
}
let artifacts = normalize_manifest_task_text_list(artifacts, "任务产物", 8, 160)?;
let acceptance_criteria =
normalize_manifest_task_text_list(acceptance_criteria, "验收标准", 8, 180)?;
let task = GameCreationAppTaskState {
id: task_id,
title,
group,
role,
status,
dependencies,
artifacts,
acceptance_criteria,
};
manifest.tasks.push(task.clone());
write_manifest(&manifest_path, &manifest)?;
Ok(task)
}
fn normalize_manifest_task_id(value: &str, fallback: &str) -> Result<String, String> {
let value = value.trim();
let source = if value.is_empty() {
fallback.trim()
} else {
value
};
let normalized = source
.to_ascii_lowercase()
.chars()
.map(|character| {
if character.is_ascii_alphanumeric() || character == '-' || character == '_' {
character
} else {
'-'
}
})
.collect::<String>();
let normalized = normalized
.split('-')
.filter(|part| !part.is_empty())
.collect::<Vec<_>>()
.join("-");
if normalized.is_empty() {
return Err("任务 ID 只能包含 ASCII 字母、数字、短横线和下划线".to_string());
}
Ok(normalized.chars().take(96).collect())
}
fn normalize_manifest_task_text(
value: &str,
label: &str,
max_chars: usize,
) -> Result<String, String> {
let value = value.trim();
if value.is_empty() {
return Err(format!("{label}不能为空"));
}
if value.chars().any(char::is_control) {
return Err(format!("{label}不能包含控制字符"));
}
Ok(value.chars().take(max_chars).collect())
}
fn normalize_manifest_task_text_list(
values: Vec<String>,
label: &str,
max_items: usize,
max_chars: usize,
) -> Result<Vec<String>, String> {
let mut output = Vec::new();
for value in values {
let value = value.trim();
if value.is_empty() {
continue;
}
if value.chars().any(char::is_control) {
return Err(format!("{label}不能包含控制字符"));
}
let item = value.chars().take(max_chars).collect::<String>();
if !output.iter().any(|existing| existing == &item) {
output.push(item);
}
if output.len() > max_items {
return Err(format!("{label}最多支持 {max_items} 项"));
}
}
Ok(output)
}
fn normalize_manifest_task_id_list(
values: Vec<String>,
label: &str,
max_items: usize,
) -> Result<Vec<String>, String> {
let mut output = Vec::new();
for value in values {
let value = value.trim();
if value.is_empty() {
continue;
}
let item = normalize_manifest_task_id(value, "")?;
if !output.iter().any(|existing| existing == &item) {
output.push(item);
}
if output.len() > max_items {
return Err(format!("{label}最多支持 {max_items} 项"));
}
}
Ok(output)
}
/// 资源标签的持久化上界:manifest 每次写入都被整份序列化重写、每次读取都被整份解析,
/// 标签数量与单标签长度若无界,客户端就能让这份文件无限膨胀,并把成本摊到之后每一次读写上。
const ASSET_CLASSIFICATION_MAX_TAGS: usize = 16;
@@ -288,22 +288,6 @@ fn project_verification_package_manager_at(
Ok("npm")
}
#[cfg(test)]
pub(crate) fn resolve_project_verification_spec_at(
root: &Path,
script: &str,
expected_command: &str,
timeout_seconds: u64,
) -> Result<ProjectVerificationSpec, String> {
resolve_project_verification_spec_with_cwd_at(
root,
script,
expected_command,
timeout_seconds,
".",
)
}
pub(crate) fn resolve_project_verification_spec_with_cwd_at(
root: &Path,
script: &str,
@@ -673,24 +657,6 @@ where
})
}
#[cfg(test)]
pub(crate) async fn run_project_verification_at(
root: &Path,
script: &str,
expected_command: &str,
timeout_seconds: u64,
) -> Result<ProjectVerificationResult, String> {
run_project_verification_with_commit_at(
root,
script,
expected_command,
timeout_seconds,
".",
|| Ok(()),
)
.await
}
pub(crate) async fn run_project_verification_with_commit_at<F>(
root: &Path,
script: &str,
@@ -12,11 +12,6 @@ use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use std::sync::{mpsc, Arc, Barrier, Mutex};
use std::time::{Duration, Instant};
use crate::{
default_agent_runtime_run_profile, unix_timestamp, AgentRuntimeTaskRecord,
ProcessSessionRecord, AGENT_RUNTIME_SCHEMA_VERSION,
};
static TEST_DIRECTORY_COUNTER: AtomicU64 = AtomicU64::new(0);
#[test]