隔离Direct埋点状态锁并修正分页注释
Project CI / AI game creator shell Rust crates (pull_request) Successful in 1m19s
Project CI / AI game creator shell Rust smoke (pull_request) Successful in 1m54s
Project CI / Backend tests (pull_request) Successful in 4m48s
Project CI / Native shell tests (pull_request) Successful in 5m55s
Project CI / AI game creator shell Rust lane 2/2 (pull_request) Successful in 7m52s
Project CI / Frontend tests (pull_request) Successful in 2m9s
Project CI / AI game creator shell Rust lane 1/2 (pull_request) Successful in 8m34s
Project CI / AI game creator shell web tests (pull_request) Successful in 1m20s
Project CI / Repository checks (pull_request) Successful in 1m43s

采集身份与成果编号改用独立纯内存锁,避免业务落盘锁竞争导致漏采
补充锁竞争、原用户归属与恢复边界回归测试
修正后台分页注释,同步埋点方案与共享决策
This commit is contained in:
2026-09-22 15:51:48 +00:00
parent bc8e25c0e3
commit 069dec7bdd
5 changed files with 164 additions and 23 deletions
@@ -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<u64>,
started: Instant,
initial_elapsed_ms: u64,
lease_started: BTreeMap<String, Instant>,
@@ -113,6 +108,17 @@ struct SessionData {
elapsed_offset_ms: u64,
}
// 与业务持久化锁分离;只保存内存状态,持锁期间不执行 I/O 或投递事件。
struct SessionAnalytics {
project_id: String,
route: Option<crate::analytics::contract::Route>,
capture: Option<(
crate::analytics::contract::Context,
crate::analytics::store::AnalyticsWriter,
)>,
output_revision: Option<u64>,
}
#[derive(Clone, PartialEq, Eq)]
struct CodexExecutorIdentity {
path: PathBuf,
@@ -165,6 +171,7 @@ pub(super) struct ExecutionSession {
state_path: PathBuf,
_owner: File,
data: Mutex<SessionData>,
analytics: Mutex<SessionAnalytics>,
changed: tokio::sync::watch::Sender<u64>,
cancellation: Arc<std::sync::atomic::AtomicBool>,
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<String> {
self.data
.try_lock()
self.analytics
.lock()
.ok()?
.analytics_output_revision
.output_revision
.map(|revision| revision.to_string())
}
@@ -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(&current.context.editor_session_id)
.join("batches");
let deadline = Instant::now() + std::time::Duration::from_secs(5);
loop {
let events: Vec<Value> = 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::<Value>(line).unwrap())
.collect::<Vec<_>>()
})
.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<ExecutionSession>) {
@@ -1,5 +1,11 @@
# 决策记录
## 2026-09-22 Direct 埋点与业务持久化锁隔离
- Direct 采集身份和最新成果编号改由独立纯内存状态保存,初始化时从最终执行账本冻结项目与原 run 身份;成果采集、预览采集上下文和终态成果读取不再争用业务落盘锁。
- 内存锁释放后才投递事件,不等待文件 I/O、后台队列或网络,不新增用户报错。缺少 run 元数据不套用当前用户身份,恢复不补造历史成果;退出仍尽力封存,不增加退出等待。
- 修正后台埋点查询 DTO 的分页注释:按 `(event_time, event_id)` 倒序,入库时间仅限制快照;接口和查询行为不变。
## 2026-09-22 新项目埋点资格支持有限恢复
- 真实创建成功先登记有界进程内待办,不依赖埋点服务或身份快照;后台服务就绪及后续真实受理可重试初始化资格。资格文件格式、每项目最多一次首次提交和上传合同不变。
@@ -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:<session_id>:<target_phase>,使用真实策划会话 ID 和本次实际进入的目标阶段;source=design_agentrevision_source=agentchange_kind=design_documentfiles_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 的源码入口。本文新增合同和参数均以草案标识,不作为已经上线的事实。
@@ -1497,7 +1497,8 @@ pub struct AdminAgcModelCatalog {
pub models: Vec<AdminAgcModel>,
}
// 客户端埋点明细;时间范围按发生时间,分页按首次入库时间
// 客户端埋点明细;时间范围按发生时间,按 (event_time, event_id) 倒序分页。
// received_at 仅用于固定分页快照的入库时间上界。
#[derive(Clone, Debug, Default, serde::Deserialize, serde::Serialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
pub struct AdminAgcTrackingEventListQuery {