From 069dec7bdd8e9c156f2701275533ddff12f9074b Mon Sep 17 00:00:00 2001 From: Linghong Date: Tue, 22 Sep 2026 15:51:48 +0000 Subject: [PATCH] =?UTF-8?q?=E9=9A=94=E7=A6=BBDirect=E5=9F=8B=E7=82=B9?= =?UTF-8?q?=E7=8A=B6=E6=80=81=E9=94=81=E5=B9=B6=E4=BF=AE=E6=AD=A3=E5=88=86?= =?UTF-8?q?=E9=A1=B5=E6=B3=A8=E9=87=8A?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 采集身份与成果编号改用独立纯内存锁,避免业务落盘锁竞争导致漏采 补充锁竞争、原用户归属与恢复边界回归测试 修正后台分页注释,同步埋点方案与共享决策 --- .../src-tauri/src/agent/direct_execution.rs | 55 ++++---- .../src/agent/direct_execution/tests.rs | 118 ++++++++++++++++++ .../shared-memory/decision-log.md | 6 + ...案】客户端本地埋点与主站入库契约-2026-09-21.md | 5 +- .../crates/shared-contracts/src/admin.rs | 3 +- 5 files changed, 164 insertions(+), 23 deletions(-) diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_execution.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_execution.rs index af3a03f1a..a31c360f1 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_execution.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_execution.rs @@ -99,11 +99,6 @@ pub(super) struct ExecutionLedger { struct SessionData { ledger: ExecutionLedger, - analytics_capture: Option<( - crate::analytics::contract::Context, - crate::analytics::store::AnalyticsWriter, - )>, - analytics_output_revision: Option, started: Instant, initial_elapsed_ms: u64, lease_started: BTreeMap, @@ -113,6 +108,17 @@ struct SessionData { elapsed_offset_ms: u64, } +// 与业务持久化锁分离;只保存内存状态,持锁期间不执行 I/O 或投递事件。 +struct SessionAnalytics { + project_id: String, + route: Option, + capture: Option<( + crate::analytics::contract::Context, + crate::analytics::store::AnalyticsWriter, + )>, + output_revision: Option, +} + #[derive(Clone, PartialEq, Eq)] struct CodexExecutorIdentity { path: PathBuf, @@ -165,6 +171,7 @@ pub(super) struct ExecutionSession { state_path: PathBuf, _owner: File, data: Mutex, + analytics: Mutex, changed: tokio::sync::watch::Sender, cancellation: Arc, abort_requested: std::sync::atomic::AtomicBool, @@ -662,10 +669,17 @@ pub(super) fn open_with_analytics_at( newly_accepted, state_path, _owner: owner, + analytics: Mutex::new(SessionAnalytics { + project_id: ledger.project_id.clone(), + route: ledger + .analytics_run + .as_ref() + .map(|run| run.context.route.clone()), + capture: None, + output_revision: None, + }), data: Mutex::new(SessionData { ledger, - analytics_capture: None, - analytics_output_revision: None, started: Instant::now(), initial_elapsed_ms, lease_started: BTreeMap::new(), @@ -824,12 +838,12 @@ impl ExecutionSession { crate::analytics::store::AnalyticsWriter, )>, ) { - let Ok(mut data) = self.data.try_lock() else { + let Ok(mut analytics) = self.analytics.lock() else { return; }; - data.analytics_capture = capture.and_then(|(mut context, writer)| { + analytics.capture = capture.and_then(|(mut context, writer)| { // 恢复或账号切换后仍归属于真实受理的原 run。 - context.route = data.ledger.analytics_run.as_ref()?.context.route.clone(); + context.route = analytics.route.clone()?; Some((context, writer)) }); } @@ -840,7 +854,7 @@ impl ExecutionSession { crate::analytics::contract::Context, crate::analytics::store::AnalyticsWriter, )> { - self.data.try_lock().ok()?.analytics_capture.clone() + self.analytics.lock().ok()?.capture.clone() } pub(super) fn record_analytics_revision( @@ -853,17 +867,16 @@ impl ExecutionSession { if files_changed_count == 0 { return; } - let Ok(mut data) = self.data.try_lock() else { + let Ok(mut analytics) = self.analytics.lock() else { return; }; - if data.ledger.analytics_run.is_none() { + if analytics.route.is_none() { return; } - data.analytics_output_revision = - Some(data.analytics_output_revision.unwrap_or(0).max(revision)); - let capture = data.analytics_capture.clone(); - let project_id = data.ledger.project_id.clone(); - drop(data); + analytics.output_revision = Some(analytics.output_revision.unwrap_or(0).max(revision)); + let capture = analytics.capture.clone(); + let project_id = analytics.project_id.clone(); + drop(analytics); crate::analytics::project::revision( capture, &project_id, @@ -878,10 +891,10 @@ impl ExecutionSession { } pub(super) fn analytics_output_revision(&self) -> Option { - self.data - .try_lock() + self.analytics + .lock() .ok()? - .analytics_output_revision + .output_revision .map(|revision| revision.to_string()) } diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_execution/tests.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_execution/tests.rs index 3f190af4e..9f8ddb1f6 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_execution/tests.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_execution/tests.rs @@ -13,6 +13,113 @@ fn analytics_metadata(user: &str) -> crate::analytics::run::Metadata { ) } +#[test] +fn analytics_survives_business_lock_contention_and_preserves_replayed_run_identity() { + use crate::analytics::{contract::ChangeKind, store::AnalyticsWriter}; + let temp = tempfile::tempdir().unwrap(); + let root = temp.path().join("project"); + let host = temp.path().join("host"); + crate::init_local_game_project_at(&root, "analytics-lock", "采集锁隔离").unwrap(); + let original = open_with_analytics_at( + &host, + &root, + "turn", + &hash(b"request"), + false, + &Default::default(), + Some(analytics_metadata("A")), + ) + .unwrap(); + drop(original); + let current = analytics_metadata("B"); + let session = open_with_analytics_at( + &host, + &root, + "turn", + &hash(b"request"), + false, + &Default::default(), + Some(current.clone()), + ) + .unwrap(); + let config = temp.path().join("config"); + std::fs::create_dir_all(&config).unwrap(); + let writer = AnalyticsWriter::start(config.clone(), current.context.editor_session_id.clone()); + let capture = (current.context.clone(), writer.clone()); + // 模拟业务提交长期占锁;采集必须在释放该锁之前完成。 + let business_lock = session.data.lock().unwrap(); + let task_session = session.clone(); + let (sender, receiver) = std::sync::mpsc::channel(); + let worker = std::thread::spawn(move || { + task_session.set_analytics_capture(Some(capture)); + task_session.record_analytics_revision(7, ChangeKind::Code, 1); + task_session.record_analytics_revision(5, ChangeKind::Code, 1); + task_session.record_analytics_revision(99, ChangeKind::Code, 0); + sender + .send(( + task_session.analytics_capture(), + task_session.analytics_output_revision(), + )) + .unwrap(); + }); + let result = receiver.recv_timeout(std::time::Duration::from_secs(5)); + // 即使回归成等待业务锁,也先释放锁和回收线程,让测试明确失败而非挂死。 + drop(business_lock); + worker.join().unwrap(); + let (capture, revision) = result.expect("采集不得等待业务持久化锁"); + let (context, _) = capture.expect("锁竞争不得丢失采集身份"); + assert_eq!(context.route.user_id.as_deref(), Some("A")); + assert_eq!(context.editor_session_id, current.context.editor_session_id); + assert_eq!(revision.as_deref(), Some("7")); + assert!(writer.flush()); + let batches = config + .join("analytics/instances") + .join(¤t.context.editor_session_id) + .join("batches"); + let deadline = Instant::now() + std::time::Duration::from_secs(5); + loop { + let events: Vec = std::fs::read_dir(&batches) + .into_iter() + .flatten() + .flatten() + .filter(|entry| !entry.file_name().to_string_lossy().starts_with('.')) + .filter_map(|entry| std::fs::read_to_string(entry.path().join("events.jsonl")).ok()) + .flat_map(|text| { + text.lines() + .map(|line| serde_json::from_str::(line).unwrap()) + .collect::>() + }) + .collect(); + if events.len() == 2 { + for (event, expected_revision) in events.iter().zip(["7", "5"]) { + assert_eq!(event["event_name"], "project_revision_created"); + assert_eq!(event["user_id"], "A"); + assert_eq!(event["project_id"], "analytics-lock"); + assert_eq!(event["properties"]["revision_id"], expected_revision); + } + break; + } + assert!(Instant::now() < deadline, "成果事件未落盘"); + std::thread::sleep(std::time::Duration::from_millis(5)); + } + drop(session); + let resumed = open_with_analytics_at( + &host, + &root, + "turn", + &hash(b"request"), + false, + &Default::default(), + Some(current), + ) + .unwrap(); + assert_eq!( + resumed.analytics_output_revision(), + None, + "恢复不补造历史成果编号" + ); +} + #[test] fn run_metadata_is_persisted_with_new_ledger_and_replay_keeps_original_identity() { let temp = tempfile::tempdir().unwrap(); @@ -67,6 +174,17 @@ fn legacy_run_without_metadata_is_not_assigned_current_users_identity() { ) .unwrap(); assert!(replay.snapshot().unwrap().analytics_run.is_none()); + let config = temp.path().join("config"); + std::fs::create_dir_all(&config).unwrap(); + let current = analytics_metadata("B"); + let writer = crate::analytics::store::AnalyticsWriter::start( + config, + current.context.editor_session_id.clone(), + ); + replay.set_analytics_capture(Some((current.context, writer))); + replay.record_analytics_revision(1, crate::analytics::contract::ChangeKind::Code, 1); + assert!(replay.analytics_capture().is_none()); + assert!(replay.analytics_output_revision().is_none()); } fn fixture(config: DirectValidationConfig) -> (tempfile::TempDir, Arc) { diff --git a/docs/project-memory/shared-memory/decision-log.md b/docs/project-memory/shared-memory/decision-log.md index 17696b9c7..4d95ac128 100644 --- a/docs/project-memory/shared-memory/decision-log.md +++ b/docs/project-memory/shared-memory/decision-log.md @@ -1,5 +1,11 @@ # 决策记录 +## 2026-09-22 Direct 埋点与业务持久化锁隔离 + +- Direct 采集身份和最新成果编号改由独立纯内存状态保存,初始化时从最终执行账本冻结项目与原 run 身份;成果采集、预览采集上下文和终态成果读取不再争用业务落盘锁。 +- 内存锁释放后才投递事件,不等待文件 I/O、后台队列或网络,不新增用户报错。缺少 run 元数据不套用当前用户身份,恢复不补造历史成果;退出仍尽力封存,不增加退出等待。 +- 修正后台埋点查询 DTO 的分页注释:按 `(event_time, event_id)` 倒序,入库时间仅限制快照;接口和查询行为不变。 + ## 2026-09-22 新项目埋点资格支持有限恢复 - 真实创建成功先登记有界进程内待办,不依赖埋点服务或身份快照;后台服务就绪及后续真实受理可重试初始化资格。资格文件格式、每项目最多一次首次提交和上传合同不变。 diff --git a/docs/technical/【技术方案】客户端本地埋点与主站入库契约-2026-09-21.md b/docs/technical/【技术方案】客户端本地埋点与主站入库契约-2026-09-21.md index a2373fb3e..3b6b8be97 100644 --- a/docs/technical/【技术方案】客户端本地埋点与主站入库契约-2026-09-21.md +++ b/docs/technical/【技术方案】客户端本地埋点与主站入库契约-2026-09-21.md @@ -1,6 +1,6 @@ # 客户端本地埋点与主站入库契约 -Version: 0.30 +Version: 0.31 Status: 本地采集及上传、入库、成功清理、后台查询已实现并完成本地隔离环境验收;未部署生产 Date: 2026-09-21 需求来源:[Game Agent 埋点设计原始方案](./【需求来源】GameAgent埋点设计原始方案-2026-09-05.md);原始方案与后续已确认决策有差异时,以本文为准。 @@ -262,6 +262,7 @@ Direct 的自动登录刷新重试仍属同一个 run。原生单次调用结束 - 当前 Direct 的 outputs_changed 与 manifest_requires_sync 要区分;不能给 revision 递增函数加一个无条件事件后就宣称满足语义。 - Direct 的宿主 file.write 与正式补丁事务使用持锁目标内容比较和本次已提交项目 revision,成功且内容确实变化时记 source=direct、revision_source=agent。运行期间的全局 fingerprint 和写许可本身不证明本 run 改了文件;末尾投影登记不再重复统计工具已经提交的成果,也不将纯登记修复计为新成果。未经过可观察宿主提交的外部/命令行写入不推断事件。 - Direct run 的 output_change_detected=true 与 revision_id 来自本 run 已确认的宿主成果提交;保存当前内存中最后一个真实成果 revision,不为埋点新增业务文件版本或同步写盘。没有观察到成果、恢复后缺少该证据时为 null,不能推断 false。后续 run 失败不撤销已提交成果事件。 +- Direct 的采集身份与最新成果编号使用独立的纯内存状态,项目 ID 和原 run 身份在执行会话创建时从最终账本冻结;采集不争用业务持久化锁,不因业务写入或 tick 占锁而漏记。内存锁只保护采集状态读写,事件投递在释放锁后进行,不等待磁盘、网络或队列,不向用户报错。恢复仍使用原 run 身份,但不补造此前未保留的成果编号;旧账本缺少 run 元数据时不套用当前用户身份。 - 宿主文件通道仅对已识别成果类型采集:策划成果目录、UI 目录以及支持的代码/媒体扩展名;私有控制面、隐藏记录、构建缓存和类型未知的文件不据此统计。该口径不承诺观察全部磁盘文件;恢复后的未知成果证据不补填。 - Design Agent:仅审批通过后实际进入下一阶段且成功持久化时记录;不严格检查文档版本或文件差异,不增加业务文件 revision。revision_id=design::,使用真实策划会话 ID 和本次实际进入的目标阶段;source=design_agent,revision_source=agent,change_kind=design_document,files_changed_count 省略。该标识只表示阶段推进事实,不代表文件版本。 - 策划同会话同目标阶段只记一次;同阶段内文件写入、普通对话、工具执行、提交审批、审批拒绝、等待澄清、没有实际阶段变化的批准,以及持久化失败均不产生该成果事件。项目重开、会话恢复或读取现有阶段不补历史。 @@ -501,6 +502,8 @@ session.json 是本地恢复元数据,不是待上传事件;事件文件不 ## 12. 实施与验证状态 +2026-09-22 Direct 采集锁隔离验证:原生 `analytics` 定向测试 60 项通过、1 项既有跳过,`agent::direct_execution::tests` 24 项全部通过(两组有 1 项重叠)。新增线程用例持续占用业务锁,验证采集身份设置、成果记录、上下文与最新成果读取在释放业务锁前完成,并实际读回原用户的成果事件;另验证恢复不补造历史成果编号、旧账本缺少元数据不套用当前用户。独立代码复核及格式、编码、文档索引、diff 检查通过。分页 DTO 仅修正注释,查询行为不变;退出机制未改,未运行完整 GUI/Provider 联调。 + 2026-09-22 新项目资格有限恢复验证:原生 `analytics` 定向测试 59 项通过、1 项既有跳过;覆盖创建时无埋点上下文、创建通知队列拒绝、独立锁竞争和临时文件系统故障后的真实受理恢复,以及资格成功落盘后不重授、进程内待办容量限制。既有损坏/备份/身份不符、跨 writer 去重测试继续通过。命令配置、Rust 格式、编码、文档索引与 diff 检查通过;独立代码复核确认业务线程不执行资格文件 I/O,内存锁不覆盖后台 I/O。未执行完整 GUI/Provider 联调或生产发布。 已做:原产品方案逐项对照;核对现役 Direct/Design 持久化、项目创建、revision、预览、UI 保存与 checkpoint 的源码入口。本文新增合同和参数均以草案标识,不作为已经上线的事实。 diff --git a/server-rs/crates/shared-contracts/src/admin.rs b/server-rs/crates/shared-contracts/src/admin.rs index ede7b790e..f9730feb4 100644 --- a/server-rs/crates/shared-contracts/src/admin.rs +++ b/server-rs/crates/shared-contracts/src/admin.rs @@ -1497,7 +1497,8 @@ pub struct AdminAgcModelCatalog { pub models: Vec, } -// 客户端埋点明细;时间范围按发生时间,分页按首次入库时间。 +// 客户端埋点明细;时间范围按发生时间,按 (event_time, event_id) 倒序分页。 +// received_at 仅用于固定分页快照的入库时间上界。 #[derive(Clone, Debug, Default, serde::Deserialize, serde::Serialize)] #[serde(rename_all = "camelCase", deny_unknown_fields)] pub struct AdminAgcTrackingEventListQuery {