DirectProject 埋点成绩改由宿主结算,删掉候选/settle 两阶段与 analyticsAttemptId

- analytics/run.rs:删 Request::DirectCandidate / Request::Settle 与 settle();direct_finished 改收放行代次,终态再读一次 identity_generation,不同则整条不记,相同则直接落 Request::Terminal
- analytics/store.rs:删 pending_runs / pending_run_bytes 与两条候选分支(连带 16 条上限),run_request 不再需要字节预算参数
- analytics/gui.rs + main.rs:删 settle_direct_run_analytics 命令与注册
- 放行侧:kick_queue_dispatch 在读 claim 之后读一次平台会话代次并随整轮传下去;user_input 命令不再收 analyticsAttemptId,队列条目不再持有埋点句柄
- 渲染层:删 beginDirectRunAnalytics、句柄表与结算 effect、settle 调用与命令参数;删除 tests/directRunAnalytics.test.ts,appSurface 两处断言同步
- 判据口径随放行搬家:从“发送时与终态同代”改成“放行时与终态同代”,入队后放行前的账号切换不再丢弃成绩
- 验证:cargo test -- analytics::(45 passed,含新增的直写与代际变化两条)、cargo test -- agent::(949 passed)、analytics_real_file_write 用例、npx vitest run tests/appSurface.test.ts(194 passed)、相关 chat/direct 单测 139 passed、typecheck 通过
This commit is contained in:
2026-09-30 16:25:05 +08:00
parent 15a495dc29
commit 0270cd601e
17 changed files with 94 additions and 579 deletions
@@ -5,7 +5,6 @@ mod validation;
mod wire;
pub(crate) use model::DirectCodexUserItem;
pub(crate) use validation::validate_direct_codex_user_item;
pub(crate) use wire::{
direct_codex_user_item_to_codex_turn_input, direct_codex_user_item_to_prompt,
direct_codex_user_item_to_response_item, freeze_direct_codex_user_item,
@@ -4309,7 +4309,7 @@ pub(crate) async fn run_direct_game_creator_turn_at_with_creation_type_and_emitt
crate::analytics::contract::Context,
crate::analytics::store::AnalyticsWriter,
)>,
analytics_attempt_id: Option<&str>,
release_identity_generation: u64,
) -> Result<String, DirectTurnError> {
let prompt = prompt.trim();
emit_direct_game_creator_progress(root, "request.accepted", "已发送消息,正在等待陶泥儿回复");
@@ -4323,7 +4323,7 @@ pub(crate) async fn run_direct_game_creator_turn_at_with_creation_type_and_emitt
turn_emitter,
direct_user_item,
capture,
analytics_attempt_id,
release_identity_generation,
)
.await
{
@@ -4531,7 +4531,7 @@ async fn run_direct_game_creator_turn_inner(
crate::analytics::contract::Context,
crate::analytics::store::AnalyticsWriter,
)>,
analytics_attempt_id: Option<&str>,
release_identity_generation: u64,
) -> Result<String, DirectTurnError> {
let requires_contract = super::direct_delivery::requires_new_web_contract(
root,
@@ -4948,7 +4948,7 @@ async fn run_direct_game_creator_turn_inner(
root,
&ledger.project_id,
metadata,
analytics_attempt_id,
release_identity_generation,
crate::analytics::run::Outcome {
turn_id: Some(ledger.client_turn_id.clone()),
end_reason,
@@ -49,19 +49,12 @@ pub(crate) async fn enqueue_direct_codex_turn(
user_item: DirectCodexUserItem,
creation_type: Option<String>,
client_turn_id: Option<String>,
analytics_attempt_id: Option<String>,
) -> Result<(), DirectTurnEnqueueFailure> {
let root = Path::new(project_path.trim());
let boundary_turn_id = client_turn_id.clone();
enqueue_direct_codex_turn_typed(
root,
user_item,
creation_type,
client_turn_id,
analytics_attempt_id,
)
.await
.map_err(|failure| direct_turn_enqueue_failure(root, boundary_turn_id.as_deref(), failure))
enqueue_direct_codex_turn_typed(root, user_item, creation_type, client_turn_id)
.await
.map_err(|failure| direct_turn_enqueue_failure(root, boundary_turn_id.as_deref(), failure))
}
/// 入队侧:全程 typed,顺序固定,**每一步失败都还是入队失败**:
@@ -76,7 +69,6 @@ async fn enqueue_direct_codex_turn_typed(
user_item: DirectCodexUserItem,
creation_type: Option<String>,
client_turn_id: Option<String>,
analytics_attempt_id: Option<String>,
) -> Result<(), DirectTurnError> {
let turn_id = normalize_direct_client_turn_id(client_turn_id.as_deref())?;
recover_direct_taonier_regeneration_workflow_at(root).map_err(|error| {
@@ -115,7 +107,6 @@ async fn enqueue_direct_codex_turn_typed(
user_item,
user_prompt,
creation_type,
analytics_attempt_id,
direct_tool_call_now_ms(),
)
.map_err(|error| DirectTurnError::InputRejected {
@@ -186,15 +177,9 @@ mod tests {
let subscription = subscribe_thread(&thread_id);
let _ = consume_thread(&subscription.subscription_id);
enqueue_direct_codex_turn_typed(
&root,
user_item("你好"),
None,
Some("turn-1".to_string()),
None,
)
.await
.expect("入队成立:命令只回报入队");
enqueue_direct_codex_turn_typed(&root, user_item("你好"), None, Some("turn-1".to_string()))
.await
.expect("入队成立:命令只回报入队");
let events = wait_for_completion(&subscription.subscription_id).await;
let terminal = events
@@ -246,7 +231,6 @@ mod tests {
user_item("hello"),
None,
Some("turn-1".to_string()),
None,
)
.await
.expect("入队成立:命令只回报入队");
@@ -323,7 +307,6 @@ mod tests {
user_item("生成一个游戏"),
None,
Some("turn-1".to_string()),
None,
)
.await
.expect("入队成立:命令只回报入队");
@@ -4202,6 +4202,14 @@ mod tests {
let temporary = tempfile::tempdir().unwrap();
let root = temporary.path().join("project");
let config = temporary.path().join("config");
// 成绩由宿主在终态直写:这一轮跑在哪个身份下,就看放行代次与终态代次是否同代。
let _platform_session = crate::platform_session::install_test_platform_session(
"A",
"token-a",
"https://dev.genarrative.world",
);
let release_identity_generation =
crate::platform_session::current_platform_session_write_state().identity_generation;
let (metadata, context, writer) = analytics_test_writer(&config);
let original_capture = Some((metadata.context.clone(), writer.clone()));
let mut lifecycle = crate::analytics::gui::LifecycleFixture::start(
@@ -4264,13 +4272,12 @@ mod tests {
);
assert_eq!(session.analytics_output_revision(), Some(revision.clone()));
lease.finish(true, true, None).unwrap();
let attempt = uuid::Uuid::new_v4().to_string();
run::direct_finished(
Some((context.clone(), writer.clone())),
&root,
&project_id,
&metadata,
Some(&attempt),
release_identity_generation,
run::Outcome {
turn_id: Some("analytics-write".into()),
end_reason: RunEndReason::Failed,
@@ -4280,7 +4287,6 @@ mod tests {
revision_id: session.analytics_output_revision(),
},
);
run::settle(Some((context.clone(), writer.clone())), &attempt, false);
// 默认项目使用 npm:预览服务读取真实构建目录,夹具提供构建入口,不调用构建器。
let served_root = crate::project_game_root(&root);
fs::create_dir_all(&served_root).unwrap();
@@ -13,12 +13,13 @@
use std::path::{Path, PathBuf};
#[cfg(test)]
use crate::agent::PendingTurn;
use crate::agent::{
append_direct_project_user_message_at, claim_pending_turn, complete_turn_if_reserved,
direct_tool_call_now_ms, redact_agent_runtime_error,
run_direct_game_creator_turn_at_with_creation_type_and_emitter, thread_id_for_project,
DirectGameCreatorTurnUpdateEmitter, DirectTurnError, DirectTurnTerminal, DispatchedTurn,
PendingTurn,
};
/// 一次放行的占用。持有它就代表这一轮还没收口。
@@ -72,7 +73,6 @@ impl TurnReservation {
.expect("canonical user item"),
"测试消息".to_string(),
None,
None,
direct_tool_call_now_ms(),
)
.expect("prepare pending turn");
@@ -106,10 +106,14 @@ pub(crate) fn kick_queue_dispatch(root: &Path) {
let Some(dispatched) = claim_pending_turn(&thread_id) else {
return;
};
// 放行这一刻的身份代次:这一轮的成绩只在"放行时与终态同代"时才入账(见 `analytics::run`)。
let release_identity_generation =
crate::platform_session::current_platform_session_write_state().identity_generation;
let reservation = TurnReservation::resume(&thread_id, &dispatched);
let root = root.to_path_buf();
tauri::async_runtime::spawn(async move {
run_dispatched_direct_turn(root, dispatched, reservation).await;
run_dispatched_direct_turn(root, dispatched, reservation, release_identity_generation)
.await;
});
}
@@ -121,6 +125,7 @@ async fn run_dispatched_direct_turn(
root: PathBuf,
dispatched: DispatchedTurn,
reservation: TurnReservation,
release_identity_generation: u64,
) {
let turn_id = dispatched.pending.client_turn_id.clone();
// 落盘即回合成立:放行之后必须留下这条用户消息,哪怕这一轮随后失败。
@@ -154,7 +159,7 @@ async fn run_dispatched_direct_turn(
Some(&emitter),
Some(dispatched.pending.canonical_user_item.clone()),
capture,
dispatched.pending.analytics_attempt_id.as_deref(),
release_identity_generation,
)
.await;
match outcome {
@@ -210,7 +215,6 @@ mod tests {
.expect("canonical user item"),
format!("消息 {client_turn_id}"),
None,
None,
direct_tool_call_now_ms(),
)
.expect("prepare pending turn")
@@ -1482,7 +1482,6 @@ mod tests {
.expect("canonical user item"),
format!("消息 {client_turn_id}"),
None,
None,
FIXED_AT_MS,
)
.expect("prepare pending turn")
@@ -35,8 +35,6 @@ pub(crate) struct PendingTurn {
/// 入队检查产出的 prompt:放行不重算。
pub(crate) prompt: String,
pub(crate) creation_type: Option<String>,
/// 这一轮的运行埋点句柄:由前端在入队时生成,放行开跑时用同一份身份。
pub(crate) analytics_attempt_id: Option<String>,
/// 入队那一刻的宿主毫秒钟。
pub(crate) at: u64,
}
@@ -51,7 +49,6 @@ impl PendingTurn {
user_item: DirectCodexUserItem,
prompt: String,
creation_type: Option<String>,
analytics_attempt_id: Option<String>,
at: u64,
) -> Result<Self, serde_json::Error> {
let canonical_user_item = serde_json::to_value(&user_item)?;
@@ -61,7 +58,6 @@ impl PendingTurn {
canonical_user_item,
prompt,
creation_type,
analytics_attempt_id,
at,
})
}
@@ -128,7 +124,6 @@ mod tests {
user_item("生成一个游戏", "direct-codex:turn-1:user"),
"生成一个游戏".to_string(),
creation_type.map(str::to_string),
None,
1_700_000_000_123,
)
.expect("prepare pending turn")
@@ -444,11 +444,6 @@ pub(crate) fn capture_writer_context() -> Option<(Context, AnalyticsWriter)> {
Some((capture_analytics_context()?, GUI.get()?.writer.clone()))
}
#[tauri::command]
pub(crate) fn settle_direct_run_analytics(attempt_id: String, discard: bool) {
super::run::settle(capture_writer_context(), &attempt_id, discard);
}
fn valid_context(context: &Context) -> bool {
GUI.get().is_some_and(|gui| {
context.editor_session_id == gui.session_id
@@ -55,15 +55,6 @@ pub(crate) struct Outcome {
#[derive(Serialize)]
pub(super) enum Request {
DirectCandidate {
attempt_id: String,
terminal: Box<Request>,
},
Settle {
editor_session_id: String,
attempt_id: String,
discard: bool,
},
Accepted {
root: PathBuf,
project_id: String,
@@ -82,34 +73,11 @@ pub(super) enum Request {
impl Request {
pub(super) fn session_id(&self) -> &str {
match self {
Self::DirectCandidate { terminal, .. } => terminal.session_id(),
Self::Settle {
editor_session_id, ..
} => editor_session_id,
Self::Accepted { metadata, .. } => &metadata.context.editor_session_id,
Self::Terminal { context, .. } => &context.editor_session_id,
}
}
pub(super) fn validate(&self) -> bool {
match self {
Self::DirectCandidate {
attempt_id,
terminal,
} => {
return uuid::Uuid::parse_str(attempt_id).is_ok()
&& matches!(terminal.as_ref(), Self::Terminal { metadata, .. } if metadata.source == Source::Direct)
&& terminal.validate();
}
Self::Settle {
editor_session_id,
attempt_id,
..
} => {
return uuid::Uuid::parse_str(editor_session_id).is_ok()
&& uuid::Uuid::parse_str(attempt_id).is_ok();
}
_ => {}
}
let (root, project_id, metadata) = match self {
Self::Accepted {
root,
@@ -122,7 +90,6 @@ impl Request {
metadata,
..
} => (root, project_id, metadata),
_ => return false,
};
root.is_absolute()
&& root.as_os_str().len() <= 32768
@@ -166,44 +133,29 @@ pub(crate) fn finished(
});
}
/// Direct 回合终态:成绩由宿主自己结算,中间没有候选表,也没有第二次 settle。
///
/// `release_identity_generation` 是放行那一刻的身份代次(由放行侧读并随这一轮传下来)。
/// 终态再读一次,不同就整条不记:换了号 / 退出之后这一轮的成绩不属于任何在册身份。
/// 判据口径因此从"发送时与终态同代"变成"放行时与终态同代"——入队后、放行前发生的账号切换
/// 不再丢弃成绩,那一轮确实是在新身份下跑的。
pub(crate) fn direct_finished(
capture: Option<(Context, AnalyticsWriter)>,
root: &Path,
project_id: &str,
metadata: &Metadata,
attempt_id: Option<&str>,
release_identity_generation: u64,
outcome: Outcome,
) {
let Some(attempt_id) = attempt_id else { return };
let Some((mut context, writer)) = capture else {
return;
};
if metadata.source != Source::Direct {
return;
}
context.route = metadata.context.route.clone();
writer.try_run(Request::DirectCandidate {
attempt_id: attempt_id.into(),
terminal: Box::new(Request::Terminal {
root: root.into(),
project_id: project_id.into(),
metadata: metadata.clone(),
context,
event_time: contract::timestamp_now(),
outcome,
}),
});
}
pub(crate) fn settle(capture: Option<(Context, AnalyticsWriter)>, attempt_id: &str, discard: bool) {
let Some((context, writer)) = capture else {
if crate::platform_session::current_platform_session_write_state().identity_generation
!= release_identity_generation
{
return;
};
writer.try_run(Request::Settle {
editor_session_id: context.editor_session_id,
attempt_id: attempt_id.into(),
discard,
});
}
finished(capture, root, project_id, metadata, outcome);
}
#[derive(Serialize, Deserialize)]
@@ -266,7 +218,6 @@ pub(super) fn process(request: Request) -> Result<Option<(Route, Event)>, String
metadata,
..
} => (root, project_id, metadata),
_ => return Ok(None),
};
let Some(_lock) =
crate::agent::try_acquire_game_creator_agent_runtime_task_lock(root, "analytics-runs")?
@@ -387,7 +338,6 @@ pub(super) fn process(request: Request) -> Result<Option<(Route, Event)>, String
slot.terminal_consumed = true;
Some((context.route.clone(), event))
}
_ => return Ok(None),
};
crate::agent::write_agent_runtime_json_sidecar_with_max_bytes(
root,
@@ -81,7 +81,7 @@ impl AnalyticsWriter {
worker_counters
.queued_bytes
.fetch_sub(bytes, Ordering::Relaxed);
store.run_request(request, bytes);
store.run_request(request);
}
Ok(Command::Goal(request, bytes)) => {
worker_counters
@@ -309,8 +309,6 @@ struct Batch {
}
struct Store {
pending_runs: VecDeque<(String, super::run::Request, u64)>,
pending_run_bytes: u64,
root: PathBuf,
_owner: File,
projects: HashMap<String, bool>,
@@ -342,8 +340,6 @@ impl Store {
let owner = session::claim(&instance)?;
session::recover(&instances, &session_id);
let mut store = Self {
pending_runs: VecDeque::new(),
pending_run_bytes: 0,
root,
_owner: owner,
projects: HashMap::new(),
@@ -370,57 +366,11 @@ impl Store {
Ok(store)
}
fn run_request(&mut self, request: super::run::Request, bytes: u64) {
use super::run::Request;
fn run_request(&mut self, request: super::run::Request) {
if !request.validate() || request.session_id() != self.session_id {
self.counters.dropped.fetch_add(1, Ordering::Relaxed);
return;
}
let request = match request {
Request::DirectCandidate {
attempt_id,
terminal,
} => {
if bytes > MAX_BYTES as u64 {
return;
}
if self.pending_runs.iter().any(|(id, _, _)| id == &attempt_id) {
return;
}
while self.pending_runs.len() >= 16
|| self.pending_run_bytes + bytes > MAX_BYTES as u64
{
let Some((_, _, removed)) = self.pending_runs.pop_front() else {
break;
};
self.pending_run_bytes -= removed;
self.counters.dropped.fetch_add(1, Ordering::Relaxed);
}
self.pending_run_bytes += bytes;
self.pending_runs.push_back((attempt_id, *terminal, bytes));
return;
}
Request::Settle {
attempt_id,
discard,
..
} => {
let Some(index) = self
.pending_runs
.iter()
.position(|(id, _, _)| id == &attempt_id)
else {
return;
};
let (_, terminal, removed) = self.pending_runs.remove(index).unwrap();
self.pending_run_bytes -= removed;
if discard {
return;
}
terminal
}
other => other,
};
match super::run::process(request) {
Ok(Some((route, event))) => {
let fact = format!("{}:terminal", event.agent_run_id.as_deref().unwrap());
@@ -1079,167 +1079,75 @@ fn missing_overwritten_or_unwritable_run_slots_do_not_fabricate_terminal_events(
fs::remove_dir(backup).unwrap();
}
/// Direct 成绩由宿主直写终态:放行代次与终态代次相同就落一条成绩事件,中间没有候选表。
#[test]
fn direct_terminal_waits_for_exact_attempt_settlement_and_discards_intermediate_failures() {
use super::super::{
contract::{ErrorCode, RunEndReason, RunSource, Source},
run,
};
let (_project_dir, root) = goal_project();
let config = tempfile::tempdir().unwrap();
let (context, writer) = goal_writer(config.path(), "A");
let metadata = run::Metadata::new(context.clone(), Source::Direct, RunSource::UserSubmit);
run::accepted(&writer, &root, "goal-project", &metadata);
let first = uuid::Uuid::new_v4().to_string();
let final_attempt = uuid::Uuid::new_v4().to_string();
let capture = || Some((context.clone(), writer.clone()));
let mut failure = run_outcome();
failure.end_reason = RunEndReason::Failed;
failure.error_code = Some(ErrorCode::ProviderTimeout);
run::direct_finished(
capture(),
&root,
"goal-project",
&metadata,
Some(&first),
failure.clone(),
);
run::settle(capture(), &uuid::Uuid::new_v4().to_string(), false);
assert!(!drain_goal_writer(config.path(), &context, &writer)
.iter()
.any(|e| e.agent_run_id.is_some()));
// 另一个被取消的尝试不消费槽位;保留第一次失败,验证最终确认不会选错它。
let cancelled = uuid::Uuid::new_v4().to_string();
run::direct_finished(
capture(),
&root,
"goal-project",
&metadata,
Some(&cancelled),
failure,
);
run::settle(capture(), &cancelled, true);
run::settle(capture(), &cancelled, false);
assert!(!drain_goal_writer(config.path(), &context, &writer)
.iter()
.any(|e| e.agent_run_id.is_some()));
run::direct_finished(
capture(),
&root,
"goal-project",
&metadata,
None,
run_outcome(),
);
run::direct_finished(
capture(),
&root,
"goal-project",
&metadata,
Some(&final_attempt),
run_outcome(),
);
assert!(!drain_goal_writer(config.path(), &context, &writer)
.iter()
.any(|e| e.agent_run_id.is_some()));
run::settle(capture(), &final_attempt, false);
run::settle(capture(), &final_attempt, false);
let events = drain_goal_writer(config.path(), &context, &writer);
let terminal: Vec<_> = events.iter().filter(|e| e.agent_run_id.is_some()).collect();
assert_eq!(terminal.len(), 1);
assert_eq!(terminal[0].event_name, "agent_run_completed");
assert_eq!(terminal[0].event_id, metadata.terminal_event_id);
}
#[test]
fn direct_pending_capacity_evicts_oldest_without_settlement_fallback() {
fn direct_terminal_writes_once_when_the_identity_generation_is_unchanged() {
use super::super::{
contract::{RunSource, Source},
run,
};
let _session = crate::platform_session::install_test_platform_session(
"A",
"token-a",
"https://dev.genarrative.world",
);
let release_generation =
crate::platform_session::current_platform_session_write_state().identity_generation;
let (_project_dir, root) = goal_project();
let config = tempfile::tempdir().unwrap();
let (context, writer) = goal_writer(config.path(), "A");
let metadata = run::Metadata::new(context.clone(), Source::Direct, RunSource::UserSubmit);
run::accepted(&writer, &root, "goal-project", &metadata);
let ids: Vec<_> = (0..17).map(|_| uuid::Uuid::new_v4().to_string()).collect();
for id in &ids {
run::direct_finished(
Some((context.clone(), writer.clone())),
&root,
"goal-project",
&metadata,
Some(id),
run_outcome(),
);
}
run::settle(Some((context.clone(), writer.clone())), &ids[0], false);
assert!(!drain_goal_writer(config.path(), &context, &writer)
.iter()
.any(|e| e.agent_run_id.is_some()));
assert!(writer.counters.dropped.load(Ordering::Relaxed) >= 1);
run::settle(Some((context.clone(), writer.clone())), &ids[16], false);
assert_eq!(
drain_goal_writer(config.path(), &context, &writer)
.iter()
.filter(|e| e.agent_run_id.is_some())
.count(),
1
run::direct_finished(
Some((context.clone(), writer.clone())),
&root,
"goal-project",
&metadata,
release_generation,
run_outcome(),
);
let events = drain_goal_writer(config.path(), &context, &writer);
let terminal: Vec<_> = events.iter().filter(|e| e.agent_run_id.is_some()).collect();
assert_eq!(terminal.len(), 1, "一轮只落一条成绩:{events:?}");
assert_eq!(terminal[0].event_name, "agent_run_completed");
assert_eq!(terminal[0].event_id, metadata.terminal_event_id);
}
/// 放行之后换号 / 退出:终态代次不再相同,这一轮的成绩整条不记。
#[test]
fn direct_pending_byte_budget_and_session_are_checked_before_consumption() {
fn direct_terminal_is_dropped_when_the_identity_generation_changed_since_release() {
use super::super::{
contract::{Context, RunSource, Source},
contract::{RunSource, Source},
run,
};
let _session = crate::platform_session::install_test_platform_session(
"A",
"token-a",
"https://dev.genarrative.world",
);
let release_generation =
crate::platform_session::current_platform_session_write_state().identity_generation;
let (_project_dir, root) = goal_project();
let config = tempfile::tempdir().unwrap();
let mut store = open(config.path());
let context = Context {
route: route("A"),
editor_session_id: store.session_id.clone(),
client_version: "1.0.0".into(),
};
let (context, writer) = goal_writer(config.path(), "A");
let metadata = run::Metadata::new(context.clone(), Source::Direct, RunSource::UserSubmit);
let mut last = String::new();
for _ in 0..10 {
last = uuid::Uuid::new_v4().to_string();
// worker接收到的序列化预算是独立的保守保留量。
store.run_request(
run::Request::DirectCandidate {
attempt_id: last.clone(),
terminal: Box::new(run::Request::Terminal {
root: config.path().into(),
project_id: "project".into(),
metadata: metadata.clone(),
context: context.clone(),
event_time: "2026-09-21T00:00:00.000Z".into(),
outcome: run_outcome(),
}),
},
128 * 1024,
);
}
assert_eq!(store.pending_runs.len(), 8);
assert_eq!(store.pending_run_bytes, MAX_BYTES as u64);
store.run_request(
run::Request::Settle {
editor_session_id: uuid::Uuid::new_v4().to_string(),
attempt_id: last.clone(),
discard: true,
},
128,
run::accepted(&writer, &root, "goal-project", &metadata);
// 这一轮还在跑的时候身份代次推进了(换号 / 退出)。
let next =
crate::platform_session::current_platform_session_write_state().identity_generation + 1;
crate::platform_session::clear_platform_session(next, next);
run::direct_finished(
Some((context.clone(), writer.clone())),
&root,
"goal-project",
&metadata,
release_generation,
run_outcome(),
);
assert_eq!(store.pending_runs.len(), 8);
store.run_request(
run::Request::Settle {
editor_session_id: context.editor_session_id,
attempt_id: last,
discard: true,
},
128,
assert!(
drain_goal_writer(config.path(), &context, &writer)
.iter()
.all(|event| event.agent_run_id.is_none()),
"终态代次不同时不能落任何成绩"
);
assert_eq!(store.pending_runs.len(), 7);
assert_eq!(store.pending_run_bytes, 7 * 128 * 1024);
}
@@ -2590,7 +2590,6 @@ fn main() {
})
.invoke_handler(tauri::generate_handler![
analytics::gui::capture_analytics_context,
analytics::gui::settle_direct_run_analytics,
analytics::gui::record_analytics_project_open,
analytics::gui::record_analytics_project_leave,
analytics::gui::record_analytics_ui_save,
@@ -1,39 +1,5 @@
import type { TauriInvoke } from '../app/types';
// 一个 run 可因自动认证刷新调用多次 native,只确认最后一次尝试。
export function beginDirectRunAnalytics(
invoke: TauriInvoke,
currentGeneration: () => number,
) {
const generation = currentGeneration();
let attemptId: string | undefined;
let settled = false;
return {
nextAttempt() {
attemptId = undefined;
try {
attemptId = crypto.randomUUID();
} catch {
// 最后一次 UUID 失败不能留下前一次失败候选的标识。
}
return attemptId;
},
settle() {
if (settled) return;
settled = true;
if (!attemptId) return;
try {
void invoke('settle_direct_run_analytics', {
attemptId,
discard: currentGeneration() !== generation,
}).catch(() => undefined);
} catch {
// 不等待埋点,不改变调用链的成功、失败和取消结果。
}
},
};
}
// 宿主签发的无凭据上下文:前端只传回原值,不从当前登录态补身份。
type AnalyticsContext = {
route: { user_id: string | null; destination_origin: string | null };
@@ -16,9 +16,7 @@ import {
directCodexUserItemFromContent,
hasMeaningfulDirectCodexContent,
} from '../../../../features/project-workspace/resourceReferences';
import { beginDirectRunAnalytics } from '../../../../services/clientAnalytics';
import { captureAgentRuntimeError } from '../../../../services/errorReporting';
import { currentPlatformSessionGeneration } from '../../../../services/platformSession';
import {
createDirectProjectTurnId,
DIRECT_CODEX_AGENT_ID,
@@ -179,11 +177,6 @@ export function useDirectProjectChatController({
const directEntries = directThread.entries;
const currentTurnRunning = directThread.turnRunning;
const currentTurnStartedAt = directThread.turnStartedAt;
const completedTurnCount = directThread.completedTurnCount;
/**
* 本轮开口条目的 canonical 身份:收口之后仍留着,用来把"刚结束的是哪一轮"讲清楚(结算埋点读它)。
*/
const currentTurnUserItemId = directThread.state.turnUserItemId;
/**
* 待发消息:**投影自事件流**,界面不持有第二份队列。入队、取消、放行都只让这条列表跟着变。
*/
@@ -197,19 +190,6 @@ export function useDirectProjectChatController({
directTurnRunningRef.current = currentTurnRunning;
const commandInFlightRef = useRef(commandInFlight);
commandInFlightRef.current = commandInFlight;
/** 已经处理过的收口回合数:与 reducer 的计数比较,识别"又有回合结束了"。 */
const handledCompletedTurnCountRef = useRef(0);
/**
* 每个已经交给宿主的回合的埋点句柄,按 `clientTurnId` 存。
*
* 必须按身份存:队列在宿主侧,同一时刻可能有好几条待发消息(各自带着自己的 attempt id),
* 而成绩是**各自那一轮**收尾时在宿主侧入账的——单个槽位只会结算到最后一个,前面的全部静默丢掉。
* 句柄活到本轮终态(见下面结算的那个 effect)。
*/
const pendingRunAnalyticsRef = useRef(
new Map<string, ReturnType<typeof beginDirectRunAnalytics>>(),
);
useEffect(() => {
setAttachmentNotice('');
setComposerNotice('');
@@ -217,8 +197,6 @@ export function useDirectProjectChatController({
setLocalMessages([]);
setHistoryHasMore(false);
historyOldestItemIdRef.current = null;
handledCompletedTurnCountRef.current = 0;
pendingRunAnalyticsRef.current.clear();
}, [projectPath]);
// 回合期间的平台会话保活由 Rust 持有:Direct 回合的占用登记时启动,占用释放(回合收口)即停止;
@@ -227,26 +205,6 @@ export function useDirectProjectChatController({
// 本地忙态不再由这里兜:入队化之后"这一轮开没开"只由事件流回答(`useDirectProjectTurnStatus` 的
// `commandInFlight` / `pendingCount` / `nativeRunning` 三者并集)。
/**
* 回合终态是**埋点结算**的唯一出口(入队失败走命令那条路,见 `runTurn` 的清理)。
*
* 判据用 reducer 的**单调计数**而不是 `turnRunning` 的下降沿:一轮可能在同一次 consume 里
* 开始并结束,那时下降沿永远不会出现,这一轮的埋点就永远没人结算。
*
* 队列**不在这里推进**:待发消息归宿主,下一条什么时候走由 Thread Manager 在占用释放时自己踢。
*/
useEffect(() => {
if (!enabled) {
handledCompletedTurnCountRef.current = completedTurnCount;
return;
}
if (completedTurnCount === handledCompletedTurnCountRef.current) return;
handledCompletedTurnCountRef.current = completedTurnCount;
// 先结算埋点:宿主此刻已经把这一轮的成绩写进候选,再晚也还是同一轮。
settlePendingRunAnalytics(currentTurnUserItemId);
// eslint-disable-next-line react-hooks/exhaustive-deps
}, [enabled, completedTurnCount, currentTurnUserItemId]);
useEffect(() => {
if (!enabled || !projectPath) return;
let disposed = false;
@@ -283,29 +241,6 @@ export function useDirectProjectChatController({
setCommandInFlight(false);
}
/**
* 结算本轮埋点。只在**回合终态**调用:成绩是回合末才在宿主侧入账的,入队返回时就结算
* 会变成一次空操作(宿主找不到候选,静默丢弃)。入队失败那一轮没有候选,句柄由 `runTurn` 自己清掉。
*
* 认领靠身份:`userItemId` 是本轮开口条目的 canonical id,与 `clientTurnId` 一一对应
* (`directCodexConversationMessageId`)。身份缺失(旧事件 / 没有开口条目)时退一步结算**最早**的
* 那一条——放行严格按队首顺序,收口顺序就是入队顺序,这一层没有歧义。
*/
function settlePendingRunAnalytics(turnUserItemId: string) {
const pending = pendingRunAnalyticsRef.current;
if (pending.size === 0) return;
const matched = [...pending.keys()].find(
(clientTurnId) =>
directCodexConversationMessageId(clientTurnId, 'user') ===
turnUserItemId,
);
const key = matched ?? pending.keys().next().value;
if (key === undefined) return;
const handle = pending.get(key);
pending.delete(key);
handle?.settle();
}
function appendLocalMessage(message: ChatMessage) {
setLocalMessages((current) => [...current, message]);
}
@@ -330,10 +265,7 @@ export function useDirectProjectChatController({
);
if (projectPathRef.current !== nextProjectPath) return;
if (outcome === 'removed') {
// 这一条不会再被放行,也就不会再有 `turn.completed` 来结算它:句柄只能在这里清掉。
// 留着不但泄漏,还会挤进"身份缺失"的兜底结算的候选里(`settlePendingRunAnalytics`
// 在没有对得上身份的候选时挑最早那一条)——取消一条消息不该改变别人的结算是谁。
pendingRunAnalyticsRef.current.delete(clientTurnId);
// 这一条不会再被放行:界面上撤掉它就行,成绩本来也只在放行后由宿主自己结算。
}
if (outcome === 'notFound') {
setComposerNotice('这条待发消息已经不在队列里');
@@ -544,18 +476,10 @@ export function useDirectProjectChatController({
// 这里的返回值只回答"草稿能不能清"——入队成立能清,用户自己能改的入队失败不能清。
let draftSafe = false;
try {
const runAnalytics = beginDirectRunAnalytics(
invoke,
currentPlatformSessionGeneration,
);
// 埋点句柄必须活过命令返回:成绩是回合末才入账的,settle 只能在回合终态发生。
// 按 `clientTurnId` 存:队列在宿主侧,同一时刻可能有好几条待发消息各自带着自己的身份。
pendingRunAnalyticsRef.current.set(input.clientTurnId, runAnalytics);
await invoke<string>('enqueue_direct_codex_turn', {
projectPath: nextProjectPath,
clientTurnId: input.clientTurnId,
userItem: input.userItem,
analyticsAttemptId: runAnalytics.nextAttempt(),
...(input.creationType ? { creationType: input.creationType } : {}),
});
draftSafe = true;
@@ -564,12 +488,6 @@ export function useDirectProjectChatController({
// 命令的入队失败是**结构化的**:命令返回 `Ok` 只说明入队成立,所以这条 catch 从入队化之后
// 只剩"入队失败"一种输入(整轮结果由 `turn.completed` 事件回答,不再回到这里)。
const rejection = readDirectTurnEnqueueFailure(error);
// 只有结构化入队失败能证明"这一轮没成立"(入队失败不产生回合事件、也就不会有埋点候选),
// 句柄才只清不发;非结构化错误(IPC 失败、命令 panic)可能发生在放行之后,那时必须留着句柄等
// `turn.completed` 来结算——提前清掉会让宿主侧这一轮的候选永远没有人结算。
if (rejection) {
pendingRunAnalyticsRef.current.delete(input.clientTurnId);
}
if (rejection) {
const notice = directTurnEnqueueFailureNotice(rejection);
if (notice) {
@@ -495,7 +495,7 @@ export function registerChatComposerControlTests() {
it('does not re-run the whole DirectProject turn when authentication fails', async () => {
// 登录态失效不再"刷新 + 重跑整轮"(重跑会重复落盘同一条用户消息):它按普通回合结果呈现,
// 整轮只 invoke 一次、埋点只结算一次。
// 整轮只 invoke 一次。
const { invoke, surface } = await openDirectCodexSurface({
enqueue_direct_codex_turn: () => {
throw new Error('authentication-required');
@@ -520,12 +520,6 @@ export function registerChatComposerControlTests() {
);
expect(attempts).toHaveLength(1);
expect(refresh).not.toHaveBeenCalled();
// 这一轮没有入队,不会有回合终态事件来驱动结算:埋点一次也不结算。
expect(
invoke.mock.calls.filter(
([command]) => command === 'settle_direct_run_analytics',
),
).toHaveLength(0);
// 认不出的入队失败 / 非结构化错误仍走既有捕获链路:横幅给用户一句可读的话。
await waitFor(() => {
expect(
@@ -647,61 +641,6 @@ export function registerChatComposerControlTests() {
).toEqual(['第一条消息', '第二条消息', '第三条消息']);
});
it('不再认领被取消那一条的埋点句柄:身份缺失的兜底结算只落到还能收口的那一条', async () => {
// 取消掉的待发消息永远不会被放行,也就永远不会有它自己的 `turn.completed`:句柄必须在取消时清掉。
// 留着的话,"身份缺失"的兜底结算(认领表里**最早**那一条)会先认出它——取消一条消息等于把下一次
// 结算算到一条根本没跑过的回合头上,真正跑着的那条反而没人结算。
let host: ReturnType<typeof hostQueue> | null = null;
const { invoke, surface, harness } = await openDirectCodexSurface({
enqueue_direct_codex_turn: (args) => host?.enqueue(args),
remove_direct_project_pending_turn: (args) => host?.cancel(args),
});
// 线程上已有一条(恢复出来的)回合在跑:发出去的两条都只会排队,事件流里也不会走 `turn.started`。
host = hostQueue(harness, { running: true });
const composer = within(surface).getByLabelText('陶泥儿对话内容');
await setComposerText(composer, '会被取消的消息');
submitComposerForm(composer);
await setComposerText(composer, '留下来的消息');
submitComposerForm(composer);
const queue = await within(surface).findByLabelText('待发送消息队列');
await waitFor(() => {
expect(within(queue).getAllByRole('listitem')).toHaveLength(2);
});
fireEvent.click(
within(queue).getByRole('button', {
name: '取消待发消息 会被取消的消息',
}),
);
await waitFor(() => {
expect(within(queue).getAllByRole('listitem')).toHaveLength(1);
});
const attempts = invoke.mock.calls
.filter(([command]) => command === 'enqueue_direct_codex_turn')
.map(([, args]) => String(args?.analyticsAttemptId ?? ''));
expect(attempts).toHaveLength(2);
const settledAttemptIds = () =>
invoke.mock.calls
.filter(([command]) => command === 'settle_direct_run_analytics')
.map(([, args]) => String(args?.attemptId ?? ''));
// 回放一条**没有身份**的失败终态:订阅重建时的生命周期锚点就是这个形状(订阅在新回合开始前
// 断过,开局那条 `turn.started` 不在手里)。兜底结算此刻只能按"认领表里最早那一条"认人——
// 被取消的那条若还留在表里,就会顶替掉真正该结算的那一条。
act(() => {
harness.emitDirectThreadEvents({
type: 'turn.completed',
status: 'failed',
failure: { message: '重连之前那一轮已经失败' },
at: 4,
});
});
await waitFor(() => {
expect(settledAttemptIds()).toHaveLength(1);
});
expect(settledAttemptIds()).toEqual([attempts[1]]);
});
it('tells the user when the chip they tried to cancel has already been dispatched', async () => {
// 已放行的那一条不能按待发消息取消:宿主的 typed 结果把这件事讲清楚,界面照说,不自作主张。
const { invoke, surface, harness } = await openDirectCodexSurface({
@@ -1537,7 +1537,6 @@ export function registerHomeProjectCreationTests() {
projectPath: automaticProjectPath,
creationType: 'game',
clientTurnId: expect.any(String),
analyticsAttemptId: expect.any(String),
userItem: {
id: expect.stringMatching(/^direct-codex:[A-Za-z0-9-]+:user$/),
type: 'message',
@@ -1657,7 +1656,6 @@ export function registerHomeProjectCreationTests() {
projectPath: automaticProjectPath,
creationType: 'game',
clientTurnId: expect.any(String),
analyticsAttemptId: expect.any(String),
userItem: {
id: expect.stringMatching(/^direct-codex:[A-Za-z0-9-]+:user$/),
type: 'message',
@@ -1696,7 +1694,6 @@ export function registerHomeProjectCreationTests() {
expect(invoke).toHaveBeenCalledWith('enqueue_direct_codex_turn', {
projectPath: automaticProjectPath,
clientTurnId: expect.any(String),
analyticsAttemptId: expect.any(String),
userItem: {
id: expect.stringMatching(/^direct-codex:[A-Za-z0-9-]+:user$/),
type: 'message',
@@ -2447,7 +2444,6 @@ export function registerHomeProjectCreationTests() {
expect(invoke).toHaveBeenCalledWith('enqueue_direct_codex_turn', {
projectPath,
clientTurnId: expect.any(String),
analyticsAttemptId: expect.any(String),
userItem: {
id: expect.stringMatching(/^direct-codex:[A-Za-z0-9-]+:user$/),
type: 'message',
@@ -2501,7 +2497,6 @@ export function registerHomeProjectCreationTests() {
expect(args).toEqual({
projectPath,
clientTurnId: expect.any(String),
analyticsAttemptId: expect.any(String),
userItem: {
id: `direct-codex:${clientTurnId}:user`,
type: 'message',
@@ -1,91 +0,0 @@
// @vitest-environment jsdom
import { afterEach, expect, it, vi } from 'vitest';
import type { TauriInvoke } from '../src/app/types';
import { beginDirectRunAnalytics } from '../src/services/clientAnalytics';
afterEach(() => vi.restoreAllMocks());
it('自动重试保持业务回合,只确认最后一次尝试且确认不等待桥接', async () => {
const first = '11111111-1111-4111-8111-111111111111';
const last = '22222222-2222-4222-8222-222222222222';
vi.spyOn(crypto, 'randomUUID')
.mockReturnValueOnce(first)
.mockReturnValueOnce(last);
const invoke = vi.fn(() => new Promise(() => {}));
const analytics = beginDirectRunAnalytics(invoke as TauriInvoke, () => 1);
const attempts: string[] = [];
const operation = async () => {
attempts.push(analytics.nextAttempt()!);
if (attempts.length === 1) throw new Error('authentication-required');
return '完成';
};
let result: string;
try {
result = await operation().catch(() => operation());
expect(invoke).not.toHaveBeenCalled();
} finally {
analytics.settle();
}
expect(result).toBe('完成');
expect(attempts).toEqual([first, last]);
expect(invoke).toHaveBeenCalledTimes(1);
expect(invoke).toHaveBeenCalledWith('settle_direct_run_analytics', {
attemptId: last,
discard: false,
});
analytics.settle();
expect(invoke).toHaveBeenCalledTimes(1);
});
it('最终 UUID 失败不能回退确认第一次失败', () => {
vi.spyOn(crypto, 'randomUUID')
.mockReturnValueOnce('11111111-1111-4111-8111-111111111111')
.mockImplementationOnce(() => {
throw new Error('UUID unavailable');
});
const invoke = vi.fn();
const analytics = beginDirectRunAnalytics(invoke as TauriInvoke, () => 1);
expect(analytics.nextAttempt()).toBeDefined();
expect(analytics.nextAttempt()).toBeUndefined();
analytics.settle();
expect(invoke).not.toHaveBeenCalled();
});
it('账号代次变化只丢弃候选,不提交成功失败结论', () => {
let generation = 1;
const invoke = vi.fn(async () => undefined);
const analytics = beginDirectRunAnalytics(
invoke as TauriInvoke,
() => generation,
);
const attemptId = analytics.nextAttempt();
generation = 2;
analytics.settle();
expect(invoke).toHaveBeenCalledWith('settle_direct_run_analytics', {
attemptId,
discard: true,
});
});
it.each(['sync', 'async'])(
'确认桥接 %s 失败不会覆盖最终业务错误',
async (mode) => {
const invoke = vi.fn(() => {
if (mode === 'sync') throw new Error('bridge');
return Promise.reject(new Error('bridge'));
});
const analytics = beginDirectRunAnalytics(invoke as TauriInvoke, () => 1);
const failure = new Error('最终执行失败');
const operation = async () => {
try {
analytics.nextAttempt();
throw failure;
} finally {
analytics.settle();
}
};
await expect(operation()).rejects.toBe(failure);
expect(invoke).toHaveBeenCalledTimes(1);
},
);