64b0d2f766
收敛外部错误体、注入提示词与锁诊断的不设界行为(PR #316 review) - resource_editor.rs:2476 [bug · medium] 已被修复:`error_details` 改为 `error.details` 优先、顶层 `details` 兜底,与 docstring 声明的取值顺序一致;新增「两种形状同时出现」用例钉住优先级。 - resource_editor.rs:2468 [performance · low] 已被修复:4xx 错误体改走 `bytes_stream` 有界读取(64KiB 上限,超限即放弃取原因),不再用 `response.text()` 无上限缓冲外部正文。 - resource_editor.rs:2514-2525 [maintainability · low] 已被修复:provider / assetKind / mediaType / code 四个追加字段各自按 80 字符截断(`MAX_EDITOR_ERROR_DETAIL_CHARS`),最终文案长度不再随外部正文增长;新增超长字段整串比对用例。 - direct_codex_references.rs:175-182 [performance · medium] 已被修复:`resourceIds` 超过 32 条直接按模块失败关闭口径拒绝,并在同一 id 的重复项上做去重,注入提示词的「关联素材 ID」行长度与 manifest 扫描次数都有界;新增去重 + 超限用例。 - direct_runtime.rs:1795-1797 [maintainability · low] 已被修复:`direct_project_history_contention_failure` 收窄为只认追加锁超时标记(项目写锁争用已在更早的专用分支判掉,是死代码),恢复提示文案相应去掉「或项目锁」;新增用例钉住两条判据互不重叠、各自命中正确提示。 - conversation.rs:984-989 [maintainability · medium] 已被修复:`DIRECT_PROJECT_HISTORY_RECORD_TYPE` 提升为 `pub(crate)`,读取侧 `is_direct_project_history_row` 引用同一常量,`project.jsonl` 的信封契约改为编译期共享。 - agent_db.rs:3570-3574 [other · medium] 已被修复:`lock_with_attempts` 的进程内锁改为有界 `try_lock` 轮询(窗口与调用方一致,2000/200 次 × 5ms),同进程读路径不再被写者的 10s 跨进程等待拖住;顺序仍是「进程内锁 → 跨进程锁」,无 ABBA;新增用例验证短窗口 <5s 失败且释放后立即可重取。 - agent_db.rs:3668-3669 [maintainability · low] 已被修复:新增与锁同级的旁路诊断文件(普通共享写入),诊断优先读它、读不到再退回锁文件,Windows 上 zero-share 独占锁文件时也能报出持锁方;新增「取锁写出旁路文件」「锁文件不可读时退到旁路」两条用例。仅存 pid/用途/时间戳,不参与判活、回收或抢占。 - agent_db.rs:3827-3828 [maintainability · low] 已被修复:仅当修复错误真的符合提权口径(`windows_acl_error_may_need_elevation`)时才追加「自动提权修复未完成」,否则原样透传底层错误,不再重复报同一个错误并谎称尝试过提权。 - main.rs:29-31 [maintainability · low] 判定为不适用/不改:实测删掉那四个 import 会产生 18 处 E0425(`direct_tool_bridge` / `canvas_generation` / `assets` / `manifest` / `resource_editor` / `resource_dependency_graph` / `recovery_tests` 等经 `use super::*` 从 crate 根复用它们),而保留它们并不产生 `unused_imports` 警告(cargo check --all-targets 无一条指向这几行);已按原样恢复 main.rs,工作树中该文件无改动。
1523 lines
57 KiB
Rust
1523 lines
57 KiB
Rust
use super::*;
|
||
|
||
use super::agent_db::{
|
||
append_jsonl_line_unlocked, ensure_conversation_message_audit_at,
|
||
is_valid_agent_db_finalization_id, project_append_lock_for,
|
||
};
|
||
|
||
#[cfg(test)]
|
||
use super::agent_db::{agent_db_record_append_class, AgentDbRecordAppendClass};
|
||
|
||
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
|
||
#[serde(rename_all = "camelCase")]
|
||
struct AgentConversationSessionCatalogFile {
|
||
schema_version: String,
|
||
agent_id: String,
|
||
active_session_id: String,
|
||
sessions: Vec<AgentConversationSessionRecord>,
|
||
}
|
||
|
||
pub(crate) fn legacy_agent_conversation_session_id(agent_id: &str) -> String {
|
||
format!("agent-session-{agent_id}")
|
||
}
|
||
|
||
pub(crate) fn normalize_agent_conversation_session_id(session_id: &str) -> Result<String, String> {
|
||
let session_id = session_id.trim();
|
||
if session_id.is_empty()
|
||
|| session_id.contains("..")
|
||
|| session_id
|
||
.chars()
|
||
.any(|ch| !(ch.is_ascii_alphanumeric() || ch == '-' || ch == '_'))
|
||
{
|
||
return Err("session id 只能包含 ASCII 字母、数字、短横线和下划线".to_string());
|
||
}
|
||
Ok(session_id.to_string())
|
||
}
|
||
|
||
fn agent_conversation_session_catalog_path(root: &Path, agent_id: &str) -> PathBuf {
|
||
root.join(".agent/runtime/sessions")
|
||
.join(format!("{agent_id}.json"))
|
||
}
|
||
|
||
fn agent_conversation_sessions_dir(root: &Path, agent_id: &str) -> PathBuf {
|
||
root.join(".agent/conversations/agents")
|
||
.join(agent_id)
|
||
.join("sessions")
|
||
}
|
||
|
||
fn conversation_file_path_for_resolved_session(
|
||
root: &Path,
|
||
agent_id: &str,
|
||
session_id: &str,
|
||
) -> PathBuf {
|
||
if session_id == legacy_agent_conversation_session_id(agent_id) {
|
||
root.join(".agent/conversations/agents")
|
||
.join(format!("{agent_id}.jsonl"))
|
||
} else {
|
||
agent_conversation_sessions_dir(root, agent_id).join(format!("{session_id}.jsonl"))
|
||
}
|
||
}
|
||
|
||
fn file_modified_timestamp(path: &Path) -> u64 {
|
||
fs::metadata(path)
|
||
.and_then(|metadata| metadata.modified())
|
||
.ok()
|
||
.and_then(|modified| modified.duration_since(UNIX_EPOCH).ok())
|
||
.map(|duration| duration.as_secs())
|
||
.unwrap_or(0)
|
||
}
|
||
|
||
fn count_conversation_messages(path: &Path) -> Result<u64, String> {
|
||
if !prepare_game_creator_private_path_for_read(path, false, "对话记录")? {
|
||
return Ok(0);
|
||
}
|
||
match File::open(path) {
|
||
Ok(file) => {
|
||
let mut count = 0_u64;
|
||
for line in BufReader::new(file).lines() {
|
||
let line =
|
||
line.map_err(|error| format!("读取对话记录失败:{}: {error}", path.display()))?;
|
||
if !line.trim().is_empty() {
|
||
count = count.saturating_add(1);
|
||
}
|
||
}
|
||
Ok(count)
|
||
}
|
||
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(0),
|
||
Err(error) => Err(format!("读取对话记录失败:{}: {error}", path.display())),
|
||
}
|
||
}
|
||
|
||
fn default_agent_conversation_session_catalog(
|
||
root: &Path,
|
||
agent_id: &str,
|
||
) -> AgentConversationSessionCatalogFile {
|
||
let legacy_session_id = legacy_agent_conversation_session_id(agent_id);
|
||
let legacy_path =
|
||
conversation_file_path_for_resolved_session(root, agent_id, &legacy_session_id);
|
||
let updated_at = file_modified_timestamp(&legacy_path);
|
||
AgentConversationSessionCatalogFile {
|
||
schema_version: AGENT_CONVERSATION_SESSION_SCHEMA_VERSION.to_string(),
|
||
agent_id: agent_id.to_string(),
|
||
active_session_id: legacy_session_id.clone(),
|
||
sessions: vec![AgentConversationSessionRecord {
|
||
session_id: legacy_session_id,
|
||
title: "默认会话".to_string(),
|
||
created_at: updated_at,
|
||
updated_at,
|
||
archived_at: None,
|
||
message_count: count_conversation_messages(&legacy_path).unwrap_or(0),
|
||
legacy: true,
|
||
forked_from_session_id: None,
|
||
forked_message_count: None,
|
||
}],
|
||
}
|
||
}
|
||
|
||
fn read_agent_conversation_session_catalog_unlocked(
|
||
root: &Path,
|
||
agent_id: &str,
|
||
) -> Result<AgentConversationSessionCatalogFile, String> {
|
||
validate_project_root(root)?;
|
||
let agent_id = normalize_conversation_agent_id(agent_id)?;
|
||
let catalog_path = agent_conversation_session_catalog_path(root, &agent_id);
|
||
prepare_game_creator_private_path_for_read(&catalog_path, false, "Agent Session 目录")?;
|
||
let mut catalog = match fs::read_to_string(&catalog_path) {
|
||
Ok(content) => serde_json::from_str::<AgentConversationSessionCatalogFile>(&content)
|
||
.map_err(|error| {
|
||
format!(
|
||
"解析 Agent Session 目录失败:{}: {error}",
|
||
catalog_path.display()
|
||
)
|
||
})?,
|
||
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
|
||
default_agent_conversation_session_catalog(root, &agent_id)
|
||
}
|
||
Err(error) => {
|
||
return Err(format!(
|
||
"读取 Agent Session 目录失败:{}: {error}",
|
||
catalog_path.display()
|
||
));
|
||
}
|
||
};
|
||
if catalog.agent_id.trim().is_empty() {
|
||
catalog.agent_id = agent_id.clone();
|
||
}
|
||
if catalog.agent_id != agent_id {
|
||
return Err(format!(
|
||
"Agent Session 目录归属不一致:期望 {agent_id},实际 {}",
|
||
catalog.agent_id
|
||
));
|
||
}
|
||
if catalog.schema_version.trim().is_empty() {
|
||
catalog.schema_version = AGENT_CONVERSATION_SESSION_SCHEMA_VERSION.to_string();
|
||
}
|
||
|
||
let legacy_session_id = legacy_agent_conversation_session_id(&agent_id);
|
||
let mut seen = std::collections::BTreeSet::new();
|
||
for session in &mut catalog.sessions {
|
||
session.session_id = normalize_agent_conversation_session_id(&session.session_id)?;
|
||
if !seen.insert(session.session_id.clone()) {
|
||
return Err(format!(
|
||
"Agent Session 目录包含重复 sessionId:{}",
|
||
session.session_id
|
||
));
|
||
}
|
||
session.legacy = session.session_id == legacy_session_id;
|
||
if session.title.trim().is_empty() {
|
||
session.title = if session.legacy {
|
||
"默认会话".to_string()
|
||
} else {
|
||
"恢复会话".to_string()
|
||
};
|
||
}
|
||
let conversation_path =
|
||
conversation_file_path_for_resolved_session(root, &agent_id, &session.session_id);
|
||
session.message_count = count_conversation_messages(&conversation_path)?;
|
||
session.updated_at = session
|
||
.updated_at
|
||
.max(file_modified_timestamp(&conversation_path));
|
||
}
|
||
if !seen.contains(&legacy_session_id) {
|
||
let legacy_path =
|
||
conversation_file_path_for_resolved_session(root, &agent_id, &legacy_session_id);
|
||
let updated_at = file_modified_timestamp(&legacy_path);
|
||
catalog.sessions.insert(
|
||
0,
|
||
AgentConversationSessionRecord {
|
||
session_id: legacy_session_id.clone(),
|
||
title: "默认会话".to_string(),
|
||
created_at: updated_at,
|
||
updated_at,
|
||
archived_at: None,
|
||
message_count: count_conversation_messages(&legacy_path)?,
|
||
legacy: true,
|
||
forked_from_session_id: None,
|
||
forked_message_count: None,
|
||
},
|
||
);
|
||
seen.insert(legacy_session_id.clone());
|
||
}
|
||
|
||
let sessions_dir = agent_conversation_sessions_dir(root, &agent_id);
|
||
match fs::read_dir(&sessions_dir) {
|
||
Ok(read_dir) => {
|
||
for entry in read_dir {
|
||
let entry = entry.map_err(|error| {
|
||
format!(
|
||
"读取 Agent Session 对话目录失败:{}: {error}",
|
||
sessions_dir.display()
|
||
)
|
||
})?;
|
||
let path = entry.path();
|
||
if path.extension().and_then(|value| value.to_str()) != Some("jsonl") {
|
||
continue;
|
||
}
|
||
let Some(session_id) = path.file_stem().and_then(|value| value.to_str()) else {
|
||
continue;
|
||
};
|
||
let session_id = normalize_agent_conversation_session_id(session_id)?;
|
||
if seen.insert(session_id.clone()) {
|
||
let updated_at = file_modified_timestamp(&path);
|
||
catalog.sessions.push(AgentConversationSessionRecord {
|
||
session_id,
|
||
title: "恢复会话".to_string(),
|
||
created_at: updated_at,
|
||
updated_at,
|
||
archived_at: None,
|
||
message_count: count_conversation_messages(&path)?,
|
||
legacy: false,
|
||
forked_from_session_id: None,
|
||
forked_message_count: None,
|
||
});
|
||
}
|
||
}
|
||
}
|
||
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
|
||
Err(error) => {
|
||
return Err(format!(
|
||
"读取 Agent Session 对话目录失败:{}: {error}",
|
||
sessions_dir.display()
|
||
));
|
||
}
|
||
}
|
||
|
||
for session in &mut catalog.sessions {
|
||
session.forked_from_session_id = session
|
||
.forked_from_session_id
|
||
.as_deref()
|
||
.map(normalize_agent_conversation_session_id)
|
||
.transpose()?;
|
||
match (
|
||
session.forked_from_session_id.as_deref(),
|
||
session.forked_message_count,
|
||
) {
|
||
(Some(source_session_id), Some(message_count)) => {
|
||
if source_session_id == session.session_id {
|
||
return Err(format!(
|
||
"Agent Session 分叉来源不能指向自身:{}",
|
||
session.session_id
|
||
));
|
||
}
|
||
if !seen.contains(source_session_id) {
|
||
return Err(format!(
|
||
"Agent Session 分叉来源不存在:{} -> {source_session_id}",
|
||
session.session_id
|
||
));
|
||
}
|
||
if message_count > session.message_count {
|
||
return Err(format!(
|
||
"Agent Session 分叉消息数超过当前消息数:{}",
|
||
session.session_id
|
||
));
|
||
}
|
||
}
|
||
(None, None) => {}
|
||
_ => {
|
||
return Err(format!(
|
||
"Agent Session 分叉来源字段不完整:{}",
|
||
session.session_id
|
||
));
|
||
}
|
||
}
|
||
}
|
||
|
||
let active_is_valid = catalog.sessions.iter().any(|session| {
|
||
session.session_id == catalog.active_session_id && session.archived_at.is_none()
|
||
});
|
||
if !active_is_valid {
|
||
catalog.active_session_id = catalog
|
||
.sessions
|
||
.iter()
|
||
.find(|session| session.archived_at.is_none())
|
||
.map(|session| session.session_id.clone())
|
||
.unwrap_or_else(|| legacy_session_id.clone());
|
||
}
|
||
catalog.sessions.sort_by(|left, right| {
|
||
right
|
||
.legacy
|
||
.cmp(&left.legacy)
|
||
.then_with(|| left.created_at.cmp(&right.created_at))
|
||
.then_with(|| left.session_id.cmp(&right.session_id))
|
||
});
|
||
Ok(catalog)
|
||
}
|
||
|
||
fn write_agent_conversation_session_catalog_unlocked(
|
||
root: &Path,
|
||
catalog: &AgentConversationSessionCatalogFile,
|
||
) -> Result<(), String> {
|
||
let path = agent_conversation_session_catalog_path(root, &catalog.agent_id);
|
||
if let Some(parent) = path.parent() {
|
||
ensure_game_creator_private_directory_tree(parent, "Agent Session 目录")?;
|
||
prepare_game_creator_private_path_for_read(parent, true, "Agent Session 目录")?;
|
||
}
|
||
let content = serde_json::to_string_pretty(catalog)
|
||
.map_err(|error| format!("序列化 Agent Session 目录失败:{error}"))?;
|
||
write_game_creator_private_file(
|
||
&path,
|
||
format!("{content}\n").as_bytes(),
|
||
"Agent Session 目录",
|
||
)
|
||
}
|
||
|
||
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)
|
||
})
|
||
}
|
||
|
||
fn agent_conversation_session_list_result(
|
||
root: &Path,
|
||
catalog: AgentConversationSessionCatalogFile,
|
||
) -> AgentConversationSessionListResult {
|
||
AgentConversationSessionListResult {
|
||
path: agent_conversation_session_catalog_path(root, &catalog.agent_id)
|
||
.to_string_lossy()
|
||
.into_owned(),
|
||
agent_id: catalog.agent_id,
|
||
active_session_id: catalog.active_session_id,
|
||
sessions: catalog.sessions,
|
||
}
|
||
}
|
||
|
||
pub(crate) fn list_game_creator_agent_sessions_at(
|
||
root: &Path,
|
||
agent_id: &str,
|
||
) -> Result<AgentConversationSessionListResult, String> {
|
||
let agent_id = normalize_game_creator_runtime_agent_id(agent_id)?;
|
||
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 目录")?;
|
||
read_agent_conversation_session_catalog_unlocked(root, &agent_id)
|
||
.map(|catalog| agent_conversation_session_list_result(root, catalog))
|
||
}
|
||
|
||
fn agent_runtime_task_is_terminal_for_session_mutation(record: &AgentRuntimeTaskRecord) -> bool {
|
||
matches!(record.status.as_str(), "completed" | "failed" | "cancelled")
|
||
&& record.phase != "needs-reconciliation"
|
||
}
|
||
|
||
pub(crate) fn ensure_agent_session_has_no_live_tasks(
|
||
root: &Path,
|
||
agent_id: &str,
|
||
session_id: &str,
|
||
) -> Result<(), String> {
|
||
let mut latest_by_run = BTreeMap::<String, AgentRuntimeTaskRecord>::new();
|
||
let task_path = root
|
||
.join(".agent/runtime/tasks")
|
||
.join(format!("{agent_id}.jsonl"));
|
||
prepare_game_creator_private_path_for_read(&task_path, false, "Agent Runtime 任务")?;
|
||
match File::open(&task_path) {
|
||
Ok(file) => {
|
||
for line in BufReader::new(file).lines() {
|
||
let line = line.map_err(|error| {
|
||
format!(
|
||
"读取 Agent Runtime 任务失败:{}: {error}",
|
||
task_path.display()
|
||
)
|
||
})?;
|
||
let line = line.trim();
|
||
if line.is_empty() {
|
||
continue;
|
||
}
|
||
let record =
|
||
serde_json::from_str::<AgentRuntimeTaskRecord>(line).map_err(|error| {
|
||
format!(
|
||
"解析 Agent Runtime 任务失败:{}: {error}",
|
||
task_path.display()
|
||
)
|
||
})?;
|
||
latest_by_run.insert(record.run_id.clone(), record);
|
||
}
|
||
}
|
||
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
|
||
Err(error) => {
|
||
return Err(format!(
|
||
"读取 Agent Runtime 任务失败:{}: {error}",
|
||
task_path.display()
|
||
));
|
||
}
|
||
}
|
||
let blocked = latest_by_run.values().find(|record| {
|
||
record.session_id == session_id
|
||
&& !agent_runtime_task_is_terminal_for_session_mutation(record)
|
||
});
|
||
if let Some(record) = blocked {
|
||
return Err(format!(
|
||
"Session {} 仍有未结束任务 {}({} / {}),不能切换、归档或分叉",
|
||
session_id, record.run_id, record.status, record.phase
|
||
));
|
||
}
|
||
let parent_run_ids = latest_by_run
|
||
.values()
|
||
.filter(|record| record.session_id == session_id)
|
||
.map(|record| record.run_id.as_str())
|
||
.collect::<Vec<_>>();
|
||
if parent_run_ids.is_empty() {
|
||
return Ok(());
|
||
}
|
||
let task_dir = root.join(".agent/runtime/tasks");
|
||
let entries = match fs::read_dir(&task_dir) {
|
||
Ok(entries) => entries,
|
||
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(()),
|
||
Err(error) => {
|
||
return Err(format!(
|
||
"读取 Agent Runtime 任务目录失败:{}: {error}",
|
||
task_dir.display()
|
||
));
|
||
}
|
||
};
|
||
for entry in entries {
|
||
let entry = entry.map_err(|error| {
|
||
format!(
|
||
"读取 Agent Runtime 任务目录项失败:{}: {error}",
|
||
task_dir.display()
|
||
)
|
||
})?;
|
||
if !entry
|
||
.file_type()
|
||
.map_err(|error| format!("读取 Agent Runtime 任务文件类型失败:{error}"))?
|
||
.is_file()
|
||
{
|
||
continue;
|
||
}
|
||
let path = entry.path();
|
||
if path.extension().and_then(|value| value.to_str()) != Some("jsonl") {
|
||
continue;
|
||
}
|
||
prepare_game_creator_private_path_for_read(&path, false, "Agent Runtime 任务")?;
|
||
let mut delegated_latest_by_run = BTreeMap::<String, AgentRuntimeTaskRecord>::new();
|
||
let file = File::open(&path)
|
||
.map_err(|error| format!("读取 Agent Runtime 任务失败:{}: {error}", path.display()))?;
|
||
for line in BufReader::new(file).lines() {
|
||
let line = line.map_err(|error| {
|
||
format!("读取 Agent Runtime 任务失败:{}: {error}", path.display())
|
||
})?;
|
||
let line = line.trim();
|
||
if line.is_empty() {
|
||
continue;
|
||
}
|
||
let record = serde_json::from_str::<AgentRuntimeTaskRecord>(line).map_err(|error| {
|
||
format!("解析 Agent Runtime 任务失败:{}: {error}", path.display())
|
||
})?;
|
||
delegated_latest_by_run.insert(record.run_id.clone(), record);
|
||
}
|
||
let delegated_child = delegated_latest_by_run.values().find(|record| {
|
||
record.parent_agent_id.as_deref() == Some(agent_id)
|
||
&& record
|
||
.parent_run_id
|
||
.as_deref()
|
||
.is_some_and(|parent_run_id| parent_run_ids.contains(&parent_run_id))
|
||
&& !agent_runtime_task_is_terminal_for_session_mutation(record)
|
||
});
|
||
if let Some(record) = delegated_child {
|
||
return Err(format!(
|
||
"Session {} 的父任务 {} 仍有委派子任务 {}({} / {}),不能切换、归档或分叉",
|
||
session_id,
|
||
record.parent_run_id.as_deref().unwrap_or("-"),
|
||
record.run_id,
|
||
record.status,
|
||
record.phase
|
||
));
|
||
}
|
||
}
|
||
Ok(())
|
||
}
|
||
|
||
pub(crate) fn create_game_creator_agent_session_at(
|
||
root: &Path,
|
||
agent_id: &str,
|
||
title: &str,
|
||
) -> Result<AgentConversationSessionListResult, String> {
|
||
let agent_id = normalize_game_creator_runtime_agent_id(agent_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)?;
|
||
ensure_agent_session_has_no_live_tasks(root, &agent_id, &catalog.active_session_id)?;
|
||
let session_id = new_agent_conversation_session_id(&catalog, &agent_id)?;
|
||
let now = unix_timestamp();
|
||
let title = title.trim();
|
||
let title = if title.is_empty() {
|
||
"新会话".to_string()
|
||
} else {
|
||
title.chars().take(80).collect::<String>()
|
||
};
|
||
let conversation_path =
|
||
conversation_file_path_for_resolved_session(root, &agent_id, &session_id);
|
||
if let Some(parent) = conversation_path.parent() {
|
||
ensure_game_creator_private_directory_tree(parent, "对话目录")?;
|
||
prepare_game_creator_private_path_for_read(parent, true, "对话目录")?;
|
||
}
|
||
prepare_game_creator_private_path_for_read(&conversation_path, false, "对话记录")?;
|
||
let mut options = fs::OpenOptions::new();
|
||
options.create_new(true).write(true);
|
||
#[cfg(windows)]
|
||
{
|
||
use std::os::windows::fs::OpenOptionsExt;
|
||
options.custom_flags(PROJECT_FILE_FLAG_OPEN_REPARSE_POINT);
|
||
}
|
||
let file = options.open(&conversation_path).map_err(|error| {
|
||
format!(
|
||
"创建 Agent Session 对话失败:{}: {error}",
|
||
conversation_path.display()
|
||
)
|
||
})?;
|
||
if let Err(error) =
|
||
harden_new_game_creator_private_path(&conversation_path, false, "对话记录")
|
||
{
|
||
drop(file);
|
||
let _ = fs::remove_file(&conversation_path);
|
||
return Err(error);
|
||
}
|
||
drop(file);
|
||
catalog.sessions.push(AgentConversationSessionRecord {
|
||
session_id: session_id.clone(),
|
||
title,
|
||
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)?;
|
||
Ok(agent_conversation_session_list_result(root, catalog))
|
||
})
|
||
}
|
||
|
||
fn new_agent_conversation_session_id(
|
||
catalog: &AgentConversationSessionCatalogFile,
|
||
agent_id: &str,
|
||
) -> Result<String, String> {
|
||
(0_u32..100)
|
||
.map(|attempt| {
|
||
format!(
|
||
"agent-session-{agent_id}-{}-{attempt}",
|
||
SystemTime::now()
|
||
.duration_since(UNIX_EPOCH)
|
||
.unwrap_or_default()
|
||
.as_nanos()
|
||
)
|
||
})
|
||
.find(|candidate| {
|
||
!catalog
|
||
.sessions
|
||
.iter()
|
||
.any(|session| session.session_id == *candidate)
|
||
})
|
||
.ok_or_else(|| "无法生成唯一 Agent Session ID".to_string())
|
||
}
|
||
|
||
pub(crate) fn fork_game_creator_agent_session_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,
|
||
write_agent_conversation_session_catalog_unlocked,
|
||
)
|
||
}
|
||
|
||
fn fork_game_creator_agent_session_with_catalog_writer_at<F>(
|
||
root: &Path,
|
||
agent_id: &str,
|
||
source_session_id: &str,
|
||
title: &str,
|
||
write_catalog: F,
|
||
) -> Result<AgentConversationSessionListResult, String>
|
||
where
|
||
F: FnOnce(&Path, &AgentConversationSessionCatalogFile) -> Result<(), String>,
|
||
{
|
||
validate_project_root(root)?;
|
||
let agent_id = normalize_game_creator_runtime_agent_id(agent_id)?;
|
||
with_agent_conversation_session_lane_at(root, &agent_id, "Agent Session 分叉", || {
|
||
let source_session_id = normalize_agent_conversation_session_id(source_session_id)?;
|
||
let source_path =
|
||
conversation_file_path_for_resolved_session(root, &agent_id, &source_session_id);
|
||
let source_lock = project_append_lock_for(&source_path)?;
|
||
let _source_guard = source_lock.lock("Agent Session 分叉源对话")?;
|
||
|
||
let catalog_path = agent_conversation_session_catalog_path(root, &agent_id);
|
||
let catalog_lock = project_append_lock_for(&catalog_path)?;
|
||
let _catalog_guard = catalog_lock.lock("Agent Session 目录")?;
|
||
let mut catalog = read_agent_conversation_session_catalog_unlocked(root, &agent_id)?;
|
||
let source = catalog
|
||
.sessions
|
||
.iter()
|
||
.find(|session| session.session_id == source_session_id)
|
||
.cloned()
|
||
.ok_or_else(|| format!("Agent Session 不存在:{source_session_id}"))?;
|
||
ensure_agent_session_has_no_live_tasks(root, &agent_id, &catalog.active_session_id)?;
|
||
if source.session_id != catalog.active_session_id {
|
||
ensure_agent_session_has_no_live_tasks(root, &agent_id, &source.session_id)?;
|
||
}
|
||
|
||
let records = read_persisted_local_conversation_records_unlocked(&source_path)?;
|
||
let session_id = new_agent_conversation_session_id(&catalog, &agent_id)?;
|
||
let conversation_path =
|
||
conversation_file_path_for_resolved_session(root, &agent_id, &session_id);
|
||
if let Some(parent) = conversation_path.parent() {
|
||
ensure_game_creator_private_directory_tree(parent, "对话目录")?;
|
||
prepare_game_creator_private_path_for_read(parent, true, "对话目录")?;
|
||
}
|
||
prepare_game_creator_private_path_for_read(
|
||
&conversation_path,
|
||
false,
|
||
"Agent Session 分叉对话",
|
||
)?;
|
||
let write_result = (|| -> Result<(), String> {
|
||
let mut options = fs::OpenOptions::new();
|
||
options.create_new(true).write(true);
|
||
#[cfg(windows)]
|
||
{
|
||
use std::os::windows::fs::OpenOptionsExt;
|
||
options.custom_flags(PROJECT_FILE_FLAG_OPEN_REPARSE_POINT);
|
||
}
|
||
let mut file = options.open(&conversation_path).map_err(|error| {
|
||
format!(
|
||
"创建 Agent Session 分叉对话失败:{}: {error}",
|
||
conversation_path.display()
|
||
)
|
||
})?;
|
||
for record in &records {
|
||
serde_json::to_writer(&mut file, record)
|
||
.map_err(|error| format!("序列化 Agent Session 分叉消息失败:{error}"))?;
|
||
file.write_all(b"\n").map_err(|error| {
|
||
format!(
|
||
"写入 Agent Session 分叉对话失败:{}: {error}",
|
||
conversation_path.display()
|
||
)
|
||
})?;
|
||
}
|
||
file.flush().map_err(|error| {
|
||
format!(
|
||
"刷新 Agent Session 分叉对话失败:{}: {error}",
|
||
conversation_path.display()
|
||
)
|
||
})
|
||
})();
|
||
if let Err(error) = write_result {
|
||
let _ = fs::remove_file(&conversation_path);
|
||
return Err(error);
|
||
}
|
||
prepare_game_creator_private_path_for_read(
|
||
&conversation_path,
|
||
false,
|
||
"Agent Session 分叉对话",
|
||
)?;
|
||
|
||
let now = unix_timestamp();
|
||
let requested_title = title.trim();
|
||
let title = if requested_title.is_empty() {
|
||
format!("分支:{}", source.title).chars().take(80).collect()
|
||
} else {
|
||
requested_title.chars().take(80).collect()
|
||
};
|
||
catalog.sessions.push(AgentConversationSessionRecord {
|
||
session_id: session_id.clone(),
|
||
title,
|
||
created_at: now,
|
||
updated_at: now,
|
||
archived_at: None,
|
||
message_count: records.len() as u64,
|
||
legacy: false,
|
||
forked_from_session_id: Some(source.session_id),
|
||
forked_message_count: Some(records.len() as u64),
|
||
});
|
||
catalog.active_session_id = session_id;
|
||
if let Err(error) = write_catalog(root, &catalog) {
|
||
let cleanup = fs::remove_file(&conversation_path);
|
||
return Err(match cleanup {
|
||
Ok(()) => error,
|
||
Err(cleanup_error) => format!(
|
||
"{error};清理未登记分叉对话失败:{}: {cleanup_error}",
|
||
conversation_path.display()
|
||
),
|
||
});
|
||
}
|
||
Ok(agent_conversation_session_list_result(root, catalog))
|
||
})
|
||
}
|
||
|
||
#[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)
|
||
},
|
||
)
|
||
}
|
||
|
||
pub(crate) fn set_active_game_creator_agent_session_at(
|
||
root: &Path,
|
||
agent_id: &str,
|
||
session_id: &str,
|
||
) -> Result<AgentConversationSessionListResult, String> {
|
||
let agent_id = normalize_game_creator_runtime_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)?;
|
||
let session = catalog
|
||
.sessions
|
||
.iter()
|
||
.find(|session| session.session_id == session_id)
|
||
.ok_or_else(|| format!("Agent Session 不存在:{session_id}"))?;
|
||
if session.archived_at.is_some() {
|
||
return Err(format!(
|
||
"Agent Session 已归档,不能设为活动会话:{session_id}"
|
||
));
|
||
}
|
||
ensure_agent_session_has_no_live_tasks(root, &agent_id, &catalog.active_session_id)?;
|
||
catalog.active_session_id = session_id;
|
||
write_agent_conversation_session_catalog_unlocked(root, &catalog)?;
|
||
Ok(agent_conversation_session_list_result(root, catalog))
|
||
})
|
||
}
|
||
|
||
pub(crate) fn archive_game_creator_agent_session_at(
|
||
root: &Path,
|
||
agent_id: &str,
|
||
session_id: &str,
|
||
) -> Result<AgentConversationSessionListResult, String> {
|
||
let agent_id = normalize_game_creator_runtime_agent_id(agent_id)?;
|
||
let session_id = normalize_agent_conversation_session_id(session_id)?;
|
||
if session_id == legacy_agent_conversation_session_id(&agent_id) {
|
||
return Err("默认 legacy Session 不能归档".to_string());
|
||
}
|
||
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)?;
|
||
ensure_agent_session_has_no_live_tasks(root, &agent_id, &session_id)?;
|
||
let session = catalog
|
||
.sessions
|
||
.iter_mut()
|
||
.find(|session| session.session_id == session_id)
|
||
.ok_or_else(|| format!("Agent Session 不存在:{session_id}"))?;
|
||
if session.archived_at.is_none() {
|
||
session.archived_at = Some(unix_timestamp());
|
||
session.updated_at = unix_timestamp();
|
||
}
|
||
if catalog.active_session_id == session_id {
|
||
catalog.active_session_id = catalog
|
||
.sessions
|
||
.iter()
|
||
.find(|candidate| candidate.archived_at.is_none())
|
||
.map(|candidate| candidate.session_id.clone())
|
||
.ok_or_else(|| "至少需要保留一个未归档 Agent Session".to_string())?;
|
||
}
|
||
write_agent_conversation_session_catalog_unlocked(root, &catalog)?;
|
||
Ok(agent_conversation_session_list_result(root, catalog))
|
||
})
|
||
}
|
||
|
||
fn resolve_agent_conversation_session_at(
|
||
root: &Path,
|
||
agent_id: &str,
|
||
session_id: Option<&str>,
|
||
) -> Result<AgentConversationSessionRecord, String> {
|
||
let catalog = read_agent_conversation_session_catalog_unlocked(root, agent_id)?;
|
||
let session_id = match session_id.map(str::trim).filter(|value| !value.is_empty()) {
|
||
Some(session_id) => normalize_agent_conversation_session_id(session_id)?,
|
||
None => catalog.active_session_id.clone(),
|
||
};
|
||
catalog
|
||
.sessions
|
||
.into_iter()
|
||
.find(|session| session.session_id == session_id)
|
||
.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,
|
||
session_id: &str,
|
||
message_count: u64,
|
||
message_appended: bool,
|
||
) -> Result<(), String> {
|
||
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)?;
|
||
let session = catalog
|
||
.sessions
|
||
.iter_mut()
|
||
.find(|session| session.session_id == session_id)
|
||
.ok_or_else(|| format!("Agent Session 不存在:{session_id}"))?;
|
||
session.message_count = message_count;
|
||
if message_appended {
|
||
session.updated_at = unix_timestamp();
|
||
}
|
||
write_agent_conversation_session_catalog_unlocked(root, &catalog)
|
||
}
|
||
|
||
#[derive(Debug, Deserialize, Eq, PartialEq, Serialize)]
|
||
#[serde(rename_all = "camelCase")]
|
||
struct PersistedLocalConversationMessageRecord {
|
||
schema_version: String,
|
||
role: String,
|
||
content: String,
|
||
agent_id: Option<String>,
|
||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||
message_id: Option<String>,
|
||
updated_at: u64,
|
||
}
|
||
|
||
impl PersistedLocalConversationMessageRecord {
|
||
fn to_public_record(&self) -> LocalConversationMessageRecord {
|
||
LocalConversationMessageRecord {
|
||
schema_version: self.schema_version.clone(),
|
||
role: self.role.clone(),
|
||
content: self.content.clone(),
|
||
agent_id: self.agent_id.clone(),
|
||
message_id: self.message_id.clone(),
|
||
updated_at: self.updated_at,
|
||
}
|
||
}
|
||
}
|
||
|
||
fn normalize_local_conversation_message_id(message_id: &str) -> Result<String, String> {
|
||
let message_id = message_id.trim();
|
||
if message_id.is_empty() {
|
||
return Err("对话 messageId 不能为空".to_string());
|
||
}
|
||
if message_id.chars().any(char::is_control) {
|
||
return Err("对话 messageId 不能包含控制字符".to_string());
|
||
}
|
||
Ok(message_id.to_string())
|
||
}
|
||
|
||
fn read_persisted_local_conversation_records_unlocked(
|
||
path: &Path,
|
||
) -> Result<Vec<PersistedLocalConversationMessageRecord>, String> {
|
||
let mut records = Vec::new();
|
||
prepare_game_creator_private_path_for_read(path, false, "对话记录")?;
|
||
match File::open(path) {
|
||
Ok(file) => {
|
||
for line in BufReader::new(file).lines() {
|
||
let line =
|
||
line.map_err(|error| format!("读取对话记录失败:{}: {error}", path.display()))?;
|
||
let line = line.trim();
|
||
if line.is_empty() {
|
||
continue;
|
||
}
|
||
let record =
|
||
match serde_json::from_str::<PersistedLocalConversationMessageRecord>(line) {
|
||
Ok(record) => record,
|
||
// 项目主对话(agent_id=None)与 DirectProject 共用同一份
|
||
// `project.jsonl`:DirectProject 写的是 Codex rollout 的
|
||
// `response_item` 行。通用对话链只认自己的 legacy 行,遇到对方那种
|
||
// 行就跳过;其余任何坏行继续失败关闭。
|
||
Err(_) if is_direct_project_history_row(line) => continue,
|
||
Err(error) => {
|
||
return Err(format!("解析对话记录失败:{}: {error}", path.display()))
|
||
}
|
||
};
|
||
records.push(record);
|
||
}
|
||
}
|
||
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
|
||
Err(error) => return Err(format!("读取对话记录失败:{}: {error}", path.display())),
|
||
}
|
||
Ok(records)
|
||
}
|
||
|
||
/// 只认 DirectProject 那一种明确枚举的行信封:`type=response_item` 且带 `payload`。
|
||
/// 其余任何形状都不算“对方的行”,仍由上方按坏行失败关闭。
|
||
///
|
||
/// 类型串引用写入侧的 `DIRECT_PROJECT_HISTORY_RECORD_TYPE`,不复制字面量:两边读写的
|
||
/// 是同一份 `project.jsonl`,一旦只在写入侧改名,这里就会把每一行 DirectProject 记录
|
||
/// 都判成坏行并以 `解析对话记录失败` 失败关闭,共享文件的读取整条断掉。
|
||
fn is_direct_project_history_row(line: &str) -> bool {
|
||
serde_json::from_str::<serde_json::Value>(line).is_ok_and(|parsed| {
|
||
parsed.get("type").and_then(serde_json::Value::as_str)
|
||
== Some(crate::agent::DIRECT_PROJECT_HISTORY_RECORD_TYPE)
|
||
&& parsed.get("payload").is_some()
|
||
})
|
||
}
|
||
|
||
fn persisted_local_conversation_message_index_by_id(
|
||
records: &[PersistedLocalConversationMessageRecord],
|
||
expected_agent_id: Option<&str>,
|
||
expected_session_id: Option<&str>,
|
||
message_id: &str,
|
||
) -> Result<Option<usize>, String> {
|
||
let mut matched_index: Option<usize> = None;
|
||
for (index, record) in records.iter().enumerate() {
|
||
if record.message_id.as_deref() != Some(message_id) {
|
||
continue;
|
||
}
|
||
let conflicts_with_scope = record.agent_id.as_deref() != expected_agent_id;
|
||
let conflicts_with_existing = matched_index.is_some_and(|matched_index| {
|
||
let matched = &records[matched_index];
|
||
matched.role != record.role
|
||
|| matched.content != record.content
|
||
|| matched.agent_id != record.agent_id
|
||
});
|
||
if conflicts_with_scope || conflicts_with_existing {
|
||
return Err(format!(
|
||
"对话 messageId 冲突:{message_id} 在 agent={} session={} 下对应不同消息",
|
||
expected_agent_id.unwrap_or("project"),
|
||
expected_session_id.unwrap_or("project")
|
||
));
|
||
}
|
||
matched_index.get_or_insert(index);
|
||
}
|
||
Ok(matched_index)
|
||
}
|
||
|
||
fn local_conversation_result_from_persisted_records(
|
||
path: &Path,
|
||
agent_id: Option<String>,
|
||
session_id: Option<String>,
|
||
records: &[PersistedLocalConversationMessageRecord],
|
||
) -> LocalConversationResult {
|
||
LocalConversationResult {
|
||
path: path.to_string_lossy().into_owned(),
|
||
agent_id,
|
||
session_id,
|
||
messages: records
|
||
.iter()
|
||
.map(PersistedLocalConversationMessageRecord::to_public_record)
|
||
.collect(),
|
||
}
|
||
}
|
||
|
||
pub(crate) fn read_local_conversation_for_session_at(
|
||
root: &Path,
|
||
agent_id: Option<&str>,
|
||
session_id: Option<&str>,
|
||
) -> Result<LocalConversationResult, String> {
|
||
validate_project_root(root)?;
|
||
let (path, normalized_agent_id, normalized_session_id) = match agent_id
|
||
.map(str::trim)
|
||
.filter(|value| !value.is_empty())
|
||
{
|
||
Some(agent_id) => {
|
||
let agent_id = normalize_conversation_agent_id(agent_id)?;
|
||
let session = resolve_agent_conversation_session_at(root, &agent_id, session_id)?;
|
||
let path =
|
||
conversation_file_path_for_resolved_session(root, &agent_id, &session.session_id);
|
||
(path, Some(agent_id), Some(session.session_id))
|
||
}
|
||
None => {
|
||
if session_id
|
||
.map(str::trim)
|
||
.is_some_and(|value| !value.is_empty())
|
||
{
|
||
return Err("项目主对话不接受 sessionId".to_string());
|
||
}
|
||
(root.join(".agent/conversations/project.jsonl"), None, None)
|
||
}
|
||
};
|
||
let append_lock = project_append_lock_for(&path)?;
|
||
// 只读路径用短窗口:整份读只要一份当前快照,调用方(面板刷新、下一轮 prompt 组装)
|
||
// 自己会重跑,长时间阻塞只会把它一起拖住;写路径继续用完整窗口。
|
||
let _append_guard = append_lock.lock_short("对话记录追加写")?;
|
||
let records = read_persisted_local_conversation_records_unlocked(&path)?;
|
||
Ok(local_conversation_result_from_persisted_records(
|
||
&path,
|
||
normalized_agent_id,
|
||
normalized_session_id,
|
||
&records,
|
||
))
|
||
}
|
||
|
||
pub(crate) fn read_local_conversation_at(
|
||
root: &Path,
|
||
agent_id: Option<&str>,
|
||
) -> Result<LocalConversationResult, String> {
|
||
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>,
|
||
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, false,
|
||
)
|
||
}
|
||
|
||
fn read_local_conversation_message_by_id_for_session_internal_at(
|
||
root: &Path,
|
||
agent_id: Option<&str>,
|
||
session_id: Option<&str>,
|
||
message_id: &str,
|
||
touch_session: bool,
|
||
) -> Result<Option<LocalConversationMessageRecord>, String> {
|
||
let message_id = normalize_local_conversation_message_id(message_id)?;
|
||
let (path, normalized_agent_id, normalized_session_id) =
|
||
conversation_file_path_for_session(root, agent_id, session_id)?;
|
||
let append_lock = project_append_lock_for(&path)?;
|
||
let _append_guard = append_lock.lock_short("对话记录追加写")?;
|
||
let records = read_persisted_local_conversation_records_unlocked(&path)?;
|
||
let matched_index = persisted_local_conversation_message_index_by_id(
|
||
&records,
|
||
normalized_agent_id.as_deref(),
|
||
normalized_session_id.as_deref(),
|
||
&message_id,
|
||
)?;
|
||
if touch_session {
|
||
if let (Some(agent_id), Some(session_id)) = (
|
||
normalized_agent_id.as_deref(),
|
||
normalized_session_id.as_deref(),
|
||
) {
|
||
touch_agent_conversation_session_at(
|
||
root,
|
||
agent_id,
|
||
session_id,
|
||
records.len() as u64,
|
||
false,
|
||
)?;
|
||
}
|
||
}
|
||
Ok(matched_index.map(|index| records[index].to_public_record()))
|
||
}
|
||
|
||
fn append_local_conversation_message_for_session_internal_at(
|
||
root: &Path,
|
||
agent_id: Option<&str>,
|
||
session_id: Option<&str>,
|
||
message: LocalConversationMessage,
|
||
message_id: Option<&str>,
|
||
finalization_id: Option<&str>,
|
||
) -> Result<(LocalConversationResult, bool), String> {
|
||
let LocalConversationMessage {
|
||
role,
|
||
content,
|
||
agent_id: _client_agent_id,
|
||
} = message;
|
||
let (path, normalized_agent_id, normalized_session_id, archived) = match agent_id
|
||
.map(str::trim)
|
||
.filter(|value| !value.is_empty())
|
||
{
|
||
Some(agent_id) => {
|
||
let agent_id = normalize_conversation_agent_id(agent_id)?;
|
||
let session = resolve_agent_conversation_session_at(root, &agent_id, session_id)?;
|
||
let path =
|
||
conversation_file_path_for_resolved_session(root, &agent_id, &session.session_id);
|
||
(
|
||
path,
|
||
Some(agent_id),
|
||
Some(session.session_id),
|
||
session.archived_at.is_some(),
|
||
)
|
||
}
|
||
None => {
|
||
if session_id
|
||
.map(str::trim)
|
||
.is_some_and(|value| !value.is_empty())
|
||
{
|
||
return Err("项目主对话不接受 sessionId".to_string());
|
||
}
|
||
(
|
||
root.join(".agent/conversations/project.jsonl"),
|
||
None,
|
||
None,
|
||
false,
|
||
)
|
||
}
|
||
};
|
||
let role = role.trim();
|
||
if !matches!(role, "user" | "assistant" | "tool") {
|
||
return Err("对话角色必须是 user、assistant 或 tool".to_string());
|
||
}
|
||
let content = content.trim();
|
||
if content.is_empty() {
|
||
if message_id.is_some() {
|
||
return Err("带 messageId 的对话内容不能为空".to_string());
|
||
}
|
||
return read_local_conversation_for_session_at(root, agent_id, session_id)
|
||
.map(|conversation| (conversation, false));
|
||
}
|
||
let message_id = message_id
|
||
.map(normalize_local_conversation_message_id)
|
||
.transpose()?;
|
||
let finalization_id = finalization_id.map(str::trim);
|
||
if let Some(finalization_id) = finalization_id {
|
||
if role != "assistant"
|
||
|| normalized_agent_id.is_none()
|
||
|| normalized_session_id.is_none()
|
||
|| !is_valid_agent_db_finalization_id(finalization_id)
|
||
|| !message_id
|
||
.as_deref()
|
||
.is_some_and(is_valid_agent_db_finalization_id)
|
||
{
|
||
return Err("finalization conversation 审计身份或角色无效".to_string());
|
||
}
|
||
}
|
||
#[cfg(test)]
|
||
take_local_conversation_append_failure_injection(root, role)?;
|
||
let record = PersistedLocalConversationMessageRecord {
|
||
schema_version: LOCAL_CONVERSATION_SCHEMA_VERSION.to_string(),
|
||
role: role.to_string(),
|
||
content: content.to_string(),
|
||
agent_id: normalized_agent_id.clone(),
|
||
message_id: message_id.clone(),
|
||
updated_at: unix_timestamp(),
|
||
};
|
||
let line =
|
||
serde_json::to_string(&record).map_err(|error| format!("序列化对话记录失败:{error}"))?;
|
||
let append_lock = project_append_lock_for(&path)?;
|
||
let _append_guard = append_lock.lock("对话记录追加写")?;
|
||
let mut records = if message_id.is_some() {
|
||
read_persisted_local_conversation_records_unlocked(&path)?
|
||
} else {
|
||
Vec::new()
|
||
};
|
||
let matched_index = match message_id.as_deref() {
|
||
Some(message_id) => persisted_local_conversation_message_index_by_id(
|
||
&records,
|
||
normalized_agent_id.as_deref(),
|
||
normalized_session_id.as_deref(),
|
||
message_id,
|
||
)?,
|
||
None => None,
|
||
};
|
||
let appended = if let Some(index) = matched_index {
|
||
let existing = &records[index];
|
||
if existing.role != role || existing.content != content {
|
||
return Err(format!(
|
||
"对话 messageId 冲突:{} 在 agent={} session={} 下对应不同 role/content",
|
||
message_id.as_deref().unwrap_or_default(),
|
||
normalized_agent_id.as_deref().unwrap_or("project"),
|
||
normalized_session_id.as_deref().unwrap_or("project")
|
||
));
|
||
}
|
||
false
|
||
} else {
|
||
if archived {
|
||
return Err(format!(
|
||
"Agent Session 已归档,只能读取:{}",
|
||
normalized_session_id.as_deref().unwrap_or_default()
|
||
));
|
||
}
|
||
append_jsonl_line_unlocked(&path, &line, "对话记录")?;
|
||
true
|
||
};
|
||
if appended || message_id.is_some() {
|
||
let mut audit_record = serde_json::json!({
|
||
"recordType": "conversation.message",
|
||
"agentId": normalized_agent_id.as_deref(),
|
||
"sessionId": normalized_session_id.as_deref(),
|
||
"role": role,
|
||
"path": relative_project_path(root, &path)?,
|
||
});
|
||
if let (Some(message_id), Some(object)) =
|
||
(message_id.as_deref(), audit_record.as_object_mut())
|
||
{
|
||
object.insert(
|
||
"messageId".to_string(),
|
||
serde_json::Value::String(message_id.to_string()),
|
||
);
|
||
if let Some(finalization_id) = finalization_id {
|
||
object.insert(
|
||
"finalizationId".to_string(),
|
||
serde_json::Value::String(finalization_id.to_string()),
|
||
);
|
||
}
|
||
}
|
||
if let Some(message_id) = message_id.as_deref() {
|
||
ensure_conversation_message_audit_at(
|
||
root,
|
||
normalized_agent_id.as_deref(),
|
||
normalized_session_id.as_deref(),
|
||
message_id,
|
||
audit_record,
|
||
)?;
|
||
} else {
|
||
append_agent_db_record(root, audit_record)?;
|
||
}
|
||
}
|
||
if message_id.is_none() {
|
||
records = read_persisted_local_conversation_records_unlocked(&path)?;
|
||
} else if appended {
|
||
records.push(record);
|
||
}
|
||
if let (Some(agent_id), Some(session_id)) = (
|
||
normalized_agent_id.as_deref(),
|
||
normalized_session_id.as_deref(),
|
||
) {
|
||
touch_agent_conversation_session_at(
|
||
root,
|
||
agent_id,
|
||
session_id,
|
||
records.len() as u64,
|
||
appended,
|
||
)?;
|
||
}
|
||
Ok((
|
||
local_conversation_result_from_persisted_records(
|
||
&path,
|
||
normalized_agent_id,
|
||
normalized_session_id,
|
||
&records,
|
||
),
|
||
appended,
|
||
))
|
||
}
|
||
|
||
#[cfg(test)]
|
||
fn take_local_conversation_append_failure_injection(root: &Path, role: &str) -> Result<(), String> {
|
||
let path = root.join(".agent/runtime/test-fail-next-conversation-append");
|
||
match fs::read_to_string(&path) {
|
||
Ok(expected_role) if expected_role.trim() == role => {
|
||
fs::remove_file(&path)
|
||
.map_err(|error| format!("清理对话写入测试失败注入标记失败:{error}"))?;
|
||
Err(format!("测试注入 {role} 对话写入失败"))
|
||
}
|
||
Ok(_) => Ok(()),
|
||
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
|
||
Err(error) => Err(format!("读取对话写入测试失败注入标记失败:{error}")),
|
||
}
|
||
}
|
||
|
||
pub(crate) fn append_local_conversation_message_for_session_at(
|
||
root: &Path,
|
||
agent_id: Option<&str>,
|
||
session_id: Option<&str>,
|
||
message: LocalConversationMessage,
|
||
) -> Result<LocalConversationResult, String> {
|
||
append_local_conversation_message_for_session_internal_at(
|
||
root, agent_id, session_id, message, None, None,
|
||
)
|
||
.map(|(conversation, _)| conversation)
|
||
}
|
||
|
||
pub(crate) fn append_local_conversation_message_for_session_idempotent_at(
|
||
root: &Path,
|
||
agent_id: Option<&str>,
|
||
session_id: Option<&str>,
|
||
message: LocalConversationMessage,
|
||
message_id: &str,
|
||
) -> Result<LocalConversationResult, String> {
|
||
append_local_conversation_message_for_session_internal_at(
|
||
root,
|
||
agent_id,
|
||
session_id,
|
||
message,
|
||
Some(message_id),
|
||
None,
|
||
)
|
||
.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,
|
||
)
|
||
}
|
||
|
||
pub(crate) fn append_local_conversation_message_for_session_idempotent_with_finalization_at(
|
||
root: &Path,
|
||
agent_id: Option<&str>,
|
||
session_id: Option<&str>,
|
||
message: LocalConversationMessage,
|
||
message_id: &str,
|
||
finalization_id: &str,
|
||
) -> Result<LocalConversationResult, String> {
|
||
append_local_conversation_message_for_session_internal_at(
|
||
root,
|
||
agent_id,
|
||
session_id,
|
||
message,
|
||
Some(message_id),
|
||
Some(finalization_id),
|
||
)
|
||
.map(|(conversation, _)| conversation)
|
||
}
|
||
|
||
pub(crate) fn append_local_conversation_message_at(
|
||
root: &Path,
|
||
agent_id: Option<&str>,
|
||
message: LocalConversationMessage,
|
||
) -> Result<LocalConversationResult, String> {
|
||
append_local_conversation_message_for_session_at(root, agent_id, None, message)
|
||
}
|
||
|
||
pub(crate) fn conversation_file_path_for_session(
|
||
root: &Path,
|
||
agent_id: Option<&str>,
|
||
session_id: Option<&str>,
|
||
) -> Result<(PathBuf, Option<String>, Option<String>), String> {
|
||
validate_project_root(root)?;
|
||
let Some(agent_id) = agent_id.map(str::trim).filter(|value| !value.is_empty()) else {
|
||
if session_id
|
||
.map(str::trim)
|
||
.is_some_and(|value| !value.is_empty())
|
||
{
|
||
return Err("项目主对话不接受 sessionId".to_string());
|
||
}
|
||
return Ok((root.join(".agent/conversations/project.jsonl"), None, None));
|
||
};
|
||
let agent_id = normalize_conversation_agent_id(agent_id)?;
|
||
let session = resolve_agent_conversation_session_at(root, &agent_id, session_id)?;
|
||
Ok((
|
||
conversation_file_path_for_resolved_session(root, &agent_id, &session.session_id),
|
||
Some(agent_id),
|
||
Some(session.session_id),
|
||
))
|
||
}
|
||
|
||
pub(crate) fn conversation_file_path(
|
||
root: &Path,
|
||
agent_id: Option<&str>,
|
||
) -> Result<(PathBuf, Option<String>), String> {
|
||
let (path, agent_id, _session_id) = conversation_file_path_for_session(root, agent_id, None)?;
|
||
Ok((path, agent_id))
|
||
}
|
||
|
||
pub(crate) fn normalize_conversation_agent_id(agent_id: &str) -> Result<String, String> {
|
||
if agent_id.is_empty()
|
||
|| agent_id.contains("..")
|
||
|| agent_id
|
||
.chars()
|
||
.any(|ch| !(ch.is_ascii_alphanumeric() || ch == '-' || ch == '_'))
|
||
{
|
||
return Err("agent id 只能包含 ASCII 字母、数字、短横线和下划线".to_string());
|
||
}
|
||
Ok(agent_id.to_string())
|
||
}
|
||
|
||
pub(crate) fn append_markdown_entry(
|
||
path: &Path,
|
||
header: &str,
|
||
entry: &str,
|
||
error_label: &str,
|
||
) -> Result<(), String> {
|
||
if let Some(parent) = path.parent() {
|
||
ensure_game_creator_private_directory_tree(parent, error_label)?;
|
||
prepare_game_creator_private_path_for_read(parent, true, error_label)?;
|
||
}
|
||
prepare_game_creator_private_path_for_read(path, false, error_label)?;
|
||
let needs_header = fs::metadata(path)
|
||
.map(|metadata| metadata.len() == 0)
|
||
.unwrap_or(true);
|
||
let bytes = if needs_header {
|
||
let mut bytes = Vec::with_capacity(header.len() + entry.len());
|
||
bytes.extend_from_slice(header.as_bytes());
|
||
bytes.extend_from_slice(entry.as_bytes());
|
||
bytes
|
||
} else {
|
||
entry.as_bytes().to_vec()
|
||
};
|
||
append_game_creator_private_file(path, &bytes, error_label)
|
||
}
|
||
|
||
pub(crate) fn append_local_permission_log_at(
|
||
root: &Path,
|
||
event: &str,
|
||
command_id: &str,
|
||
) -> Result<(), String> {
|
||
validate_project_root(root)?;
|
||
if !matches!(
|
||
event,
|
||
"permission.pending" | "permission.confirm" | "permission.cancel" | "command.auto"
|
||
) {
|
||
return Err("不支持的命令日志事件".to_string());
|
||
}
|
||
let Some(command) = GAME_CREATION_APP_COMMANDS
|
||
.iter()
|
||
.find(|command| command.id == command_id)
|
||
else {
|
||
return Err("不支持的内置命令".to_string());
|
||
};
|
||
if event == "command.auto" && command.permission != GameCreationAppPermission::Auto {
|
||
return Err("自动命令日志只能记录 auto 权限命令".to_string());
|
||
}
|
||
|
||
let log_path = root.join(".agent/logs/command.log");
|
||
if let Some(parent) = log_path.parent() {
|
||
ensure_game_creator_private_directory_tree(parent, "命令日志目录")?;
|
||
prepare_game_creator_private_path_for_read(parent, true, "命令日志目录")?;
|
||
}
|
||
prepare_game_creator_private_path_for_read(&log_path, false, "命令日志")?;
|
||
let line = format!("{} {event} {command_id}\n", unix_timestamp());
|
||
append_game_creator_private_file(&log_path, line.as_bytes(), "命令日志")
|
||
}
|
||
|
||
#[cfg(test)]
|
||
mod tests;
|