修复新项目埋点资格丢失后的持续漏计

创建成功先登记有界内存待办,不依赖埋点服务与身份快照就绪
后台初始化与后续真实提交重试资格,保持静默且不等待业务磁盘操作
补充恢复与容量边界测试,同步采集合同和共享决策
This commit is contained in:
2026-09-22 14:49:50 +00:00
parent 5f7c168841
commit 1e9a9c9674
6 changed files with 296 additions and 46 deletions
@@ -1940,7 +1940,7 @@ mod tests {
let (analytics_dir, context, writer) = analytics_capture();
concept_artifacts(&root);
let mut session = new_design_session(design_project_id(&root).unwrap(), "quality");
crate::analytics::goal::created(&writer, &root, &session.project_id);
crate::analytics::goal::created(Some(&writer), &root, &session.project_id);
let approval = submit_design_phase_for_approval(&root, &mut session).unwrap();
let old_phase = session.current_phase.clone();
assert!(
@@ -1979,7 +1979,7 @@ mod tests {
async fn design_first_submit_keeps_accepting_user_after_model_failure_and_replay() {
let (_temp, root, resources) = init_design_project();
let (analytics_dir, context, writer) = analytics_capture();
crate::analytics::goal::created(&writer, &root, &design_project_id(&root).unwrap());
crate::analytics::goal::created(Some(&writer), &root, &design_project_id(&root).unwrap());
let _fake = fake_provider::install(
vec![Err(platform_llm::LlmError::Upstream {
status_code: 500,
@@ -2209,7 +2209,7 @@ mod tests {
)
.await
.unwrap();
crate::analytics::goal::created(&writer, &root, &design_project_id(&root).unwrap());
crate::analytics::goal::created(Some(&writer), &root, &design_project_id(&root).unwrap());
continue_design_agent_with_capture_at(
&root,
&resources,
@@ -2230,7 +2230,7 @@ mod tests {
async fn design_invalid_input_preserves_qualification_for_later_accepting_user() {
let (_temp, root, resources) = init_design_project();
let (analytics_dir, context, writer) = analytics_capture();
crate::analytics::goal::created(&writer, &root, &design_project_id(&root).unwrap());
crate::analytics::goal::created(Some(&writer), &root, &design_project_id(&root).unwrap());
assert!(continue_design_agent_with_capture_at(
&root,
&resources,
@@ -3141,7 +3141,7 @@ mod tests {
.map(|item| item.request_id.as_str()),
Some(request.as_str())
);
crate::analytics::goal::created(&writer, &root, &persisted.project_id);
crate::analytics::goal::created(Some(&writer), &root, &persisted.project_id);
let rejected = decide_design_phase_with_capture_at(
&root,
&resources,
@@ -2,10 +2,54 @@
use super::contract::{Context, EmptyProperties, Event, EventData, Route, Source};
use super::store::AnalyticsWriter;
use serde::{Deserialize, Serialize};
use std::collections::HashSet;
use std::path::{Path, PathBuf};
use std::sync::{Mutex, OnceLock};
const MARKER: &str = ".agent/analytics-goal.json";
const MAX_MARKER_BYTES: usize = 4096;
const MAX_PENDING_PROJECTS: usize = 1024;
const MAX_PENDING_BYTES: usize = 1024 * 1024;
#[derive(Default)]
struct PendingGoals {
projects: HashSet<(PathBuf, String)>,
bytes: usize,
}
impl PendingGoals {
fn remember(&mut self, root: &Path, project_id: &str) {
let key = (root.to_path_buf(), project_id.to_owned());
let bytes = root.as_os_str().as_encoded_bytes().len() + project_id.len();
if self.projects.len() < MAX_PENDING_PROJECTS
&& self.bytes.saturating_add(bytes) <= MAX_PENDING_BYTES
&& self.projects.insert(key)
{
self.bytes += bytes;
}
}
fn forget(&mut self, root: &Path, project_id: &str) {
if let Some((root, project_id)) = self
.projects
.take(&(root.to_path_buf(), project_id.to_owned()))
{
self.bytes -= root.as_os_str().as_encoded_bytes().len() + project_id.len();
}
}
}
fn pending() -> &'static Mutex<PendingGoals> {
static PENDING: OnceLock<Mutex<PendingGoals>> = OnceLock::new();
PENDING.get_or_init(|| Mutex::new(PendingGoals::default()))
}
fn forget_pending(root: &Path, project_id: &str) {
pending()
.lock()
.unwrap_or_else(|error| error.into_inner())
.forget(root, project_id);
}
#[derive(Serialize)]
pub(super) enum Request {
@@ -52,11 +96,73 @@ struct Marker {
submitted: bool,
}
pub(crate) fn created(writer: &AnalyticsWriter, root: &Path, project_id: &str) {
writer.try_goal(Request::Created {
pub(crate) fn created(writer: Option<&AnalyticsWriter>, root: &Path, project_id: &str) {
let request = Request::Created {
root: root.into(),
project_id: project_id.into(),
});
};
if !request.validate() {
return;
}
// 只串行修改有界内存;任何路径检查、文件锁和落盘均留在后台。
pending()
.lock()
.unwrap_or_else(|error| error.into_inner())
.remember(root, project_id);
if let Some(writer) = writer {
writer.try_goal(request);
}
}
pub(super) fn retry_pending(writer: &AnalyticsWriter) {
let projects: Vec<_> = pending()
.lock()
.unwrap_or_else(|error| error.into_inner())
.projects
.iter()
.cloned()
.collect();
for (root, project_id) in projects {
// 投递失败不移除资格;后续真实受理仍可重试。
writer.try_goal(Request::Created { root, project_id });
}
}
// 调用方已取得项目独立埋点锁;不得在文件 I/O 期间持有 pending 内存锁。
fn initialize_pending(root: &Path, project_id: &str, path: &Path) -> Result<(), String> {
let known_created = pending()
.lock()
.unwrap_or_else(|error| error.into_inner())
.projects
.contains(&(root.to_path_buf(), project_id.to_owned()));
if !known_created {
return Ok(());
}
let backup = crate::agent::agent_runtime_json_sidecar_backup_path(path);
for candidate in [path, backup.as_path()] {
match std::fs::symlink_metadata(candidate) {
Ok(_) => {
// 已有主标记或恢复副本,无论完好与否都不重授资格。
forget_pending(root, project_id);
return Ok(());
}
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
Err(error) => return Err(error.to_string()),
}
}
crate::agent::write_agent_runtime_json_sidecar_with_max_bytes(
root,
MARKER,
"埋点目标资格",
&Marker {
schema_version: 1,
project_id: project_id.into(),
submitted: false,
},
MAX_MARKER_BYTES,
)?;
forget_pending(root, project_id);
Ok(())
}
pub(crate) fn accepted(
@@ -102,35 +208,17 @@ pub(super) fn process(request: Request) -> Result<Option<(Route, Event)>, String
};
let manifest_path = crate::project::resolve_local_project_path(root, ".agent/manifest.json")?;
if crate::project::read_manifest(&manifest_path)?.project_id != project_id {
forget_pending(root, project_id);
return Ok(None);
}
let path = root.join(MARKER);
initialize_pending(root, project_id, &path)?;
let path = crate::project::resolve_local_project_path(root, MARKER)?;
let primary = std::fs::symlink_metadata(&path);
match &request {
Request::Created { .. } => {
// 缺失主文件但存在恢复副本仍属未知,不能重新授予资格。
let backup = crate::agent::agent_runtime_json_sidecar_backup_path(&path);
if !matches!(primary, Err(ref e) if e.kind() == std::io::ErrorKind::NotFound)
|| !matches!(std::fs::symlink_metadata(backup), Err(ref e) if e.kind() == std::io::ErrorKind::NotFound)
{
return Ok(None);
}
crate::agent::write_agent_runtime_json_sidecar_with_max_bytes(
root,
MARKER,
"埋点目标资格",
&Marker {
schema_version: 1,
project_id: project_id.into(),
submitted: false,
},
MAX_MARKER_BYTES,
)?;
Ok(None)
}
Request::Created { .. } => Ok(None),
Request::Accepted { .. } => {
// 明确要求主文件存在,不启用 sidecar 的 previous 自动回退。
let Ok(metadata) = primary else {
let Ok(metadata) = std::fs::symlink_metadata(&path) else {
return Ok(None);
};
if !metadata.is_file() || metadata.file_type().is_symlink() {
@@ -164,3 +252,34 @@ pub(super) fn process(request: Request) -> Result<Option<(Route, Event)>, String
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn pending_goals_are_bounded_and_duplicate_creation_does_not_charge_bytes() {
let mut pending = PendingGoals::default();
let root = std::env::temp_dir().join("analytics-pending-bound");
pending.remember(&root, "project");
let bytes = pending.bytes;
pending.remember(&root, "project");
assert_eq!(pending.bytes, bytes);
pending.forget(&root, "other-project");
assert_eq!(pending.bytes, bytes);
pending.forget(&root, "project");
assert_eq!(pending.bytes, 0);
for index in 0..MAX_PENDING_PROJECTS + 1 {
pending.remember(&root.join(index.to_string()), "project");
}
assert_eq!(pending.projects.len(), MAX_PENDING_PROJECTS);
let mut pending = PendingGoals::default();
let large_root = root.join("x".repeat(30_000));
for index in 0..100 {
pending.remember(&large_root, &index.to_string());
}
assert!(pending.bytes <= MAX_PENDING_BYTES);
assert!(pending.projects.len() < 100);
}
}
@@ -61,7 +61,7 @@ impl LifecycleFixture {
pub(crate) fn create_and_open(&mut self, root: &std::path::Path, project_id: &str) {
crate::init_local_game_project_at(root, project_id, "采集完整链路测试").unwrap();
super::goal::created(&self.0.writer, root, project_id);
super::goal::created(Some(&self.0.writer), root, project_id);
self.0.emit(
EventData::ProjectCreateSuccess(ProjectCreated {
creation_source: CreationSource::HomeGame,
@@ -383,6 +383,9 @@ pub(crate) fn initialize(app: tauri::AppHandle, config_dir: PathBuf, version: St
initialize_with_identity(config_dir, version, route, sequence);
}
});
if let Some(gui) = GUI.get() {
super::goal::retry_pending(&gui.writer);
}
let handle = app.clone();
let _ = app.run_on_main_thread(move || observe_windows(&handle, None, false, true));
});
@@ -461,13 +464,13 @@ pub(crate) fn created(
source: CreationSource,
template_id: Option<String>,
) {
super::goal::created(GUI.get().map(|gui| &gui.writer), project_root, &project_id);
let (Some(context), Some(gui)) = (context, GUI.get()) else {
return;
};
if !valid_context(&context) {
return;
}
super::goal::created(&gui.writer, project_root, &project_id);
if let Ok(event) = context.capture(
EventData::ProjectCreateSuccess(ProjectCreated {
creation_source: source,
@@ -631,7 +631,7 @@ fn project_goal_is_once_across_sources_writers_and_batch_retention_with_frozen_i
let (_project_dir, root) = goal_project();
let config = tempfile::tempdir().unwrap();
let (a, writer) = goal_writer(config.path(), "A");
goal::created(&writer, &root, "goal-project");
goal::created(Some(&writer), &root, "goal-project");
goal::accepted(
Some((a.clone(), writer.clone())),
&root,
@@ -659,7 +659,7 @@ fn project_goal_is_once_across_sources_writers_and_batch_retention_with_frozen_i
remove_batch(&batch.path).unwrap();
}
let (restarted, next) = goal_writer(config.path(), "B");
goal::created(&next, &root, "goal-project"); // 不覆盖 consumed
goal::created(Some(&next), &root, "goal-project"); // 不覆盖 consumed
goal::accepted(
Some((restarted.clone(), next.clone())),
&root,
@@ -692,7 +692,7 @@ fn goal_unknown_corrupt_backup_and_mismatched_projects_do_not_submit() {
.any(|e| e.event_name == "creative_task_submit"));
let marker = root.join(".agent/analytics-goal.json");
fs::write(&marker, b"broken").unwrap();
goal::created(&writer, &root, "goal-project");
goal::created(Some(&writer), &root, "goal-project");
accept("goal-project");
assert!(!drain_goal_writer(config.path(), &context, &writer)
.iter()
@@ -705,20 +705,121 @@ fn goal_unknown_corrupt_backup_and_mismatched_projects_do_not_submit() {
br#"{"schema_version":1,"project_id":"goal-project","submitted":false}"#,
)
.unwrap();
goal::created(&writer, &root, "goal-project");
goal::created(Some(&writer), &root, "goal-project");
accept("goal-project");
assert!(!drain_goal_writer(config.path(), &context, &writer)
.iter()
.any(|e| e.event_name == "creative_task_submit"));
assert!(!marker.exists());
fs::remove_file(backup).unwrap();
goal::created(&writer, &root, "goal-project");
goal::created(Some(&writer), &root, "goal-project");
accept("other-project");
assert!(!drain_goal_writer(config.path(), &context, &writer)
.iter()
.any(|e| e.event_name == "creative_task_submit"));
}
#[test]
fn goal_creation_without_analytics_context_recovers_on_later_acceptance() {
use super::super::{
contract::{CreationSource, Source},
goal, gui,
};
let (_project_dir, root) = goal_project();
gui::created(
None,
&root,
"goal-project".into(),
CreationSource::HomeGame,
None,
);
assert!(!root.join(".agent/analytics-goal.json").exists());
let config = tempfile::tempdir().unwrap();
let (context, writer) = goal_writer(config.path(), "A");
goal::accepted(
Some((context.clone(), writer.clone())),
&root,
"goal-project",
Source::Direct,
);
let events = drain_goal_writer(config.path(), &context, &writer);
assert_eq!(
events
.iter()
.filter(|event| event.event_name == "creative_task_submit")
.count(),
1
);
// 已落盘资格之后即使标记丢失,也不能凭旧的进程内待办重新授予。
let marker = root.join(".agent/analytics-goal.json");
fs::remove_file(&marker).unwrap();
let backup = crate::agent::agent_runtime_json_sidecar_backup_path(&marker);
if backup.exists() {
fs::remove_file(backup).unwrap();
}
goal::accepted(
Some((context.clone(), writer.clone())),
&root,
"goal-project",
Source::Direct,
);
let events = drain_goal_writer(config.path(), &context, &writer);
assert_eq!(
events
.iter()
.filter(|event| event.event_name == "creative_task_submit")
.count(),
1
);
assert!(!marker.exists());
}
#[test]
fn goal_creation_lock_and_io_failures_recover_on_later_acceptance() {
use super::super::{contract::Source, goal};
for io_failure in [false, true] {
let (_project_dir, root) = goal_project();
let config = tempfile::tempdir().unwrap();
let (context, writer) = goal_writer(config.path(), "A");
let lock_path = root.join(".agent/runtime/locks/analytics-goal.lock");
let lock = if io_failure {
// 用目录占据埋点锁文件,模拟临时文件系统故障,不改业务 manifest。
fs::create_dir_all(&lock_path).unwrap();
None
} else {
Some(
crate::agent::try_acquire_game_creator_agent_runtime_task_lock(
&root,
"analytics-goal",
)
.unwrap()
.unwrap(),
)
};
goal::created(Some(&writer), &root, "goal-project");
drain_goal_writer(config.path(), &context, &writer);
assert!(!root.join(".agent/analytics-goal.json").exists());
drop(lock);
if io_failure {
fs::remove_dir(&lock_path).unwrap();
}
goal::accepted(
Some((context.clone(), writer.clone())),
&root,
"goal-project",
Source::Direct,
);
let events = drain_goal_writer(config.path(), &context, &writer);
assert_eq!(
events
.iter()
.filter(|event| event.event_name == "creative_task_submit")
.count(),
1
);
}
}
#[test]
fn goal_queue_capacity_rejection_is_silent_and_does_not_charge_bytes() {
use super::super::goal;
@@ -729,22 +830,39 @@ fn goal_queue_capacity_rejection_is_silent_and_does_not_charge_bytes() {
counters: counters.clone(),
};
let (_dir, root) = goal_project();
goal::created(&writer, &root, "goal-project");
assert!(writer.flush()); // 占满队列,让创建通知从第一次起就无法投递。
goal::created(Some(&writer), &root, "goal-project");
let charged = counters.queued_bytes.load(Ordering::Relaxed);
assert!(charged > 0);
goal::created(&writer, &root, "goal-project");
assert_eq!(charged, 0);
goal::created(Some(&writer), &root, "goal-project");
assert_eq!(counters.queued_bytes.load(Ordering::Relaxed), charged);
assert_eq!(counters.dropped.load(Ordering::Relaxed), 1);
assert_eq!(counters.dropped.load(Ordering::Relaxed), 2);
assert!(!root.join(".agent/analytics-goal.json").exists());
drop(receiver);
counters
.queued_bytes
.store(MAX_BYTES as u64, Ordering::Relaxed);
goal::created(&writer, &root, "goal-project");
goal::created(Some(&writer), &root, "goal-project");
assert_eq!(
counters.queued_bytes.load(Ordering::Relaxed),
MAX_BYTES as u64
);
let config = tempfile::tempdir().unwrap();
let (context, recovered) = goal_writer(config.path(), "A");
goal::accepted(
Some((context.clone(), recovered.clone())),
&root,
"goal-project",
super::super::contract::Source::Direct,
);
let events = drain_goal_writer(config.path(), &context, &recovered);
assert_eq!(
events
.iter()
.filter(|event| event.event_name == "creative_task_submit")
.count(),
1
);
}
#[test]
@@ -754,7 +872,7 @@ fn later_real_acceptance_can_consume_after_lock_contention_and_two_writers_do_no
let config = tempfile::tempdir().unwrap();
let (a, first) = goal_writer(config.path(), "A");
let (b, second) = goal_writer(config.path(), "B");
goal::created(&first, &root, "goal-project");
goal::created(Some(&first), &root, "goal-project");
drain_goal_writer(config.path(), &a, &first);
let lock =
crate::agent::try_acquire_game_creator_agent_runtime_task_lock(&root, "analytics-goal")
@@ -1,5 +1,11 @@
# 决策记录
## 2026-09-22 新项目埋点资格支持有限恢复
- 真实创建成功先登记有界进程内待办,不依赖埋点服务或身份快照;后台服务就绪及后续真实受理可重试初始化资格。资格文件格式、每项目最多一次首次提交和上传合同不变。
- 业务线程仅更新内存和非阻塞投递;资格文件读写和独立锁操作仍在后台,不向用户报错或要求介入。队列满、暂时 I/O 失败和锁竞争保留待办;已有标记、备份或项目身份不符不重新授予资格。
- 内存最多 1024 项、路径与 ID 合计 1 MiB;超限或持久化前退出仍允许漏记,不补旧项目历史,不承诺零丢失。细节与验证见客户端埋点主规范。
## 2026-09-22 清理无调用方的 GUI 文件写入命令
- 确认 `write_local_project_file` 无现役前端或业务调用方后,删除命令、Tauri 注册及命令检查豁免;移除对应测试片段,保留 checkpoint、UI 保存和记忆写入的既有测试。
@@ -1,6 +1,6 @@
# 客户端本地埋点与主站入库契约
Version: 0.29
Version: 0.30
Status: 本地采集及上传、入库、成功清理、后台查询已实现并完成本地隔离环境验收;未部署生产
Date: 2026-09-21
需求来源:[Game Agent 埋点设计原始方案](./【需求来源】GameAgent埋点设计原始方案-2026-09-05.md);原始方案与后续已确认决策有差异时,以本文为准。
@@ -226,7 +226,9 @@ Direct 的自动登录刷新重试仍属同一个 run。原生单次调用结束
### 6.8 creative_task_submit
- 触发:本版实际新建项目的创作请求通过基本校验并被业务受理,且后台成功消费尚未提交的项目资格;不是创建项目、点击按钮、未提交草稿或请求重放。每项目最多采集一次。
- 项目资格只在真实创建成功后初始化,记录格式版本、project_id 与已提交状态,随项目保留,不进入 7 天批次清理。后台互斥消费后才投递事件;此前候选丢弃而资格尚未消费时,后续真实受理可记录。资格已消费但事件未落盘时允许漏记,不反向重建资格。资格缺失、损坏或身份不符均为未知
- 项目资格只在真实创建成功后初始化,记录格式版本、project_id 与已提交状态,随项目保留,不进入 7 天批次清理。创建成功先登记进程内待初始化资格,不依赖埋点服务或身份快照就绪;登记仅包含项目目录与 project_id,最多 1024 项且路径/ID 合计最多 1 MiB,超限静默放弃新增,不影响创建。后台互斥消费后才投递事件;此前候选丢弃而资格尚未消费时,后续真实受理可记录。
- 埋点服务就绪后尝试投递待办;创建通知和后续真实受理均可在现有后台写入器中重试初始化。不增加定时器或项目扫描,不在业务线程执行资格文件 I/O,也不等待队列、磁盘或网络;失败不弹窗、不返回业务错误、不要求用户介入。队列拒绝、独立锁竞争与暂时 I/O 失败保留待办;只有明确观察到本进程创建、主标记和恢复副本均不存在时才允许初始化。初始化成功、发现任意已有标记/恢复副本或 manifest 项目身份不符后移除待办,不重置已消费状态,不修复损坏标记。
- 没有待办依据的旧项目缺失资格仍为未知,不补历史。待办未持久化前进程退出或超出内存限额时允许漏记;资格已消费但事件未落盘时同样允许漏记,不反向重建资格。此恢复不承诺零丢失。
- 并发候选以成功持久消费资格的顺序决定保留哪次受理,不保证业务最早请求或最早 event_time;最终事件可能为零条。项目标记属于持久的观测状态,不保存用户身份、正文或路径,也不作为业务执行门禁。
- 公共字段:status=successproject_id、creative_task_id 必填;source=direct / design_agentrun/turn ID 均 null。
- properties:空对象。公共字段不重复,不增加 Prompt、长度、附件或任务分类。
@@ -499,6 +501,8 @@ session.json 是本地恢复元数据,不是待上传事件;事件文件不
## 12. 实施与验证状态
2026-09-22 新项目资格有限恢复验证:原生 `analytics` 定向测试 59 项通过、1 项既有跳过;覆盖创建时无埋点上下文、创建通知队列拒绝、独立锁竞争和临时文件系统故障后的真实受理恢复,以及资格成功落盘后不重授、进程内待办容量限制。既有损坏/备份/身份不符、跨 writer 去重测试继续通过。命令配置、Rust 格式、编码、文档索引与 diff 检查通过;独立代码复核确认业务线程不执行资格文件 I/O,内存锁不覆盖后台 I/O。未执行完整 GUI/Provider 联调或生产发布。
已做:原产品方案逐项对照;核对现役 Direct/Design 持久化、项目创建、revision、预览、UI 保存与 checkpoint 的源码入口。本文新增合同和参数均以草案标识,不作为已经上线的事实。
合同与本地队列、会话窗口与项目接入、策划阶段推进成果事件、首次提交、两类 Agent run、Direct 宿主文件/补丁成果、checkpoint 与 UI 保存,以及正式 Web 预览 ready 已完成实现、独立审查和定向验收。最终按原始需求和已确认口径核对全部 12 类事件入口,并补齐同一 writer、会话、项目和目标的宿主组件集成链路验证,本期本地采集验收通过。资源操作全面接线计划因超出原始需求撤销,未进入资源实现,不作为必做待办保留。本期无主站 API 或数据库行为;验收不表示已经提交、发布或上线。