diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent.rs b/apps/ai-game-creator-shell/src-tauri/src/agent.rs index a828e4928..9abc27465 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent.rs @@ -15,7 +15,6 @@ mod codex_provider_proxy; mod design_runtime; pub(crate) mod design_tools; mod direct_codex_attachments; -mod direct_codex_audit; mod direct_codex_user_item; mod direct_execution; pub(crate) use direct_execution::WritePermit; @@ -34,7 +33,6 @@ mod direct_thread_wire; mod direct_tool_bridge; mod direct_tool_calls; mod direct_tools_mcp; -mod direct_turn_metrics; mod direct_turn_stream; mod direct_validation; mod generation; @@ -61,7 +59,6 @@ pub(crate) use codex_cli::{ pub(crate) use codex_provider_proxy::*; pub(crate) use design_runtime::*; pub(crate) use direct_codex_attachments::*; -pub(crate) use direct_codex_audit::*; pub(crate) use direct_codex_user_item::*; pub(crate) use direct_project_history::*; pub(crate) use direct_project_turn_history::*; @@ -71,7 +68,6 @@ pub(crate) use direct_thread_wire::*; pub(crate) use direct_tool_bridge::*; pub(crate) use direct_tool_calls::*; pub(crate) use direct_tools_mcp::*; -pub(crate) use direct_turn_metrics::*; pub(crate) use direct_turn_stream::*; pub(crate) use direct_validation::DirectValidationConfig; pub(crate) use generation::*; diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/codex_app_server/mod.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/codex_app_server/mod.rs index b30f5ef2b..1fec303bc 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/codex_app_server/mod.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/codex_app_server/mod.rs @@ -3219,15 +3219,8 @@ impl CodexAppServerConnection { request: LlmRunRequest, on_agent_message_delta: Option<&mut (dyn FnMut(&platform_llm::LlmStreamDelta) + Send)>, ) -> Result { - self.run_turn_with_direct_observer( - snapshot, - llm, - request, - on_agent_message_delta, - None, - None, - ) - .await + self.run_turn_with_direct_observer(snapshot, llm, request, on_agent_message_delta, None) + .await } async fn run_turn_with_direct_observer( @@ -3237,7 +3230,6 @@ impl CodexAppServerConnection { request: LlmRunRequest, on_agent_message_delta: Option<&mut (dyn FnMut(&platform_llm::LlmStreamDelta) + Send)>, direct_observer: Option<&mut (dyn FnMut(DirectCodexTurnObservation) + Send)>, - audit: Option<&mut DirectCodexTurnAudit>, ) -> Result { self.run_turn_with_direct_observer_and_history( snapshot, @@ -3249,8 +3241,6 @@ impl CodexAppServerConnection { DirectCodexTurnKind::User, on_agent_message_delta, direct_observer, - audit, - None, ) .await } @@ -3266,27 +3256,8 @@ impl CodexAppServerConnection { turn_kind: DirectCodexTurnKind, mut on_agent_message_delta: Option<&mut (dyn FnMut(&platform_llm::LlmStreamDelta) + Send)>, mut direct_observer: Option<&mut (dyn FnMut(DirectCodexTurnObservation) + Send)>, - mut audit: Option<&mut DirectCodexTurnAudit>, - metrics_attempt: Option, ) -> Result { - let mut gate_timing = metrics_attempt - .as_ref() - .map(|attempt| attempt.span("local-turn-gate")); let _turn_guard = self.inner.turn_gate.lock().await; - if let Some(timing) = gate_timing.as_mut() { - timing.finish("acquired"); - } - // Scope only after acquiring the per-connection gate. Requests clone the binding - // at ingress, so a late body never borrows the next turn's identity. - let _metrics_binding = metrics_attempt.as_ref().and_then(|attempt| { - match self.inner._provider_proxy.as_ref() { - Some(proxy) => Some(proxy.bind_metrics(attempt.clone())), - None => { - attempt.route(DirectMetricRoute::AppServerAuth); - None - } - } - }); let history_root = direct_history_root.unwrap_or(&self.inner.workspace_path); // 工具调用卡片的 turnId 用 AGC 客户端回合 id(与实时事件、落盘条目同一口径), // 不用 Codex app-server 自己的 turnId——前端要按它把卡片挂回对应的那一轮。 @@ -3403,21 +3374,11 @@ impl CodexAppServerConnection { if thread_created && self.inner.workspace_mode == CodexAppServerWorkspaceMode::DirectProject { if let Some(client_turn_id) = direct_client_turn_id { - let mut prefetch_timing = metrics_attempt - .as_ref() - .map(|attempt| attempt.span("project-context-prefetch")); let prefetched = super::direct_project_context::prefetch_turn_input( history_root, client_turn_id, ) .await; - if let Some(timing) = prefetch_timing.as_mut() { - timing.finish(if prefetched.is_ok() { - "completed" - } else { - "failed" - }); - } match prefetched { Ok(Some(context)) => { if let Some(parts) = input.as_array_mut() { @@ -3505,9 +3466,6 @@ impl CodexAppServerConnection { cancellation: Arc::clone(&turn_start_cancellation), armed: true, }; - let mut start_timing = metrics_attempt - .as_ref() - .map(|attempt| attempt.span("turn-start-ack")); let result = match self .request_with_turn_start_cancellation( "turn/start", @@ -3516,16 +3474,8 @@ impl CodexAppServerConnection { ) .await { - Ok(result) => { - if let Some(timing) = start_timing.as_mut() { - timing.finish("acknowledged"); - } - result - } + Ok(result) => result, Err(error) => { - if let Some(timing) = start_timing.as_mut() { - timing.finish("failed"); - } if let Some(adapter) = approval_adapter .as_ref() .filter(|adapter| adapter.is_host_ending()) @@ -3661,18 +3611,8 @@ impl CodexAppServerConnection { .await); } }; - if event.is_some() { - if let Some(attempt) = metrics_attempt.as_ref() { - attempt.observe_app_event("first-event"); - } - } match event { Some(CodexTurnEvent::AgentMessageDelta { item_id, delta }) => { - if !delta.is_empty() { - if let Some(attempt) = metrics_attempt.as_ref() { - attempt.observe_app_event("first-content-delta"); - } - } if self.inner.workspace_mode == CodexAppServerWorkspaceMode::DirectProject { direct_project_history.observe_delta(&item_id, &delta); append_direct_thread_event( @@ -3716,11 +3656,6 @@ impl CodexAppServerConnection { } } Some(CodexTurnEvent::ReasoningDelta { item_id, delta }) => { - if !delta.is_empty() { - if let Some(attempt) = metrics_attempt.as_ref() { - attempt.observe_app_event("first-reasoning-delta"); - } - } if self.inner.workspace_mode == CodexAppServerWorkspaceMode::DirectProject { append_direct_thread_event( &direct_thread_id, @@ -3739,9 +3674,6 @@ impl CodexAppServerConnection { } } Some(CodexTurnEvent::RawItem(item)) => { - if let Some(attempt) = metrics_attempt.as_ref() { - attempt.observe_raw_item(&item); - } if self.inner.workspace_mode == CodexAppServerWorkspaceMode::DirectProject { if item.is_null() { return Err(platform_llm::LlmError::Deserialize( @@ -3796,9 +3728,6 @@ impl CodexAppServerConnection { } Some(CodexTurnEvent::Item { completed, params }) => { if let Some(item) = params.get("item") { - if let Some(attempt) = metrics_attempt.as_ref() { - attempt.observe_item(item, completed); - } let item_type = item .get("type") .and_then(serde_json::Value::as_str) @@ -3849,11 +3778,6 @@ impl CodexAppServerConnection { } } } - if completed { - if let Some(audit) = audit.as_mut() { - audit.observe_item(¶ms); - } - } } if item_type == "agentMessage" { // 某些 app-server 实现会在工具开始后停止发送 agentMessage delta, @@ -3987,9 +3911,6 @@ impl CodexAppServerConnection { } return execution::outcome_text(adapter.wait_outcome().await); } - if let Some(attempt) = metrics_attempt.as_ref() { - attempt.finish("interrupted"); - } return Err(platform_llm::LlmError::InvalidRequest( "Codex app-server turn 已中断".to_string(), )); @@ -5073,7 +4994,6 @@ pub(crate) async fn direct_game_creator_codex_chat_at( None, None, None, - None, ) .await } @@ -5092,7 +5012,6 @@ pub(crate) async fn direct_game_creator_codex_chat_at_with_observer( None, Some(observer), None, - None, ) .await } @@ -5104,7 +5023,6 @@ pub(crate) async fn direct_game_creator_codex_chat_at_with_optional_observer( turn_kind: DirectCodexTurnKind, client_turn_id: Option<&str>, observer: Option<&mut (dyn FnMut(DirectCodexTurnObservation) + Send)>, - audit: Option<&mut DirectCodexTurnAudit>, direct_user_item: Option, ) -> Result { // Resolve project authority before deriving the pool/thread identity. A @@ -5157,16 +5075,6 @@ pub(crate) async fn direct_game_creator_codex_chat_at_with_optional_observer( }; let api_kind = parse_game_creator_llm_api_kind(&config.llm.api_kind).map_err(|error| error.to_string())?; - let metrics_attempt = audit.as_ref().map(|audit| { - audit.metrics().attempt( - &config.llm.model, - &config.llm.model, - &config.llm.reasoning_effort, - ) - }); - let mut connection_timing = metrics_attempt - .as_ref() - .map(|attempt| attempt.span("connection-preparation")); let connection = Box::pin(CodexAppServerConnection::acquire_at_workspace( &snapshot, &config.llm, @@ -5174,26 +5082,14 @@ pub(crate) async fn direct_game_creator_codex_chat_at_with_optional_observer( CodexAppServerWorkspaceMode::DirectProject, effective_client_turn_id, )) - .await; - if let Some(timing) = connection_timing.as_mut() { - timing.finish(if connection.is_ok() { - "ready" - } else { - "failed" - }); - } - let connection = connection.map_err(|error| { - if let Some(attempt) = metrics_attempt.as_ref() { - attempt.finish("failed"); - } - error.to_string() - })?; + .await + .map_err(|error| error.to_string())?; let request = LlmRunRequest::single_turn(system_prompt, user_prompt) .with_api_kind(api_kind) .with_model(config.llm.model.clone()) .with_request_timeout_ms(config.llm.request_timeout_ms) .with_max_output_tokens(16_000); - let result = connection + connection .run_turn_with_direct_observer_and_history( &snapshot, &config.llm, @@ -5204,20 +5100,10 @@ pub(crate) async fn direct_game_creator_codex_chat_at_with_optional_observer( turn_kind, None, observer, - audit, - metrics_attempt.clone(), ) .await .map(|value| value.text) - .map_err(|error| error.to_string()); - if let Some(attempt) = metrics_attempt.as_ref() { - attempt.finish(if result.is_ok() { - "completed" - } else { - "failed" - }); - } - result + .map_err(|error| error.to_string()) } /// Direct home-page chat never binds Codex to a user project. It gets a @@ -5383,7 +5269,6 @@ mod tests { Some("not-executed"), None, None, - None, ); let sizes = ( std::mem::size_of_val(&spawn), @@ -7612,7 +7497,6 @@ while IFS= read -r line; do :; done tool_request(), Some(&mut on_delta), Some(&mut observer), - None, ) .await .expect("run fake app-server turn"); @@ -7741,7 +7625,6 @@ while IFS= read -r line; do :; done tool_request(), None, Some(&mut observer), - None, ) .await .expect("run fake app-server turn"); @@ -7861,8 +7744,6 @@ done DirectCodexTurnKind::User, None, Some(&mut observer), - None, - None, ) .await .expect("run direct-project turn"); diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/codex_provider_proxy.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/codex_provider_proxy.rs index 63773936a..a74bbb434 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/codex_provider_proxy.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/codex_provider_proxy.rs @@ -1,4 +1,3 @@ -use super::{DirectMetricAttempt, DirectMetricRoute, DirectRequestTiming}; use axum::body::{to_bytes, Body}; use axum::extract::State; use axum::http::{HeaderMap, HeaderName, Request, Response, StatusCode}; @@ -28,7 +27,6 @@ struct CodexProviderProxyState { downstream_bearer_token: String, main_site_upstream: bool, client: reqwest::Client, - metrics_scope: Arc>>, parallel_tool_calls: bool, model_usage: ActiveModelUsage, } @@ -37,8 +35,6 @@ pub(crate) struct CodexProviderProxy { base_url: String, downstream_bearer_token: String, task: tokio::task::JoinHandle<()>, - metrics_scope: Arc>>, - main_site_upstream: bool, model_usage: ActiveModelUsage, } @@ -71,21 +67,6 @@ impl CodexProviderProxy { &self.downstream_bearer_token } - pub(crate) fn bind_metrics(&self, attempt: DirectMetricAttempt) -> CodexProviderMetricsBinding { - attempt.route(if self.main_site_upstream { - DirectMetricRoute::MainSite - } else { - DirectMetricRoute::ProviderProxy - }); - if let Ok(mut scope) = self.metrics_scope.lock() { - *scope = Some(attempt.clone()); - } - CodexProviderMetricsBinding { - scope: Arc::clone(&self.metrics_scope), - attempt_id: attempt.id().to_string(), - } - } - pub(crate) fn begin_model_usage( &self, context: crate::project::ProjectModelUsageContext, @@ -102,32 +83,12 @@ impl CodexProviderProxy { } } -/// A late stream owns its original attempt; releasing a binding cannot clear a new one. -pub(crate) struct CodexProviderMetricsBinding { - scope: Arc>>, - attempt_id: String, -} - -impl Drop for CodexProviderMetricsBinding { - fn drop(&mut self) { - if let Ok(mut scope) = self.scope.lock() { - if scope - .as_ref() - .is_some_and(|attempt| attempt.id() == self.attempt_id) - { - *scope = None; - } - } - } -} - -struct MeasuredResponseStream { +struct ObservedResponseStream { inner: Pin>, - timing: Option, - observer: Option, + observer: ModelResponseObserver, } -impl Stream for MeasuredResponseStream +impl Stream for ObservedResponseStream where S: Stream>, { @@ -137,32 +98,17 @@ where let this = self.get_mut(); match this.inner.as_mut().poll_next(cx) { Poll::Ready(Some(Ok(bytes))) => { - if let Some(timing) = this.timing.as_mut() { - timing.chunk(&bytes); - } - if let Some(observer) = this.observer.as_mut() { - observer.observe(&bytes); - } + this.observer.observe(&bytes); Poll::Ready(Some(Ok(bytes))) } Poll::Ready(Some(Err(_))) => { - if let Some(timing) = this.timing.as_mut() { - timing.finish("stream-error"); - } - if let Some(observer) = this.observer.as_mut() { - observer.failed(); - } + this.observer.failed(); Poll::Ready(Some(Err(std::io::Error::other( "provider response stream failed", )))) } Poll::Ready(None) => { - if let Some(timing) = this.timing.as_mut() { - timing.finish("eof"); - } - if let Some(observer) = this.observer.as_mut() { - observer.finish(); - } + this.observer.finish(); Poll::Ready(None) } Poll::Pending => Poll::Pending, @@ -277,12 +223,6 @@ async fn proxy_codex_provider_request( if request.method() != axum::http::Method::POST || request.uri().path() != "/responses" { return proxy_error(StatusCode::NOT_FOUND, "provider proxy route not found"); } - let mut timing = state - .metrics_scope - .lock() - .ok() - .and_then(|scope| scope.clone()) - .map(DirectRequestTiming::new); // 在读请求体或等待上游之前冻结归属,迟到响应不能使用下一回合的项目上下文。 let model_usage = state .model_usage @@ -293,32 +233,21 @@ async fn proxy_codex_provider_request( let (parts, body) = request.into_parts(); let body = match to_bytes(body, CODEX_PROVIDER_PROXY_MAX_REQUEST_BYTES).await { Ok(body) => body, - Err(_) => { - if let Some(timing) = timing.as_mut() { - timing.finish("request-body-error"); - } - return proxy_error(StatusCode::PAYLOAD_TOO_LARGE, "provider request too large"); - } + Err(_) => return proxy_error(StatusCode::PAYLOAD_TOO_LARGE, "provider request too large"), }; let body = if state.parallel_tool_calls { match tokio::task::spawn_blocking(move || parallel_direct_request(&body)).await { Ok(Ok(bytes)) => axum::body::Bytes::from(bytes), _ => { - if let Some(timing) = timing.as_mut() { - timing.finish("request-body-error"); - } return proxy_error( StatusCode::BAD_REQUEST, "provider request JSON invalid or oversized", - ); + ) } } } else { body }; - if let Some(timing) = timing.as_mut() { - timing.request_body(&body); - } let mut headers = HeaderMap::new(); for (name, value) in &parts.headers { if !is_hop_by_hop_header(name) && name != axum::http::header::AUTHORIZATION { @@ -336,19 +265,13 @@ async fn proxy_codex_provider_request( let upstream_authorization = match format!("Bearer {}", state.upstream_bearer_token).parse() { Ok(value) => value, Err(_) => { - if let Some(timing) = timing.as_mut() { - timing.finish("invalid-credential"); - } return proxy_error( StatusCode::INTERNAL_SERVER_ERROR, "provider proxy credential invalid", - ); + ) } }; headers.insert(axum::http::header::AUTHORIZATION, upstream_authorization); - if let Some(timing) = timing.as_mut() { - timing.dispatched(); - } let upstream = match state .client .request(parts.method, upstream_url) @@ -358,32 +281,14 @@ async fn proxy_codex_provider_request( .await { Ok(response) => response, - Err(_) => { - if let Some(timing) = timing.as_mut() { - timing.finish("upstream-error"); - } - return proxy_error(StatusCode::BAD_GATEWAY, "provider upstream unavailable"); - } + Err(_) => return proxy_error(StatusCode::BAD_GATEWAY, "provider upstream unavailable"), }; let status = upstream.status(); let upstream_headers = upstream.headers().clone(); - if let Some(timing) = timing.as_mut() { - let sse = upstream_headers - .get("content-type") - .and_then(|value| value.to_str().ok()) - .is_some_and(|value| { - value - .split(';') - .next() - .is_some_and(|value| value.trim().eq_ignore_ascii_case("text/event-stream")) - }); - timing.headers(status.as_u16(), sse); - } let observer = ModelResponseObserver::new(model_usage, status, &upstream_headers); - let stream = MeasuredResponseStream { + let stream = ObservedResponseStream { inner: Box::pin(upstream.bytes_stream()), - timing, - observer: Some(observer), + observer, }; let mut response = Response::builder().status(status); if let Some(headers) = response.headers_mut() { @@ -453,7 +358,6 @@ pub(crate) async fn start_codex_provider_proxy_with_parallel( let address = listener .local_addr() .map_err(|error| format!("读取 Codex Provider 代理地址失败:{error}"))?; - let metrics_scope = Arc::new(Mutex::new(None)); let model_usage = Arc::new(Mutex::new(None)); let state = Arc::new(CodexProviderProxyState { upstream_base_url, @@ -461,7 +365,6 @@ pub(crate) async fn start_codex_provider_proxy_with_parallel( downstream_bearer_token: downstream_bearer_token.clone(), main_site_upstream, client, - metrics_scope: Arc::clone(&metrics_scope), parallel_tool_calls, model_usage: Arc::clone(&model_usage), }); @@ -475,8 +378,6 @@ pub(crate) async fn start_codex_provider_proxy_with_parallel( base_url: format!("http://127.0.0.1:{}", address.port()), downstream_bearer_token, task, - metrics_scope, - main_site_upstream, model_usage, }) } @@ -488,30 +389,8 @@ mod tests { use futures::StreamExt; use std::sync::atomic::{AtomicUsize, Ordering}; - fn timing_log_path(root: &std::path::Path) -> std::path::PathBuf { - root.join(".agent/runtime/direct-codex/turns/turn.jsonl") - } - - fn timing_records(root: &std::path::Path) -> Vec { - std::fs::read_to_string(timing_log_path(root)) - .unwrap() - .lines() - .map(|line| serde_json::from_str(line).unwrap()) - .collect() - } - #[tokio::test] - async fn measured_stream_preserves_bytes_and_records_eof_after_fragmented_sse() { - let root = tempfile::tempdir().unwrap(); - let metrics = - super::super::DirectTurnMetrics::new(timing_log_path(root.path()), "turn-stream"); - let attempt = metrics.attempt("gpt-5.6-sol", "gpt-5.6-sol", "high"); - let mut timing = DirectRequestTiming::new(attempt.clone()); - timing.request_body( - br#"{"model":"gpt-5.6-sol","reasoning":{"effort":"high"},"input":"private"}"#, - ); - timing.dispatched(); - timing.headers(200, true); + async fn observed_stream_preserves_fragmented_sse_bytes() { let chunks = [ b"data: {\"type\":\"response.created\",\"response\":{\"model\":\"gpt-5.6-sol\"}}\n\n" .as_slice(), @@ -519,151 +398,41 @@ mod tests { b"ta\":\"private content\"}\n\ndata: {\"type\":\"response.completed\"}\n\n".as_slice(), ]; let expected: Vec = chunks.concat(); - let mut stream = MeasuredResponseStream { + let mut headers = HeaderMap::new(); + headers.insert("content-type", "text/event-stream".parse().unwrap()); + let mut stream = ObservedResponseStream { inner: Box::pin(futures::stream::iter(chunks.into_iter().map(|bytes| { Ok::<_, reqwest::Error>(axum::body::Bytes::copy_from_slice(bytes)) }))), - timing: Some(timing), - observer: None, + observer: ModelResponseObserver::new(None, StatusCode::OK, &headers), }; let mut actual = Vec::new(); while let Some(chunk) = stream.next().await { actual.extend_from_slice(&chunk.unwrap()); } assert_eq!(actual, expected); - drop(stream); - assert!( - metrics.wait_for_test_writes().await, - "writer failed: {}", - metrics.snapshot() - ); - let records = timing_records(root.path()); - let requests: Vec<_> = records - .iter() - .filter(|row| row["recordType"] == "direct.codex.request_timing") - .collect(); - assert_eq!(requests.len(), 1); - let request = requests[0]; - assert_eq!(request["transportStatus"], "eof"); - assert_eq!(request["responseStatus"], "completed"); - assert_eq!(request["responseReportedModel"], "gpt-5.6-sol"); - assert!(request["firstSseEventOffsetMs"].is_number()); - assert!(request["firstContentDeltaOffsetMs"].is_number()); - assert!(!serde_json::to_string(&records).unwrap().contains("private")); - assert_eq!( - metrics.snapshot()["categories"]["http-request"]["activeCount"], - 0 - ); } #[tokio::test] - async fn measured_stream_records_errors_and_unpolled_body_drop_without_fake_first_chunk() { - let root = tempfile::tempdir().unwrap(); - let metrics = - super::super::DirectTurnMetrics::new(timing_log_path(root.path()), "turn-errors"); - let attempt = metrics.attempt("gpt-5.6-sol", "gpt-5.6-sol", "high"); + async fn observed_stream_propagates_upstream_error() { // Invalid URL fails in reqwest's request builder; no network call is made. let error = reqwest::Client::new() .get("not a URL") .send() .await .unwrap_err(); - let mut stream = MeasuredResponseStream { + let mut stream = ObservedResponseStream { inner: Box::pin(futures::stream::iter(vec![Err::( error, )])), - timing: Some(DirectRequestTiming::new(attempt.clone())), - observer: None, + observer: ModelResponseObserver::new(None, StatusCode::OK, &HeaderMap::new()), }; - assert!(stream.next().await.unwrap().is_err()); - drop(stream); - let never_polled = MeasuredResponseStream { - inner: Box::pin(futures::stream::pending::< - Result, - >()), - timing: Some(DirectRequestTiming::new(attempt)), - observer: None, - }; - drop(never_polled); - assert!( - metrics.wait_for_test_writes().await, - "writer failed: {}", - metrics.snapshot() - ); - let records = timing_records(root.path()); - let requests: Vec<_> = records - .iter() - .filter(|row| row["recordType"] == "direct.codex.request_timing") - .collect(); - assert_eq!(requests.len(), 2); - assert_eq!(requests[0]["transportStatus"], "stream-error"); - assert_eq!(requests[1]["transportStatus"], "dropped"); - assert!(requests - .iter() - .all(|row| row["firstBodyChunkOffsetMs"].is_null())); assert_eq!( - metrics.snapshot()["categories"]["http-request"]["activeCount"], - 0 + stream.next().await.unwrap().unwrap_err().to_string(), + "provider response stream failed" ); } - #[tokio::test] - async fn loopback_timing_keeps_original_scope_and_does_not_invent_sse_for_json() { - let root = tempfile::tempdir().unwrap(); - let metrics = - super::super::DirectTurnMetrics::new(timing_log_path(root.path()), "turn-proxy"); - let calls = Arc::new(AtomicUsize::new(0)); - let listener = tokio::net::TcpListener::bind((std::net::Ipv4Addr::LOCALHOST, 0)) - .await - .unwrap(); - let address = listener.local_addr().unwrap(); - let app = Router::new() - .route("/responses", post(fake_upstream)) - .with_state(calls); - let task = tokio::spawn(async move { - let _ = axum::serve(listener, app).await; - }); - let proxy = - start_codex_provider_proxy(&format!("http://{address}"), "fixture-provider-key", false) - .await - .unwrap(); - let first = metrics.attempt("gpt-5.6-sol", "gpt-5.6-sol", "high"); - let binding = proxy.bind_metrics(first.clone()); - let response = reqwest::Client::new() - .post(format!("{}/responses", proxy.base_url())) - .bearer_auth(proxy.downstream_bearer_token()) - .body(r#"{"model":"gpt-5.6-sol","input":"keep secret"}"#) - .send() - .await - .unwrap(); - let second = metrics.attempt("gpt-5.6-sol", "gpt-5.6-sol", "high"); - let _second_binding = proxy.bind_metrics(second.clone()); - drop(binding); - assert_eq!( - proxy.metrics_scope.lock().unwrap().as_ref().unwrap().id(), - second.id() - ); - assert!(response.text().await.unwrap().contains("keep secret")); - assert!( - metrics.wait_for_test_writes().await, - "writer failed: {}", - metrics.snapshot() - ); - let records = timing_records(root.path()); - let request = records - .iter() - .find(|row| row["recordType"] == "direct.codex.request_timing") - .unwrap(); - assert_eq!(request["attemptId"], first.id()); - assert_eq!(request["transportStatus"], "eof"); - assert!(request["firstSseEventOffsetMs"].is_null()); - assert!(request["firstContentDeltaOffsetMs"].is_null()); - assert!(!serde_json::to_string(&records) - .unwrap() - .contains("keep secret")); - task.abort(); - } - #[derive(Clone)] struct ModelFixture { status: StatusCode, diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_codex_audit.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_codex_audit.rs deleted file mode 100644 index 3d0c82c08..000000000 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_codex_audit.rs +++ /dev/null @@ -1,1625 +0,0 @@ -//! Direct Codex GUI 回合行为账本:把 item/completed 抽成项目内有界时间线。 -//! 不灌附件正文、不落 stdout / patch / MCP result,不进入前端观察者。 - -use super::*; -use serde_json::{json, Map, Value}; -use sha2::{Digest, Sha256}; -use std::fs; -use std::path::{Path, PathBuf}; - -const DIRECT_CODEX_AUDIT_HASH_MAX_BYTES: u64 = 2 * 1024 * 1024; -const DIRECT_CODEX_AUDIT_MAX_ITEMS: usize = 256; -const DIRECT_CODEX_AUDIT_COMMAND_CHARS: usize = 240; -const DIRECT_CODEX_AUDIT_BRIEF_CHARS: usize = 4000; -const DIRECT_CODEX_AUDIT_PREVIEW_CHARS: usize = 240; -const DIRECT_CODEX_AUDIT_QUERY_CHARS: usize = 400; -const DIRECT_CODEX_AUDIT_LIST_QUERY_CHARS: usize = 120; -const DIRECT_CODEX_AUDIT_ID_LIST_MAX: usize = 8; -const DIRECT_CODEX_AUDIT_TURN_LOG_DIR: &str = ".agent/runtime/direct-codex/turns"; - -const SKIPPED_ITEM_TYPES: &[&str] = &[ - "agentMessage", - "userMessage", - "plan", - "reasoning", - "contextCompaction", - "hookPrompt", -]; - -const DESIGN_MCP_TOOLS: &[&str] = &[ - "taonier_prepare_game_art", - "agc_generate_image", - "agc_edit_image", - "agc_create_or_derive_resource", -]; - -struct OfferedAttachment { - local_path: String, - content_sha256: Option, - read: bool, - content_sha256_match: Option, -} - -pub(crate) struct DirectCodexTurnAudit { - root: PathBuf, - client_turn_id: String, - turn_log_relative: String, - sidecar_present: bool, - offered: Vec, - item_count: usize, - items_truncated: bool, - truncated_written: bool, - first_design: Option, - audit_write_failed: bool, - finished: bool, - metrics: DirectTurnMetrics, -} - -impl DirectCodexTurnAudit { - pub(crate) fn start( - root: &Path, - client_turn_id: &str, - original_prompt: &str, - attachments: &[DirectCodexTurnAttachment], - ) -> Self { - let turn_log_relative = format!("{DIRECT_CODEX_AUDIT_TURN_LOG_DIR}/{client_turn_id}.jsonl"); - let log_path = root.join(&turn_log_relative); - let metrics = DirectTurnMetrics::new(log_path.clone(), client_turn_id); - let sidecar_present = attachments_use_project_mapping(attachments); - let (attachment_values, offered) = project_audit_attachments(root, attachments); - let mut audit = Self { - root: root.to_path_buf(), - client_turn_id: client_turn_id.to_string(), - turn_log_relative, - sidecar_present, - offered, - item_count: 0, - items_truncated: false, - truncated_written: false, - first_design: None, - audit_write_failed: false, - finished: false, - metrics, - }; - let omitted = attachments - .len() - .saturating_sub(MAX_DIRECT_CODEX_ATTACHMENTS); - let mut record = json!({ - "recordType": "direct.codex.turn_start", - "clientTurnId": client_turn_id, - "sidecarPresent": sidecar_present, - "promptSha256": sha256_hex(original_prompt.as_bytes()), - "promptChars": original_prompt.chars().count(), - "attachments": attachment_values, - }); - if omitted > 0 { - record["attachmentsOmitted"] = json!(omitted); - } - audit.append_record(record); - audit - } - - pub(crate) fn metrics(&self) -> DirectTurnMetrics { - self.metrics.clone() - } - - pub(crate) async fn flush(&self) { - self.metrics.flush().await; - } - - pub(crate) fn observe_item(&mut self, params: &Value) { - if self.finished { - return; - } - let Some(item) = params.get("item") else { - return; - }; - let item_type = item - .get("type") - .and_then(Value::as_str) - .unwrap_or("unknown"); - if SKIPPED_ITEM_TYPES.contains(&item_type) { - return; - } - if self.item_count >= DIRECT_CODEX_AUDIT_MAX_ITEMS { - self.items_truncated = true; - if !self.truncated_written { - self.truncated_written = true; - self.append_record(json!({ - "recordType": "direct.codex.items_truncated", - "clientTurnId": self.client_turn_id, - "droppedAfter": DIRECT_CODEX_AUDIT_MAX_ITEMS, - })); - } - return; - } - - self.item_count = self.item_count.saturating_add(1); - let seq = self.item_count; - let mut record = json!({ - "recordType": "direct.codex.item", - "clientTurnId": self.client_turn_id, - "seq": seq, - "itemType": item_type, - }); - if let Some(item_id) = item - .get("id") - .and_then(Value::as_str) - .filter(|id| !id.is_empty()) - { - record["itemId"] = json!(item_id); - } - if let Some(status) = item.get("status").and_then(Value::as_str) { - record["status"] = json!(status); - } else { - record["status"] = json!("completed"); - } - - match item_type { - "commandExecution" => self.fill_command_execution(&mut record, item), - "mcpToolCall" => self.fill_mcp_tool_call(&mut record, item, seq), - "fileChange" => self.fill_file_change(&mut record, item, seq), - "imageView" => self.fill_image_view(&mut record, item), - "functionCallOutput" => fill_function_call_output(&mut record, item), - "webSearch" => fill_web_search(&mut record, item), - _ => {} - } - - self.append_record(record); - } - - pub(crate) fn finish(&mut self, completed: bool) { - if self.finished { - return; - } - self.finished = true; - let offered_read = self.offered_read_values(); - let record = json!({ - "recordType": "direct.codex.turn_end", - "clientTurnId": self.client_turn_id, - "completed": completed, - "itemCount": self.item_count, - "itemsTruncated": self.items_truncated, - "offeredRead": offered_read, - "firstDesign": self.first_design.clone(), - }); - self.append_record(record); - #[cfg(test)] - self.metrics.flush_for_test(); - let summary = json!({ - "recordType": "direct.codex.turn", - "clientTurnId": self.client_turn_id, - "turnLog": self.turn_log_relative, - "sidecarPresent": self.sidecar_present, - "offeredCount": self.offered.len(), - "offeredRead": offered_read, - "firstDesign": self.first_design.clone(), - "itemCount": self.item_count, - "itemsTruncated": self.items_truncated, - "completed": completed, - "auditWriteFailed": self.audit_write_failed, - }); - if append_agent_db_record(&self.root, summary).is_err() { - self.audit_write_failed = true; - } - } - - fn fill_command_execution(&mut self, record: &mut Value, item: &Value) { - let mut path_rejected = false; - let mut actions = Vec::new(); - if let Some(raw_actions) = item.get("commandActions").and_then(Value::as_array) { - for action in raw_actions { - let action_type = action - .get("type") - .and_then(Value::as_str) - .unwrap_or("unknown"); - match action_type { - "read" => { - let (entry, rejected) = self.read_action_entry(action); - path_rejected |= rejected; - actions.push(entry); - } - "listFiles" => { - let mut entry = json!({ "type": "listFiles" }); - match optional_action_path(&self.root, action) { - ActionPath::Missing => {} - ActionPath::Rejected => { - path_rejected = true; - entry["pathRejected"] = json!(true); - } - ActionPath::Ok(path) => entry["path"] = json!(path), - } - actions.push(entry); - } - "search" => { - let mut entry = json!({ "type": "search" }); - if let Some(query) = action.get("query").and_then(Value::as_str) { - entry["query"] = - json!(truncate_chars(query, DIRECT_CODEX_AUDIT_QUERY_CHARS)); - } - match optional_action_path(&self.root, action) { - ActionPath::Missing => {} - ActionPath::Rejected => { - path_rejected = true; - entry["pathRejected"] = json!(true); - } - ActionPath::Ok(path) => entry["path"] = json!(path), - } - actions.push(entry); - } - _ => actions.push(json!({ "type": "unknown" })), - } - } - } - record["actions"] = json!(actions); - if let Some(exit_code) = item.get("exitCode").and_then(Value::as_i64) { - record["exitCode"] = json!(exit_code); - } - if let Some(duration_ms) = item.get("durationMs").and_then(Value::as_i64) { - record["durationMs"] = json!(duration_ms); - } - let command = item.get("command").and_then(Value::as_str).unwrap_or(""); - if path_rejected || command_contains_host_absolute_path(command) { - record["commandRedacted"] = json!(true); - } else if !command.is_empty() { - record["command"] = json!(truncate_chars(command, DIRECT_CODEX_AUDIT_COMMAND_CHARS)); - } - } - - fn read_action_entry(&mut self, action: &Value) -> (Value, bool) { - let Some(raw_path) = action.get("path").and_then(Value::as_str) else { - return (json!({ "type": "read", "pathRejected": true }), true); - }; - match relativize_project_path(&self.root, raw_path) { - Some(path) => { - let (content_sha256, hash_skipped) = hash_project_file(&self.root, &path); - self.mark_offered_read(&path, content_sha256.as_deref()); - let mut entry = json!({ "type": "read", "path": path }); - insert_hash_fields(&mut entry, content_sha256, hash_skipped); - (entry, false) - } - None => (json!({ "type": "read", "pathRejected": true }), true), - } - } - - fn fill_mcp_tool_call(&mut self, record: &mut Value, item: &Value, seq: usize) { - let tool = item.get("tool").and_then(Value::as_str).unwrap_or(""); - record["tool"] = json!(tool); - if let Some(server) = item - .get("server") - .and_then(Value::as_str) - .filter(|server| !server.is_empty() && *server != "agc_tools") - { - record["server"] = json!(server); - } - if let Some(duration_ms) = item.get("durationMs").and_then(Value::as_i64) { - record["durationMs"] = json!(duration_ms); - } - if item.get("error").is_some_and(|error| !error.is_null()) { - record["errorKind"] = json!(item - .pointer("/error/code") - .and_then(Value::as_str) - .or_else(|| item.pointer("/error/type").and_then(Value::as_str)) - .unwrap_or("error")); - } - let arguments = item.get("arguments").cloned().unwrap_or(Value::Null); - let extracted = extract_mcp_arguments(&self.root, tool, &arguments); - if let Some(path) = extracted - .get("path") - .and_then(Value::as_str) - .map(str::to_string) - { - self.mark_offered_read(&path, None); - } - if let Some(local_paths) = extracted.get("localPaths").and_then(Value::as_array) { - for path in local_paths { - if let Some(path) = path.as_str() { - self.mark_offered_read(path, None); - } - } - } - if extracted - .as_object() - .is_some_and(|object| !object.is_empty()) - { - record["arguments"] = extracted.clone(); - } - if self.first_design.is_none() { - if DESIGN_MCP_TOOLS.contains(&tool) { - let mut design = json!({ - "kind": format!("mcp:{tool}"), - "seq": seq, - "tool": tool, - }); - let preview = extracted - .get("brief") - .or_else(|| extracted.get("prompt")) - .and_then(Value::as_str) - .map(|text| truncate_chars(text, DIRECT_CODEX_AUDIT_PREVIEW_CHARS)); - if let Some(preview) = preview { - design["briefPreview"] = json!(preview); - } - self.first_design = Some(design); - } else if tool == "agc_write_file" { - if let Some(path) = extracted.get("path").and_then(Value::as_str) { - if is_design_write_path(path) { - self.first_design = Some(json!({ - "kind": format!("write:{path}"), - "seq": seq, - "path": path, - })); - } - } - } - } - } - - fn fill_file_change(&mut self, record: &mut Value, item: &Value, seq: usize) { - let mut changes = Vec::new(); - if let Some(raw_changes) = item.get("changes").and_then(Value::as_array) { - for change in raw_changes { - let kind = file_change_kind(change); - let mut entry = json!({ "kind": kind }); - match change.get("path").and_then(Value::as_str) { - Some(raw) => match relativize_project_path(&self.root, raw) { - Some(path) => { - if self.first_design.is_none() && is_design_write_path(&path) { - self.first_design = Some(json!({ - "kind": format!("fileChange:{path}"), - "seq": seq, - "path": path, - })); - } - entry["path"] = json!(path); - } - None => entry["pathRejected"] = json!(true), - }, - None => entry["pathRejected"] = json!(true), - } - changes.push(entry); - } - } - record["changes"] = json!(changes); - } - - fn fill_image_view(&mut self, record: &mut Value, item: &Value) { - match item.get("path").and_then(Value::as_str) { - Some(raw) => match relativize_project_path(&self.root, raw) { - Some(path) => { - let (content_sha256, hash_skipped) = hash_project_file(&self.root, &path); - self.mark_offered_read(&path, content_sha256.as_deref()); - record["path"] = json!(path); - insert_hash_fields(record, content_sha256, hash_skipped); - } - None => record["pathRejected"] = json!(true), - }, - None => record["pathRejected"] = json!(true), - } - } - - fn mark_offered_read(&mut self, path: &str, content_sha256: Option<&str>) { - for offered in &mut self.offered { - if offered.local_path != path { - continue; - } - offered.read = true; - match (offered.content_sha256.as_deref(), content_sha256) { - (Some(expected), Some(actual)) => { - let matches = expected == actual; - offered.content_sha256_match = - Some(offered.content_sha256_match.unwrap_or(true) && matches); - } - _ => {} - } - } - } - - fn offered_read_values(&self) -> Vec { - self.offered - .iter() - .map(|offered| { - let mut value = json!({ - "localPath": offered.local_path, - "read": offered.read, - }); - if let Some(matches) = offered.content_sha256_match { - value["contentSha256Match"] = json!(matches); - } - value - }) - .collect() - } - - fn append_record(&mut self, mut record: Value) { - #[cfg(test)] - if test_fail_audit_write(&self.root) { - self.audit_write_failed = true; - return; - } - if let Some(object) = record.as_object_mut() { - object.insert( - "recordedAtMs".to_string(), - json!(u64::try_from(unix_millis()).unwrap_or(u64::MAX)), - ); - } - if !self.metrics.append_audit_record(record) { - self.audit_write_failed = true; - } - } -} - -impl Drop for DirectCodexTurnAudit { - fn drop(&mut self) { - if !self.finished { - self.finish(false); - } - } -} - -enum ActionPath { - Missing, - Rejected, - Ok(String), -} - -fn optional_action_path(root: &Path, action: &Value) -> ActionPath { - let Some(raw) = action.get("path").and_then(Value::as_str) else { - return ActionPath::Missing; - }; - if raw.trim().is_empty() { - return ActionPath::Missing; - } - match relativize_project_path(root, raw) { - Some(path) => ActionPath::Ok(path), - None => ActionPath::Rejected, - } -} - -fn project_audit_attachments( - root: &Path, - attachments: &[DirectCodexTurnAttachment], -) -> (Vec, Vec) { - let mut values = Vec::new(); - let mut offered = Vec::new(); - for attachment in attachments.iter().take(MAX_DIRECT_CODEX_ATTACHMENTS) { - let name = sanitize_attachment_name(&attachment.name); - let media_type = sanitize_attachment_media_type(&attachment.media_type); - let raw_path = attachment - .local_path - .as_deref() - .map(str::trim) - .filter(|value| !value.is_empty()); - let sanitized_path = raw_path.and_then(sanitize_attachment_local_path); - let path_rejected = raw_path.is_some() && sanitized_path.is_none(); - let status = if path_rejected { - Some("failed") - } else { - sanitize_attachment_status(attachment.status.as_deref()) - }; - let mut value = json!({ - "name": name, - "mediaType": media_type, - "size": attachment.size, - }); - if let Some(path) = sanitized_path { - value["localPath"] = json!(path.clone()); - let (content_sha256, hash_skipped) = hash_project_file(root, &path); - insert_hash_fields(&mut value, content_sha256.clone(), hash_skipped); - offered.push(OfferedAttachment { - local_path: path, - content_sha256, - read: false, - content_sha256_match: None, - }); - } - if let Some(status) = status { - value["status"] = json!(status); - } - values.push(value); - } - (values, offered) -} - -fn extract_mcp_arguments(root: &Path, tool: &str, arguments: &Value) -> Value { - let Some(object) = arguments.as_object() else { - return json!({}); - }; - let mut out = Map::new(); - match tool { - "agc_list_project_files" => { - copy_sanitized_path(root, object, "path", &mut out); - copy_truncated_string( - object, - "query", - DIRECT_CODEX_AUDIT_LIST_QUERY_CHARS, - &mut out, - ); - copy_string(object, "kind", &mut out); - copy_number(object, "offset", &mut out); - copy_number(object, "limit", &mut out); - } - "agc_write_file" => { - copy_sanitized_path(root, object, "path", &mut out); - if let Some(content) = object.get("content").and_then(Value::as_str) { - out.insert("contentChars".to_string(), json!(content.chars().count())); - } - } - "agc_apply_patch" => { - if let Some(patch) = object.get("patch").and_then(Value::as_str) { - out.insert("patchHash".to_string(), json!(sha256_hex(patch.as_bytes()))); - out.insert("patchBytes".to_string(), json!(patch.len())); - out.insert("patchChars".to_string(), json!(patch.chars().count())); - // 只提取有界语法计数,不保留路径、上下文行或补丁正文。 - if patch.len() <= 64 * 1024 { - if let Ok(parsed) = codex_patch_parser::parse_patch(patch) { - out.insert("patchOperations".to_string(), json!(parsed.hunks.len())); - } - } - } - } - "agc_update_plan" => { - if let Some(plan) = object.get("plan").and_then(Value::as_array) { - out.insert("planSteps".to_string(), json!(plan.len())); - } - // 摘要包含 explanation 与完整 plan;审计中不保存任何自然语言预览。 - if let Ok(bytes) = serde_json::to_vec(arguments) { - out.insert("planHash".to_string(), json!(sha256_hex(&bytes))); - out.insert("planBytes".to_string(), json!(bytes.len())); - } - } - "taonier_prepare_game_art" => { - copy_string(object, "mode", &mut out); - copy_text_with_hash( - object, - "brief", - "brief", - DIRECT_CODEX_AUDIT_BRIEF_CHARS, - &mut out, - ); - } - "agc_generate_image" => { - copy_string(object, "kind", &mut out); - copy_string(object, "sliceMode", &mut out); - copy_number(object, "sliceCount", &mut out); - copy_string(object, "screenColor", &mut out); - copy_string(object, "aspectRatio", &mut out); - copy_string(object, "imageSize", &mut out); - copy_string(object, "assetName", &mut out); - copy_sanitized_path(root, object, "outputPath", &mut out); - copy_text_with_hash( - object, - "prompt", - "prompt", - DIRECT_CODEX_AUDIT_BRIEF_CHARS, - &mut out, - ); - } - "agc_edit_image" => { - copy_string(object, "sourceLocalAssetId", &mut out); - copy_string(object, "assetName", &mut out); - copy_text_with_hash( - object, - "prompt", - "prompt", - DIRECT_CODEX_AUDIT_BRIEF_CHARS, - &mut out, - ); - } - "agc_create_or_derive_resource" => { - copy_string(object, "kind", &mut out); - copy_string(object, "mode", &mut out); - copy_string(object, "sourceLocalAssetId", &mut out); - copy_string(object, "assetName", &mut out); - copy_text_with_hash( - object, - "prompt", - "prompt", - DIRECT_CODEX_AUDIT_BRIEF_CHARS, - &mut out, - ); - } - "agc_list_registered_assets" => { - copy_string(object, "kind", &mut out); - copy_string(object, "assetId", &mut out); - if let Some(flag) = object.get("includeSequenceFrames").and_then(Value::as_bool) { - out.insert("includeSequenceFrames".to_string(), json!(flag)); - } - copy_number(object, "offset", &mut out); - copy_number(object, "limit", &mut out); - } - "agc_list_account_assets" => { - copy_string(object, "folderId", &mut out); - copy_truncated_string( - object, - "query", - DIRECT_CODEX_AUDIT_LIST_QUERY_CHARS, - &mut out, - ); - copy_number(object, "offset", &mut out); - copy_number(object, "limit", &mut out); - } - "agc_import_account_assets" => { - if let Some(ids) = object.get("assetIds").and_then(Value::as_array) { - let kept: Vec = ids - .iter() - .filter_map(Value::as_str) - .take(DIRECT_CODEX_AUDIT_ID_LIST_MAX) - .map(Value::from) - .collect(); - let omitted = ids.len().saturating_sub(kept.len()); - out.insert("assetIds".to_string(), json!(kept)); - if omitted > 0 { - out.insert("assetIdsOmitted".to_string(), json!(omitted)); - } - } - if let Some(paths) = object.get("localPaths").and_then(Value::as_array) { - let kept: Vec = paths - .iter() - .filter_map(Value::as_str) - .filter_map(|path| relativize_project_path(root, path)) - .take(DIRECT_CODEX_AUDIT_ID_LIST_MAX) - .map(Value::from) - .collect(); - out.insert("localPaths".to_string(), json!(kept)); - } - } - "agc_remove_background" => { - copy_string(object, "sourceLocalAssetId", &mut out); - copy_string(object, "assetName", &mut out); - } - "agc_browser_playtest" => { - for (key, allowed) in [ - ("mode", &["visual", "gameplay"][..]), - ( - "scenario", - &["generic-v1", "tetris-v1", "lane-defense-v1"][..], - ), - ] { - if let Some(value) = object - .get(key) - .and_then(Value::as_str) - .filter(|value| allowed.contains(value)) - { - out.insert(key.to_string(), json!(value)); - } - } - } - "agc_environment_check" => {} - "agc_run_validation" => { - if let Some(program) = object - .get("program") - .and_then(Value::as_str) - .filter(|value| matches!(*value, "node" | "npm")) - { - out.insert("program".to_string(), json!(program)); - } - copy_sanitized_path(root, object, "cwd", &mut out); - copy_number(object, "timeoutSeconds", &mut out); - if let Some(args) = object.get("arguments").and_then(Value::as_array) { - out.insert("argsCount".to_string(), json!(args.len())); - if let Ok(bytes) = serde_json::to_vec(args) { - out.insert("argsHash".to_string(), json!(sha256_hex(&bytes))); - } - } - } - "agc_web_search" => { - copy_truncated_string(object, "query", DIRECT_CODEX_AUDIT_QUERY_CHARS, &mut out); - copy_number(object, "maxResults", &mut out); - } - _ => {} - } - Value::Object(out) -} - -fn fill_function_call_output(record: &mut Value, item: &Value) { - if let Some(name) = item.get("name").and_then(Value::as_str) { - record["name"] = json!(name); - } - if let Some(namespace) = item.get("namespace").and_then(Value::as_str) { - record["namespace"] = json!(namespace); - } -} - -fn fill_web_search(record: &mut Value, item: &Value) { - if let Some(query) = item.get("query").and_then(Value::as_str) { - record["query"] = json!(truncate_chars(query, DIRECT_CODEX_AUDIT_QUERY_CHARS)); - } -} - -fn copy_string(source: &Map, key: &str, out: &mut Map) { - if let Some(value) = source - .get(key) - .and_then(Value::as_str) - .filter(|value| !value.is_empty()) - { - out.insert(key.to_string(), json!(value)); - } -} - -fn copy_truncated_string( - source: &Map, - key: &str, - max_chars: usize, - out: &mut Map, -) { - if let Some(value) = source.get(key).and_then(Value::as_str) { - out.insert(key.to_string(), json!(truncate_chars(value, max_chars))); - } -} - -fn copy_number(source: &Map, key: &str, out: &mut Map) { - if let Some(value) = source.get(key).and_then(Value::as_i64) { - out.insert(key.to_string(), json!(value)); - } -} - -fn copy_sanitized_path( - root: &Path, - source: &Map, - key: &str, - out: &mut Map, -) { - let Some(raw) = source.get(key).and_then(Value::as_str) else { - return; - }; - match relativize_project_path(root, raw) { - Some(path) => { - out.insert(key.to_string(), json!(path)); - } - None => { - out.insert(format!("{key}Rejected"), json!(true)); - } - } -} - -fn copy_text_with_hash( - source: &Map, - source_key: &str, - dest_key: &str, - max_chars: usize, - out: &mut Map, -) { - let Some(text) = source.get(source_key).and_then(Value::as_str) else { - return; - }; - out.insert(format!("{dest_key}Chars"), json!(text.chars().count())); - out.insert( - format!("{dest_key}Sha256"), - json!(sha256_hex(text.as_bytes())), - ); - out.insert(dest_key.to_string(), json!(truncate_chars(text, max_chars))); -} - -fn file_change_kind(change: &Value) -> &'static str { - let kind = change.get("kind"); - let label = kind - .and_then(Value::as_str) - .or_else(|| { - kind.and_then(|value| value.get("type")) - .and_then(Value::as_str) - }) - .unwrap_or("update"); - match label { - "add" => "add", - "delete" => "delete", - _ => "update", - } -} - -fn is_design_write_path(path: &str) -> bool { - path == "index.html" - || path.starts_with("game/") - || path.rsplit('/').next() == Some("index.html") -} - -fn insert_hash_fields( - target: &mut Value, - content_sha256: Option, - hash_skipped: Option<&str>, -) { - if let Some(content_sha256) = content_sha256 { - target["contentSha256"] = json!(content_sha256); - } - if let Some(hash_skipped) = hash_skipped { - target["hashSkipped"] = json!(hash_skipped); - } -} - -fn sha256_hex(bytes: &[u8]) -> String { - format!("{:x}", Sha256::digest(bytes)) -} - -fn truncate_chars(value: &str, max_chars: usize) -> String { - value.chars().take(max_chars).collect() -} - -fn hash_project_file(root: &Path, relative: &str) -> (Option, Option<&'static str>) { - if reject_agent_runtime_private_control_path(relative).is_err() - || reject_sensitive_project_file_read(relative).is_err() - { - return (None, Some("missing")); - } - let path = match resolve_local_project_path(root, relative) { - Ok(path) => path, - Err(_) => return (None, Some("missing")), - }; - let metadata = match fs::metadata(&path) { - Ok(metadata) if metadata.is_file() => metadata, - _ => return (None, Some("missing")), - }; - if metadata.len() > DIRECT_CODEX_AUDIT_HASH_MAX_BYTES { - return (None, Some("too-large")); - } - match fs::read(&path) { - Ok(bytes) => (Some(sha256_hex(&bytes)), None), - Err(_) => (None, Some("missing")), - } -} - -fn posix_path_text(path: &Path) -> String { - let text = path.to_string_lossy(); - let text = text - .strip_prefix(r"\\?\") - .or_else(|| text.strip_prefix("//?/")) - .unwrap_or(&text); - text.replace('\\', "/").trim_end_matches('/').to_string() -} - -fn relativize_project_path(root: &Path, raw: &str) -> Option { - let trimmed = raw.trim(); - if trimmed.is_empty() { - return None; - } - if let Some(relative) = sanitize_attachment_local_path(trimmed) { - return accept_relative_path(&relative); - } - if let Some(relative) = strip_absolute_root_prefix(root, trimmed) { - return sanitize_attachment_local_path(&relative) - .and_then(|path| accept_relative_path(&path)); - } - None -} - -fn strip_absolute_root_prefix(root: &Path, raw: &str) -> Option { - if let (Ok(root_canon), Ok(raw_canon)) = (root.canonicalize(), Path::new(raw).canonicalize()) { - if let Ok(stripped) = raw_canon.strip_prefix(&root_canon) { - let relative = posix_path_text(stripped); - if !relative.is_empty() { - return Some(relative); - } - } - } - let root_text = posix_path_text(root); - let raw_text = posix_path_text(Path::new(raw)); - let rest = if cfg!(windows) { - let root_lower = root_text.to_ascii_lowercase(); - let raw_lower = raw_text.to_ascii_lowercase(); - let suffix = raw_lower.strip_prefix(&root_lower)?; - raw_text - .get(raw_text.len().saturating_sub(suffix.len())..) - .unwrap_or(suffix) - .to_string() - } else { - raw_text.strip_prefix(&root_text)?.to_string() - }; - let rest = rest.trim_start_matches('/').to_string(); - (!rest.is_empty()).then_some(rest) -} - -fn accept_relative_path(relative: &str) -> Option { - if reject_agent_runtime_private_control_path(relative).is_err() - || reject_sensitive_project_file_read(relative).is_err() - { - return None; - } - Some(relative.to_string()) -} - -fn command_contains_host_absolute_path(command: &str) -> bool { - if command.contains("\\\\") || command.contains("/Users/") || command.contains("/home/") { - return true; - } - let bytes = command.as_bytes(); - let mut index = 0; - while index + 2 < bytes.len() { - if bytes[index].is_ascii_alphabetic() - && bytes[index + 1] == b':' - && matches!(bytes[index + 2], b'\\' | b'/') - { - return true; - } - index += 1; - } - false -} - -#[cfg(test)] -fn test_fail_audit_write(root: &Path) -> bool { - root.join(".agent/runtime/test-fail-direct-codex-audit") - .is_file() -} - -#[cfg(test)] -mod tests { - use super::*; - - fn fixture_project(name: &str) -> tempfile::TempDir { - let directory = tempfile::tempdir().expect("temp project"); - init_local_game_project_at(directory.path(), name, "审计测试项目").expect("init project"); - directory - } - - fn attachment_json( - name: &str, - media_type: &str, - size: u64, - local_path: Option<&str>, - status: Option<&str>, - ) -> DirectCodexTurnAttachment { - let mut value = json!({ - "name": name, - "mediaType": media_type, - "size": size, - }); - if let Some(local_path) = local_path { - value["localPath"] = json!(local_path); - } - if let Some(status) = status { - value["status"] = json!(status); - } - serde_json::from_value(value).expect("attachment") - } - - fn read_turn_log(root: &Path, client_turn_id: &str) -> Vec { - let path = root - .join(DIRECT_CODEX_AUDIT_TURN_LOG_DIR) - .join(format!("{client_turn_id}.jsonl")); - fs::read_to_string(path) - .unwrap_or_default() - .lines() - .filter(|line| !line.trim().is_empty()) - .map(|line| serde_json::from_str::(line).expect("audit jsonl")) - .collect() - } - - fn read_agent_db(root: &Path) -> Vec { - fs::read_to_string(root.join(".agent/agent.db")) - .unwrap_or_default() - .lines() - .filter(|line| !line.trim().is_empty()) - .filter_map(|line| serde_json::from_str::(line).ok()) - .collect() - } - - fn start_audit( - root: &Path, - prompt: &str, - attachments: &[DirectCodexTurnAttachment], - ) -> DirectCodexTurnAudit { - DirectCodexTurnAudit::start(root, "turn-01", prompt, attachments) - } - - #[test] - fn turn_start_hashes_original_prompt_and_sanitizes_attachment_paths() { - let project = fixture_project("audit-start"); - let root = project.path(); - let upload = "assets/uploads/upload-1-fast_gdd.md"; - fs::create_dir_all(root.join("assets/uploads")).expect("uploads dir"); - fs::write(root.join(upload), "脉冲余烬").expect("write gdd"); - let attachments = vec![ - attachment_json( - "fast_gdd.md", - "text/markdown", - 12, - Some(upload), - Some("imported"), - ), - attachment_json( - "secret.md", - "text/markdown", - 1, - Some("../secret.md"), - Some("imported"), - ), - ]; - let mut audit = start_audit(root, "请根据附件做游戏", &attachments); - audit.finish(true); - let records = read_turn_log(root, "turn-01"); - let start = records - .iter() - .find(|record| record["recordType"] == "direct.codex.turn_start") - .expect("turn_start"); - assert_eq!( - start["promptSha256"], - json!(sha256_hex("请根据附件做游戏".as_bytes())) - ); - assert_eq!(start["sidecarPresent"], json!(true)); - let listed = start["attachments"].as_array().expect("attachments"); - assert_eq!(listed[0]["localPath"], json!(upload)); - assert_eq!( - listed[0]["contentSha256"], - json!(sha256_hex("脉冲余烬".as_bytes())) - ); - assert!(listed[1].get("localPath").is_none()); - assert_eq!(listed[1]["status"], json!("failed")); - assert!(!serde_json::to_string(start).expect("json").contains("..")); - } - - #[test] - fn turn_start_without_attachments_sets_sidecar_absent() { - let project = fixture_project("audit-empty"); - let mut audit = start_audit(project.path(), "继续改游戏", &[]); - audit.finish(true); - let start = &read_turn_log(project.path(), "turn-01")[0]; - assert_eq!(start["sidecarPresent"], json!(false)); - assert_eq!(start["attachments"], json!([])); - } - - #[test] - fn attachment_hash_skips_missing_and_too_large_files() { - let project = fixture_project("audit-hash"); - let root = project.path(); - fs::create_dir_all(root.join("assets/uploads")).expect("uploads"); - let large_path = "assets/uploads/upload-1-big.bin"; - let missing_path = "assets/uploads/upload-1-missing.md"; - fs::write( - root.join(large_path), - vec![0_u8; (DIRECT_CODEX_AUDIT_HASH_MAX_BYTES as usize) + 1], - ) - .expect("large file"); - let attachments = vec![ - attachment_json( - "big.bin", - "application/octet-stream", - 3, - Some(large_path), - Some("imported"), - ), - attachment_json( - "missing.md", - "text/markdown", - 1, - Some(missing_path), - Some("imported"), - ), - ]; - let mut audit = start_audit(root, "x", &attachments); - audit.finish(true); - let start = &read_turn_log(root, "turn-01")[0]; - let listed = start["attachments"].as_array().expect("attachments"); - assert_eq!(listed[0]["hashSkipped"], json!("too-large")); - assert!(listed[0].get("contentSha256").is_none()); - assert_eq!(listed[1]["hashSkipped"], json!("missing")); - } - - #[test] - fn command_read_relativizes_absolute_path_and_drops_stdout() { - let project = fixture_project("audit-read"); - let root = project.path(); - let relative = "assets/uploads/upload-1-fast_gdd.md"; - fs::create_dir_all(root.join("assets/uploads")).expect("uploads"); - fs::write(root.join(relative), "裂脉炮").expect("gdd"); - let absolute = root.join(relative); - let mut audit = start_audit( - root, - "做游戏", - &[attachment_json( - "fast_gdd.md", - "text/markdown", - 9, - Some(relative), - Some("imported"), - )], - ); - audit.observe_item(&json!({ - "item": { - "id": "item-read", - "type": "commandExecution", - "command": "type assets/uploads/upload-1-fast_gdd.md", - "status": "completed", - "exitCode": 0, - "aggregatedOutput": "裂脉炮 SECRET", - "commandActions": [{ - "type": "read", - "name": "fast_gdd.md", - "path": absolute.to_string_lossy(), - "command": format!("type {}", absolute.display()) - }] - } - })); - audit.finish(true); - let item = read_turn_log(root, "turn-01") - .into_iter() - .find(|record| record["recordType"] == "direct.codex.item") - .expect("item"); - let dumped = serde_json::to_string(&item).expect("item json"); - assert!(!dumped.contains("aggregatedOutput")); - assert!(!dumped.contains("裂脉炮 SECRET")); - assert_eq!(item["actions"][0]["path"], json!(relative)); - assert_eq!( - item["actions"][0]["contentSha256"], - json!(sha256_hex("裂脉炮".as_bytes())) - ); - let end = read_turn_log(root, "turn-01") - .into_iter() - .find(|record| record["recordType"] == "direct.codex.turn_end") - .expect("end"); - assert_eq!(end["offeredRead"][0]["read"], json!(true)); - assert_eq!(end["offeredRead"][0]["contentSha256Match"], json!(true)); - } - - #[test] - fn rejected_read_path_does_not_persist_host_absolute_path() { - let project = fixture_project("audit-reject"); - let mut audit = start_audit(project.path(), "x", &[]); - audit.observe_item(&json!({ - "item": { - "type": "commandExecution", - "command": r"type C:\Users\secret\fast_gdd.md", - "aggregatedOutput": "nope", - "commandActions": [{ - "type": "read", - "path": r"C:\Users\secret\fast_gdd.md" - }] - } - })); - audit.finish(true); - let dumped = fs::read_to_string( - project - .path() - .join(DIRECT_CODEX_AUDIT_TURN_LOG_DIR) - .join("turn-01.jsonl"), - ) - .expect("log"); - assert!(!dumped.contains(r"C:\Users")); - assert!(!dumped.contains("Users")); - let item = read_turn_log(project.path(), "turn-01") - .into_iter() - .find(|record| record["recordType"] == "direct.codex.item") - .expect("item"); - assert_eq!(item["actions"][0]["pathRejected"], json!(true)); - assert_eq!(item["commandRedacted"], json!(true)); - assert!(item.get("command").is_none()); - assert!(item.get("aggregatedOutput").is_none()); - } - - #[test] - fn mcp_art_brief_is_kept_and_result_is_dropped() { - let project = fixture_project("audit-art"); - let mut audit = start_audit(project.path(), "做游戏", &[]); - audit.observe_item(&json!({ - "item": { - "id": "art-1", - "type": "mcpToolCall", - "server": "agc_tools", - "tool": "taonier_prepare_game_art", - "status": "completed", - "arguments": { "brief": "俯视角收集冒险小游戏", "mode": "reuse-or-create" }, - "result": { "secret": "do-not-store" }, - "error": null - } - })); - audit.finish(true); - let item = read_turn_log(project.path(), "turn-01") - .into_iter() - .find(|record| record["recordType"] == "direct.codex.item") - .expect("item"); - assert_eq!(item["arguments"]["brief"], json!("俯视角收集冒险小游戏")); - assert!(item.get("result").is_none()); - assert!(item.get("errorKind").is_none(), "error:null 不能被记为失败"); - let dumped = serde_json::to_string(&item).expect("json"); - assert!(!dumped.contains("do-not-store")); - let end = read_turn_log(project.path(), "turn-01") - .into_iter() - .find(|record| record["recordType"] == "direct.codex.turn_end") - .expect("end"); - assert_eq!( - end["firstDesign"]["kind"], - json!("mcp:taonier_prepare_game_art") - ); - assert_eq!( - end["firstDesign"]["briefPreview"], - json!("俯视角收集冒险小游戏") - ); - assert_eq!(end["timing"]["schemaVersion"], "agc-direct-timing.v1"); - assert!(end["timing"]["modelInferenceMs"].is_null()); - } - - #[test] - fn agc_write_file_keeps_path_and_content_chars_not_body() { - let project = fixture_project("audit-write"); - let mut audit = start_audit(project.path(), "x", &[]); - audit.observe_item(&json!({ - "item": { - "type": "mcpToolCall", - "tool": "agc_write_file", - "arguments": { - "path": "game/index.html", - "content": "秘密正文" - } - } - })); - audit.finish(true); - let item = read_turn_log(project.path(), "turn-01") - .into_iter() - .find(|record| record["recordType"] == "direct.codex.item") - .expect("item"); - assert_eq!(item["arguments"]["path"], json!("game/index.html")); - assert_eq!( - item["arguments"]["contentChars"], - json!("秘密正文".chars().count()) - ); - let dumped = serde_json::to_string(&item).expect("json"); - assert!(!dumped.contains("秘密正文")); - assert!(item["arguments"].get("content").is_none()); - } - - #[test] - fn validation_audit_keeps_only_safe_mode_and_hashed_process_arguments() { - let project = fixture_project("audit-validation"); - let result = extract_mcp_arguments( - project.path(), - "agc_run_validation", - &json!({ - "program":"node", "cwd":"game", "timeoutSeconds":60, - "arguments":["verify.mjs", "--token", "sk-private-fixture", "https://private.example/?key=secret"] - }), - ); - assert_eq!(result["program"], "node"); - assert_eq!(result["cwd"], "game"); - assert_eq!(result["timeoutSeconds"], 60); - assert_eq!(result["argsCount"], 4); - assert_eq!(result["argsHash"].as_str().unwrap().len(), 64); - let text = serde_json::to_string(&result).unwrap(); - assert!(!text.contains("private")); - assert!(!text.contains("secret")); - assert!(result.get("args").is_none()); - assert!(result.get("arguments").is_none()); - let playtest = extract_mcp_arguments( - project.path(), - "agc_browser_playtest", - &json!({ - "attempt":-999, "mode":"visual", "scenario":"generic-v1" - }), - ); - assert_eq!(playtest, json!({"mode":"visual", "scenario":"generic-v1"})); - assert_eq!( - extract_mcp_arguments( - project.path(), - "agc_environment_check", - &json!({"token":"private"}) - ), - json!({}) - ); - } - - #[test] - fn host_patch_and_plan_audit_persists_only_hashes_and_counts() { - let project = fixture_project("audit-host-edit"); - let patch = "*** Begin Patch\n*** Add File: game/private-path-sentinel.txt\n+sk-patch-body-sentinel\n*** End Patch"; - let plan = json!({ - "explanation":"private-explanation-sentinel", - "plan":[{"step":"sk-plan-step-sentinel","status":"completed"}] - }); - let mut audit = start_audit(project.path(), "x", &[]); - for (tool, arguments) in [ - ("agc_apply_patch", json!({"patch":patch})), - ("agc_update_plan", plan.clone()), - ] { - audit.observe_item(&json!({"item":{ - "type":"mcpToolCall", "server":"agc_tools", "tool":tool, - "arguments":arguments, "result":{"content":[{"type":"text","text":"private-result-sentinel"}]} - }})); - } - audit.finish(true); - let records = read_turn_log(project.path(), "turn-01"); - let patch_item = records - .iter() - .find(|item| item["tool"] == "agc_apply_patch") - .unwrap(); - assert_eq!( - patch_item["arguments"], - json!({ - "patchHash":sha256_hex(patch.as_bytes()), "patchBytes":patch.len(), - "patchChars":patch.chars().count(), "patchOperations":1 - }) - ); - let plan_item = records - .iter() - .find(|item| item["tool"] == "agc_update_plan") - .unwrap(); - let plan_bytes = serde_json::to_vec(&plan).unwrap(); - assert_eq!( - plan_item["arguments"], - json!({ - "planHash":sha256_hex(&plan_bytes), "planBytes":plan_bytes.len(), "planSteps":1 - }) - ); - let persisted = serde_json::to_string(&records).unwrap(); - for sentinel in [ - "private-path-sentinel", - "sk-patch-body-sentinel", - "private-explanation-sentinel", - "sk-plan-step-sentinel", - "private-result-sentinel", - ] { - assert!(!persisted.contains(sentinel), "audit leaked {sentinel}"); - } - } - - #[test] - fn generate_image_prompt_is_truncated_with_hash() { - let project = fixture_project("audit-image"); - let prompt = "收".repeat(5000); - let mut audit = start_audit(project.path(), "x", &[]); - audit.observe_item(&json!({ - "item": { - "type": "mcpToolCall", - "tool": "agc_generate_image", - "arguments": { - "prompt": prompt, - "kind": "icon-spritesheet", - "sliceMode": "connected-components", - "sliceCount": 8, - "screenColor": "#CFEFFF" - } - } - })); - audit.finish(true); - let item = read_turn_log(project.path(), "turn-01") - .into_iter() - .find(|record| record["recordType"] == "direct.codex.item") - .expect("item"); - let stored = item["arguments"]["prompt"].as_str().expect("prompt"); - assert_eq!(stored.chars().count(), DIRECT_CODEX_AUDIT_BRIEF_CHARS); - assert_eq!(item["arguments"]["promptChars"], json!(5000)); - assert_eq!(item["arguments"]["sliceCount"], json!(8)); - assert_eq!(item["arguments"]["screenColor"], json!("#CFEFFF")); - assert_eq!( - item["arguments"]["promptSha256"], - json!(sha256_hex("收".repeat(5000).as_bytes())) - ); - } - - #[test] - fn file_change_keeps_path_and_kind_without_diff() { - let project = fixture_project("audit-patch"); - let mut audit = start_audit(project.path(), "x", &[]); - audit.observe_item(&json!({ - "item": { - "type": "fileChange", - "status": "completed", - "changes": [{ - "path": "game/index.html", - "kind": { "type": "add" }, - "diff": "*** SECRET PATCH" - }] - } - })); - audit.finish(true); - let item = read_turn_log(project.path(), "turn-01") - .into_iter() - .find(|record| record["recordType"] == "direct.codex.item") - .expect("item"); - assert_eq!(item["changes"][0]["path"], json!("game/index.html")); - assert_eq!(item["changes"][0]["kind"], json!("add")); - let dumped = serde_json::to_string(&item).expect("json"); - assert!(!dumped.contains("SECRET PATCH")); - assert!(!dumped.contains("diff")); - } - - #[test] - fn list_or_search_does_not_count_as_reading_offered_attachment() { - let project = fixture_project("audit-list"); - let root = project.path(); - let relative = "assets/uploads/upload-1-fast_gdd.md"; - fs::create_dir_all(root.join("assets/uploads")).expect("uploads"); - fs::write(root.join(relative), "x").expect("gdd"); - let mut audit = start_audit( - root, - "做游戏", - &[attachment_json( - "fast_gdd.md", - "text/markdown", - 1, - Some(relative), - Some("imported"), - )], - ); - audit.observe_item(&json!({ - "item": { - "type": "commandExecution", - "command": "rg fast_gdd assets", - "commandActions": [{ - "type": "search", - "query": "fast_gdd", - "path": "assets" - }] - } - })); - audit.observe_item(&json!({ - "item": { - "type": "commandExecution", - "command": "ls assets/uploads", - "commandActions": [{ "type": "listFiles", "path": "assets/uploads" }] - } - })); - audit.finish(true); - let end = read_turn_log(root, "turn-01") - .into_iter() - .find(|record| record["recordType"] == "direct.codex.turn_end") - .expect("end"); - assert_eq!(end["offeredRead"][0]["read"], json!(false)); - assert!(end["offeredRead"][0].get("contentSha256Match").is_none()); - assert!(end["firstDesign"].is_null()); - } - - #[test] - fn first_design_skips_reads_and_uses_later_art_item_seq() { - let project = fixture_project("audit-order"); - let root = project.path(); - let relative = "assets/uploads/upload-1-fast_gdd.md"; - fs::create_dir_all(root.join("assets/uploads")).expect("uploads"); - fs::write(root.join(relative), "x").expect("gdd"); - let mut audit = start_audit( - root, - "做游戏", - &[attachment_json( - "fast_gdd.md", - "text/markdown", - 1, - Some(relative), - Some("imported"), - )], - ); - audit.observe_item(&json!({ - "item": { - "type": "commandExecution", - "command": "type assets/uploads/upload-1-fast_gdd.md", - "commandActions": [{ "type": "read", "path": relative }] - } - })); - audit.observe_item(&json!({ - "item": { - "type": "mcpToolCall", - "tool": "taonier_prepare_game_art", - "arguments": { "brief": "收集冒险" } - } - })); - audit.finish(true); - let end = read_turn_log(root, "turn-01") - .into_iter() - .find(|record| record["recordType"] == "direct.codex.turn_end") - .expect("end"); - assert_eq!( - end["firstDesign"]["kind"], - json!("mcp:taonier_prepare_game_art") - ); - assert_eq!(end["firstDesign"]["seq"], json!(2)); - assert_eq!(end["offeredRead"][0]["read"], json!(true)); - } - - #[test] - fn item_cap_writes_truncated_marker() { - let project = fixture_project("audit-cap"); - let mut audit = start_audit(project.path(), "x", &[]); - for index in 0..(DIRECT_CODEX_AUDIT_MAX_ITEMS + 2) { - audit.observe_item(&json!({ - "item": { - "id": format!("item-{index}"), - "type": "commandExecution", - "command": "ls", - "commandActions": [{ "type": "unknown" }] - } - })); - } - audit.finish(true); - let records = read_turn_log(project.path(), "turn-01"); - let items = records - .iter() - .filter(|record| record["recordType"] == "direct.codex.item") - .count(); - assert_eq!(items, DIRECT_CODEX_AUDIT_MAX_ITEMS); - assert!(records - .iter() - .any(|record| record["recordType"] == "direct.codex.items_truncated")); - let end = records - .iter() - .find(|record| record["recordType"] == "direct.codex.turn_end") - .expect("end"); - assert_eq!(end["itemsTruncated"], json!(true)); - assert_eq!(end["itemCount"], json!(DIRECT_CODEX_AUDIT_MAX_ITEMS)); - } - - #[test] - fn agent_db_summary_points_at_relative_turn_log() { - let project = fixture_project("audit-db"); - let mut audit = start_audit(project.path(), "做游戏", &[]); - audit.observe_item(&json!({ - "item": { - "type": "mcpToolCall", - "tool": "taonier_prepare_game_art", - "arguments": { "brief": "俯视角收集冒险小游戏" } - } - })); - audit.finish(true); - let summary = read_agent_db(project.path()) - .into_iter() - .rev() - .find(|record| record["recordType"] == "direct.codex.turn") - .expect("summary"); - assert_eq!( - summary["turnLog"], - json!(format!("{DIRECT_CODEX_AUDIT_TURN_LOG_DIR}/turn-01.jsonl")) - ); - assert_eq!( - summary["firstDesign"]["briefPreview"], - json!("俯视角收集冒险小游戏") - ); - assert_eq!(summary["completed"], json!(true)); - assert!(!summary["turnLog"].as_str().expect("path").contains('\\')); - } - - #[test] - fn write_failure_does_not_panic_or_surface_through_finish() { - let project = fixture_project("audit-fail"); - fs::create_dir_all(project.path().join(".agent/runtime")).expect("runtime"); - fs::write( - project - .path() - .join(".agent/runtime/test-fail-direct-codex-audit"), - "1", - ) - .expect("marker"); - let mut audit = start_audit(project.path(), "x", &[]); - audit.observe_item(&json!({ - "item": { "type": "unknownTool" } - })); - audit.finish(true); - assert!(!project - .path() - .join(DIRECT_CODEX_AUDIT_TURN_LOG_DIR) - .join("turn-01.jsonl") - .is_file()); - } - - #[test] - fn unknown_item_type_keeps_only_public_fields() { - let project = fixture_project("audit-unknown"); - let mut audit = start_audit(project.path(), "x", &[]); - audit.observe_item(&json!({ - "item": { - "id": "mystery", - "type": "secretNewItem", - "payload": { "token": "leak-me" }, - "aggregatedOutput": "nope" - } - })); - audit.finish(true); - let item = read_turn_log(project.path(), "turn-01") - .into_iter() - .find(|record| record["recordType"] == "direct.codex.item") - .expect("item"); - assert_eq!(item["itemType"], json!("secretNewItem")); - assert_eq!(item["itemId"], json!("mystery")); - let dumped = serde_json::to_string(&item).expect("json"); - assert!(!dumped.contains("leak-me")); - assert!(!dumped.contains("payload")); - assert!(!dumped.contains("aggregatedOutput")); - } - - #[test] - fn agent_messages_are_skipped() { - let project = fixture_project("audit-skip"); - let mut audit = start_audit(project.path(), "x", &[]); - audit.observe_item(&json!({ - "item": { "type": "agentMessage", "text": "已按 GDD 完成" } - })); - audit.finish(true); - let items = read_turn_log(project.path(), "turn-01") - .into_iter() - .filter(|record| record["recordType"] == "direct.codex.item") - .count(); - assert_eq!(items, 0); - } -} diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_runtime/mod.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_runtime/mod.rs index a466faa8c..b33ddc4fc 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_runtime/mod.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_runtime/mod.rs @@ -4816,7 +4816,6 @@ pub(crate) async fn run_direct_game_creator_turn_at_with_creation_type( None, None, None, - None, ) .await } @@ -4826,7 +4825,6 @@ async fn run_direct_game_creator_turn_at_with_creation_type_and_emitter( prompt: &str, creation_type: Option<&str>, turn_emitter: Option<&DirectGameCreatorTurnUpdateEmitter>, - audit: Option<&mut DirectCodexTurnAudit>, direct_user_item: Option, capture: Option<( crate::analytics::contract::Context, @@ -4853,7 +4851,6 @@ async fn run_direct_game_creator_turn_at_with_creation_type_and_emitter( prompt, creation_type, turn_emitter, - audit, direct_user_item, capture, analytics_attempt_id, @@ -5055,7 +5052,6 @@ async fn run_direct_game_creator_turn_inner( prompt: &str, creation_type: Option<&str>, turn_emitter: Option<&DirectGameCreatorTurnUpdateEmitter>, - audit: Option<&mut DirectCodexTurnAudit>, direct_user_item: Option, capture: Option<( crate::analytics::contract::Context, @@ -5284,7 +5280,6 @@ async fn run_direct_game_creator_turn_inner( }; let mut feedback_prompt = prompt.to_string(); let mut turn_kind = DirectCodexTurnKind::User; - let mut audit = audit; let mut attempt = 1; let reply_result = loop { let result = direct_game_creator_codex_chat_at_with_optional_observer( @@ -5294,7 +5289,6 @@ async fn run_direct_game_creator_turn_inner( turn_kind, Some(&client_turn_id), Some(&mut observer), - audit.as_deref_mut(), Some(direct_user_item.clone()), ) .await; @@ -5343,7 +5337,6 @@ async fn run_direct_game_creator_turn_inner( } else { let mut feedback_prompt = prompt.to_string(); let mut turn_kind = DirectCodexTurnKind::User; - let mut audit = audit; let mut response = None; let mut attempt = 1; loop { @@ -5354,7 +5347,6 @@ async fn run_direct_game_creator_turn_inner( turn_kind, None, None, - audit.as_deref_mut(), Some(direct_user_item.clone()), ) .await; diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_runtime/user_input.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_runtime/user_input.rs index 6b3a85455..c29455565 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_runtime/user_input.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_runtime/user_input.rs @@ -61,8 +61,6 @@ pub(crate) async fn chat_with_game_creator_direct_codex( &user_prompt, creation_type.as_deref(), Some(&turn_emitter), - // DirectProject 的完整回合权威已经落在 project.jsonl;不再创建平行审计日志。 - None, canonical_user_item, capture, analytics_attempt_id.as_deref(), diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_turn_metrics.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_turn_metrics.rs deleted file mode 100644 index 45bdbdd17..000000000 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_turn_metrics.rs +++ /dev/null @@ -1,1355 +0,0 @@ -//! Direct 回合观察计时:只保存边界和安全元数据,不保存请求、响应或思考正文。 -//! 耗时使用单调钟,并发工具和请求按活动集合计算区间并集。 - -use serde_json::{json, Value}; -use std::collections::{BTreeMap, HashMap, HashSet, VecDeque}; -use std::path::PathBuf; -use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; -use std::sync::{Arc, Condvar, Mutex}; -use std::time::{Duration, Instant}; - -const MAX_TIMING_RECORDS: usize = 512; -const MAX_SSE_EVENT_BYTES: usize = 64 * 1024; -const WRITER_CAPACITY: usize = 1024; -const WRITER_BATCH_RECORDS: usize = 64; -const WRITER_BATCH_BYTES: usize = 256 * 1024; -const FLUSH_TIMEOUT: Duration = Duration::from_millis(1_500); - -#[derive(Clone, Copy, PartialEq)] -enum RecordPriority { - Detail, - TurnEnd, - LatestSummary, -} -struct PendingRecord { - sequence: u64, - line: String, - priority: RecordPriority, -} -#[derive(Default)] -struct WriterQueue { - records: VecDeque, - submitted: u64, - completed: u64, - closed: bool, -} -fn take_writer_batch(queue: &mut WriterQueue) -> Vec { - let mut batch = Vec::new(); - let mut bytes = 0_usize; - while let Some(next) = queue.records.front() { - if !batch.is_empty() - && (batch.len() >= WRITER_BATCH_RECORDS - || bytes.saturating_add(next.line.len() + 1) > WRITER_BATCH_BYTES) - { - break; - } - bytes = bytes.saturating_add(next.line.len() + 1); - if let Some(record) = queue.records.pop_front() { - batch.push(record); - } - } - batch -} -struct WriterShared { - queue: Mutex, - changed: Condvar, - failed: AtomicBool, - dropped: AtomicU64, -} -struct WriterOwner { - shared: Arc, -} -impl Drop for WriterOwner { - fn drop(&mut self) { - if let Ok(mut queue) = self.shared.queue.lock() { - queue.closed = true; - self.shared.changed.notify_all(); - #[cfg(test)] - { - let deadline = Instant::now() + Duration::from_secs(5); - while queue.completed < queue.submitted { - let remaining = deadline.saturating_duration_since(Instant::now()); - if remaining.is_zero() { - break; - } - queue = match self.shared.changed.wait_timeout(queue, remaining) { - Ok((next, _)) => next, - Err(_) => break, - }; - } - } - } - self.shared.changed.notify_all(); - } -} -#[derive(Clone)] -struct TimingWriter(Arc); -impl TimingWriter { - fn new(path: PathBuf) -> Self { - let shared = Arc::new(WriterShared { - queue: Mutex::new(WriterQueue::default()), - changed: Condvar::new(), - failed: AtomicBool::new(false), - dropped: AtomicU64::new(0), - }); - let worker = Arc::clone(&shared); - let spawned = std::thread::Builder::new() - .name("agc-turn-timing-writer".into()) - .spawn(move || { - loop { - let batch = { - let Ok(mut queue) = worker.queue.lock() else { - return; - }; - while queue.records.is_empty() && !queue.closed { - queue = match worker.changed.wait(queue) { - Ok(queue) => queue, - Err(_) => return, - }; - } - let batch = take_writer_batch(&mut queue); - if batch.is_empty() { - return; - } - batch - }; - // Disk locks/fsync happen only on this worker, with neither statistics - // nor queue mutex held. HTTP/body polling never waits for disk. - let lines: Vec<&str> = - batch.iter().map(|record| record.line.as_str()).collect(); - if crate::append_jsonl_lines(&path, &lines, "Direct 回合计时").is_err() { - worker.failed.store(true, Ordering::Relaxed); - } - if let Ok(mut queue) = worker.queue.lock() { - if let Some(last) = batch.last() { - queue.completed = last.sequence; - } - } - worker.changed.notify_all(); - } - }); - if spawned.is_err() { - shared.failed.store(true, Ordering::Relaxed); - if let Ok(mut queue) = shared.queue.lock() { - queue.closed = true; - } - } - Self(Arc::new(WriterOwner { shared })) - } - fn enqueue(&self, line: String, priority: RecordPriority) -> bool { - let shared = &self.0.shared; - let Ok(mut queue) = shared.queue.lock() else { - shared.failed.store(true, Ordering::Relaxed); - return false; - }; - if queue.closed { - shared.failed.store(true, Ordering::Relaxed); - return false; - } - // Only the newest post-terminal summary is needed. The turn_end itself is - // never evicted; a critical record evicts an ordinary detail if capacity fills. - if priority == RecordPriority::LatestSummary { - queue - .records - .retain(|record| record.priority != RecordPriority::LatestSummary); - } - if queue.records.len() >= WRITER_CAPACITY { - let evict = if priority != RecordPriority::Detail { - queue - .records - .iter() - .position(|record| record.priority == RecordPriority::Detail) - } else { - None - }; - if let Some(index) = evict { - queue.records.remove(index); - } else { - shared.dropped.fetch_add(1, Ordering::Relaxed); - shared.failed.store(true, Ordering::Relaxed); - return false; - } - shared.dropped.fetch_add(1, Ordering::Relaxed); - shared.failed.store(true, Ordering::Relaxed); - } - queue.submitted += 1; - let sequence = queue.submitted; - queue.records.push_back(PendingRecord { - sequence, - line, - priority, - }); - shared.changed.notify_one(); - true - } - fn flush(&self, timeout: Duration) -> bool { - let shared = &self.0.shared; - let Ok(mut queue) = shared.queue.lock() else { - return false; - }; - let target = queue.submitted; - let deadline = Instant::now() + timeout; - while queue.completed < target && !queue.closed { - let remaining = deadline.saturating_duration_since(Instant::now()); - if remaining.is_zero() { - shared.failed.store(true, Ordering::Relaxed); - return false; - } - match shared.changed.wait_timeout(queue, remaining) { - Ok((next, _)) => queue = next, - Err(_) => { - shared.failed.store(true, Ordering::Relaxed); - return false; - } - } - } - queue.completed >= target - } - fn failed(&self) -> bool { - self.0.shared.failed.load(Ordering::Relaxed) - } - fn dropped(&self) -> u64 { - self.0.shared.dropped.load(Ordering::Relaxed) - } -} - -pub(crate) fn direct_safe_model_identifier(value: &str) -> Option { - let lower = value.to_ascii_lowercase(); - if value.is_empty() - || value.len() > 96 - || !value - .bytes() - .all(|b| b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_' | b'.')) - || [ - "sk-", "pk-", "ghp_", "eyj", "bearer", "token", "secret", "password", - ] - .iter() - .any(|prefix| lower.starts_with(*prefix)) - || value.split(['-', '_', '.']).any(|part| part.len() >= 32) - { - return None; - } - Some(value.to_string()) -} - -#[derive(Clone, Copy)] -pub(crate) enum DirectMetricRoute { - MainSite, - ProviderProxy, - AppServerAuth, -} - -impl DirectMetricRoute { - fn name(self) -> &'static str { - match self { - Self::MainSite => "main-site", - Self::ProviderProxy => "provider-proxy", - Self::AppServerAuth => "app-server-auth", - } - } -} - -#[derive(Default)] -struct Coverage { - active: HashMap, - open_at: Option, - union_ms: u64, - completed_sum_ms: u64, - started_count: u64, - completed_count: u64, - closed_without_completion: u64, - incomplete_observed_sum_ms: u64, -} - -impl Coverage { - fn start(&mut self, id: &str, at: u64) { - if self.active.contains_key(id) { - return; - } - if self.active.is_empty() { - self.open_at = Some(at); - } - self.active.insert(id.to_string(), at); - self.started_count += 1; - } - fn end(&mut self, id: &str, at: u64) -> Option { - let started = self.active.remove(id)?; - self.completed_count += 1; - let duration = at.saturating_sub(started); - self.completed_sum_ms = self.completed_sum_ms.saturating_add(duration); - if self.active.is_empty() { - self.union_ms = self - .union_ms - .saturating_add(at.saturating_sub(self.open_at.take().unwrap_or(at))); - } - Some(duration) - } - fn snapshot(&self, at: u64) -> Value { - json!({ - "startedCount": self.started_count, "completedCount": self.completed_count, - "activeCount": self.active.len(), "completedSumMs": self.completed_sum_ms, - "closedWithoutCompletionCount": self.closed_without_completion, - "incompleteObservedSumMs": self.incomplete_observed_sum_ms, - "observedUnionMs": self.union_ms.saturating_add(self.open_at.map(|start| at.saturating_sub(start)).unwrap_or(0)), - }) - } - fn close_incomplete(&mut self, id: &str, at: u64) { - if let Some(duration) = self.end(id, at) { - self.completed_count = self.completed_count.saturating_sub(1); - self.completed_sum_ms = self.completed_sum_ms.saturating_sub(duration); - self.incomplete_observed_sum_ms = - self.incomplete_observed_sum_ms.saturating_add(duration); - self.closed_without_completion += 1; - } - } -} - -struct MetricsState { - origin: Instant, - started_at_ms: u64, - turn_end: Option<(u64, u64)>, - turn_end_coverage: Option, - coverage: BTreeMap<&'static str, Coverage>, - all: Coverage, - records: usize, - records_truncated: bool, - write_failed: bool, - attempts: u64, - unknown_item_starts: u64, - items: HashMap, - http_phases: BTreeMap<&'static str, PhaseAggregate>, -} - -#[derive(Default)] -struct PhaseAggregate { - observed_count: u64, - total_ms: u64, - max_ms: u64, -} -impl PhaseAggregate { - fn observe(&mut self, value: Option) { - if let Some(value) = value { - self.observed_count += 1; - self.total_ms = self.total_ms.saturating_add(value); - self.max_ms = self.max_ms.max(value); - } - } - fn snapshot(&self) -> Value { - json!({ - "observedCount": self.observed_count, - "totalMs": (self.observed_count > 0).then_some(self.total_ms), - "maxMs": (self.observed_count > 0).then_some(self.max_ms), - }) - } -} - -struct MetricsInner { - writer: TimingWriter, - client_turn_id: String, - state: Mutex, -} - -#[derive(Clone)] -pub(crate) struct DirectTurnMetrics(Arc); - -fn now_ms() -> u64 { - u64::try_from(crate::unix_millis()).unwrap_or(u64::MAX) -} - -impl DirectTurnMetrics { - pub(crate) fn new(log_path: PathBuf, client_turn_id: &str) -> Self { - Self(Arc::new(MetricsInner { - writer: TimingWriter::new(log_path), - client_turn_id: client_turn_id.to_string(), - state: Mutex::new(MetricsState { - origin: Instant::now(), - started_at_ms: now_ms(), - turn_end: None, - turn_end_coverage: None, - coverage: BTreeMap::new(), - all: Coverage::default(), - records: 0, - records_truncated: false, - write_failed: false, - attempts: 0, - unknown_item_starts: 0, - items: HashMap::new(), - http_phases: BTreeMap::new(), - }), - })) - } - fn elapsed(state: &MetricsState) -> u64 { - u64::try_from(state.origin.elapsed().as_millis()).unwrap_or(u64::MAX) - } - // Boundary writes only; never called for each text/body delta. - fn record_locked(&self, state: &mut MetricsState, mut value: Value) { - let summary = value["recordType"] == "direct.codex.timing_summary"; - if !summary && state.records >= MAX_TIMING_RECORDS { - state.records_truncated = true; - return; - } - if !summary { - state.records += 1; - } - value["clientTurnId"] = json!(self.0.client_turn_id); - value["recordedAtMs"] = json!(now_ms()); - if let Ok(line) = serde_json::to_string(&value) { - let priority = if summary { - RecordPriority::LatestSummary - } else { - RecordPriority::Detail - }; - if !self.0.writer.enqueue(line, priority) { - state.write_failed = true; - } - } - } - pub(crate) fn attempt( - &self, - configured_model: &str, - requested_model: &str, - effort: &str, - ) -> DirectMetricAttempt { - let attempt = DirectMetricAttempt(Arc::new(AttemptInner { - metrics: self.clone(), - id: uuid::Uuid::new_v4().to_string(), - finished: Mutex::new(false), - observed_events: Mutex::new(HashSet::new()), - })); - if let Ok(mut state) = self.0.state.lock() { - state.attempts += 1; - self.record_locked( - &mut state, - json!({ - "recordType": "direct.codex.attempt_started", "attemptId": attempt.id(), - "configuredModel": direct_safe_model_identifier(configured_model), - "requestedModel": direct_safe_model_identifier(requested_model), - "configuredReasoningEffort": safe_effort(effort), - }), - ); - } - attempt - } - fn summary_locked(&self, state: &MetricsState) -> Value { - let at = Self::elapsed(state); - let total = state.turn_end.map(|(elapsed, _)| elapsed).unwrap_or(at); - let observed = state.all.snapshot(at); - let all_finished = state.all.active.is_empty(); - let within_turn = state.turn_end_coverage.clone().unwrap_or_else(|| json!({ - "all": observed.clone(), - "categories": state.coverage.iter().map(|(k,v)| ((*k).to_string(), v.snapshot(at))).collect::>(), - })); - let union = within_turn["all"]["observedUnionMs"].as_u64().unwrap_or(0); - json!({ - "schemaVersion": "agc-direct-timing.v1", - "startedAtMs": state.started_at_ms, - "endedAtMs": state.turn_end.map(|(_, wall)| wall), "wallDurationMs": total, - "attemptCount": state.attempts, - "httpPhases": state.http_phases.iter().map(|(key,value)| ((*key).to_string(), value.snapshot())).collect::>(), - "httpPhaseTotalsOverlap": true, - "categories": state.coverage.iter().map(|(k,v)| ((*k).to_string(), v.snapshot(at))).collect::>(), - "observed": observed, - "withinTurn": within_turn, - // A late body can outlive the logical turn. Do not call its residual inference. - "unattributedMs": if all_finished && union <= total { Some(total - union) } else { None }, - "complete": all_finished && state.turn_end.is_some() && state.unknown_item_starts == 0 && state.all.closed_without_completion == 0, - "timingSource": "host-monotonic-observation", - "unknownItemStartCount": state.unknown_item_starts, - "detailsTruncated": state.records_truncated || self.0.writer.dropped() > 0, - "writeFailed": state.write_failed || self.0.writer.failed(), - "writerDroppedRecords": self.0.writer.dropped(), - "upstreamQueueMs": Value::Null, "modelInferenceMs": Value::Null, - }) - } - fn mark_finished(state: &mut MetricsState) { - let elapsed = Self::elapsed(state); - if state.turn_end.is_none() { - state.turn_end_coverage = Some(json!({ - "all": state.all.snapshot(elapsed), - "categories": state.coverage.iter().map(|(k,v)| ((*k).to_string(), v.snapshot(elapsed))).collect::>(), - })); - } - state.turn_end.get_or_insert((elapsed, now_ms())); - } - #[cfg(test)] - pub(crate) fn finish(&self) -> Value { - let Ok(mut state) = self.0.state.lock() else { - return json!({"available": false}); - }; - Self::mark_finished(&mut state); - self.summary_locked(&state) - } - pub(crate) fn append_audit_record(&self, mut record: Value) -> bool { - let Ok(mut state) = self.0.state.lock() else { - return false; - }; - let terminal = record["recordType"] == "direct.codex.turn_end"; - if terminal { - Self::mark_finished(&mut state); - record["timing"] = self.summary_locked(&state); - } - let Ok(line) = serde_json::to_string(&record) else { - return false; - }; - self.0.writer.enqueue( - line, - if terminal { - RecordPriority::TurnEnd - } else { - RecordPriority::Detail - }, - ) - } - pub(crate) async fn flush(&self) { - let writer = self.0.writer.clone(); - let result = tokio::time::timeout( - FLUSH_TIMEOUT + Duration::from_millis(250), - tokio::task::spawn_blocking(move || writer.flush(FLUSH_TIMEOUT)), - ) - .await; - if !matches!(result, Ok(Ok(true))) || self.0.writer.failed() { - self.0.writer.0.shared.failed.store(true, Ordering::Relaxed); - if let Ok(mut state) = self.0.state.lock() { - state.write_failed = true; - let summary = self.summary_locked(&state); - self.record_locked( - &mut state, - json!({"recordType":"direct.codex.timing_summary", "timing":summary}), - ); - } - } - } - #[cfg(test)] - pub(crate) fn flush_for_test(&self) -> bool { - self.0.writer.flush(Duration::from_secs(30)) && !self.0.writer.failed() - } - #[cfg(test)] - pub(crate) async fn wait_for_test_writes(&self) -> bool { - let metrics = self.clone(); - tokio::task::spawn_blocking(move || metrics.flush_for_test()) - .await - .unwrap_or(false) - } - #[cfg(test)] - pub(crate) fn snapshot(&self) -> Value { - self.summary_locked(&self.0.state.lock().unwrap()) - } -} - -fn safe_effort(value: &str) -> Option<&str> { - matches!( - value, - "none" | "minimal" | "low" | "medium" | "high" | "xhigh" | "max" | "ultra" - ) - .then_some(value) -} - -struct AttemptInner { - metrics: DirectTurnMetrics, - id: String, - finished: Mutex, - observed_events: Mutex>, -} -impl Drop for AttemptInner { - fn drop(&mut self) { - self.finish("dropped"); - } -} -impl AttemptInner { - fn finish(&self, status: &'static str) { - let Ok(mut finished) = self.finished.lock() else { - return; - }; - if *finished { - return; - } - *finished = true; - if let Ok(mut state) = self.metrics.0.state.lock() { - let at = DirectTurnMetrics::elapsed(&state); - let prefix = format!("{}:", self.id); - let ids: Vec<_> = state - .items - .keys() - .filter(|id| id.starts_with(&prefix)) - .cloned() - .collect(); - for id in ids { - if let Some(kind) = state.items.remove(&id) { - state - .coverage - .entry(kind) - .or_default() - .close_incomplete(&id, at); - state.all.close_incomplete(&id, at); - } - } - self.metrics.record_locked(&mut state, json!({ - "recordType": "direct.codex.attempt_finished", "attemptId": self.id, "status": status, - })); - if state.turn_end.is_some() { - let summary = self.metrics.summary_locked(&state); - self.metrics.record_locked( - &mut state, - json!({ - "recordType": "direct.codex.timing_summary", "timing": summary, - }), - ); - } - } - } -} - -#[derive(Clone)] -pub(crate) struct DirectMetricAttempt(Arc); -impl DirectMetricAttempt { - pub(crate) fn id(&self) -> &str { - &self.0.id - } - pub(crate) fn finish(&self, status: &'static str) { - self.0.finish(status); - } - pub(crate) fn observe_app_event(&self, event: &'static str) { - let Ok(mut events) = self.0.observed_events.lock() else { - return; - }; - if !events.insert(event) { - return; - } - if let Ok(mut state) = self.0.metrics.0.state.lock() { - let offset = DirectTurnMetrics::elapsed(&state); - self.0.metrics.record_locked( - &mut state, - json!({ - "recordType": "direct.codex.app_server_observation", - "attemptId": self.id(), "event": event, "offsetMsFromTurn": offset, - }), - ); - } - } - pub(crate) fn route(&self, route: DirectMetricRoute) { - if let Ok(mut state) = self.0.metrics.0.state.lock() { - self.0.metrics.record_locked(&mut state, json!({ - "recordType": "direct.codex.attempt_route", "attemptId": self.id(), - "route": route.name(), "httpTelemetryAvailable": !matches!(route, DirectMetricRoute::AppServerAuth), - })); - } - } - pub(crate) fn span(&self, kind: &'static str) -> DirectMetricSpan { - let id = uuid::Uuid::new_v4().to_string(); - let mut start = 0; - if let Ok(mut state) = self.0.metrics.0.state.lock() { - start = DirectTurnMetrics::elapsed(&state); - state.coverage.entry(kind).or_default().start(&id, start); - state.all.start(&id, start); - } - DirectMetricSpan { - attempt: self.clone(), - id, - kind, - start, - started_at_ms: now_ms(), - ended: false, - } - } - pub(crate) fn observe_raw_item(&self, item: &Value) { - let completed = match item.get("type").and_then(Value::as_str) { - Some("function_call") => false, - Some("function_call_output") => true, - _ => return, - }; - if let Some(id) = item.get("call_id").and_then(Value::as_str) { - self.item_boundary(id, "tool-pending", completed); - if !completed { - self.item_boundary(id, "tool-dispatch-wait", false); - } else { - self.end_dispatch_wait_if_present(id, false); - } - } - } - pub(crate) fn observe_item(&self, item: &Value, completed: bool) { - let kind = match item.get("type").and_then(Value::as_str) { - Some("contextCompaction") => "context-compaction", - Some("commandExecution" | "mcpToolCall" | "fileChange" | "imageView" | "webSearch") => { - "tool-execution" - } - _ => return, - }; - if let Some(id) = item.get("id").and_then(Value::as_str) { - if !completed && kind == "tool-execution" { - self.end_dispatch_wait_if_present(id, true); - } - self.item_boundary(id, kind, completed); - } - } - fn end_dispatch_wait_if_present(&self, item_id: &str, execution_observed: bool) { - let id = format!("{}:tool-dispatch-wait:{item_id}", self.id()); - let Ok(mut state) = self.0.metrics.0.state.lock() else { - return; - }; - if state.items.remove(&id).is_some() { - let at = DirectTurnMetrics::elapsed(&state); - if execution_observed { - state - .coverage - .entry("tool-dispatch-wait") - .or_default() - .end(&id, at); - state.all.end(&id, at); - } else { - state - .coverage - .entry("tool-dispatch-wait") - .or_default() - .close_incomplete(&id, at); - state.all.close_incomplete(&id, at); - } - } - } - fn item_boundary(&self, item_id: &str, kind: &'static str, completed: bool) { - // Correlation IDs remain in memory; they never become paths or logged text. - let id = format!("{}:{kind}:{item_id}", self.id()); - let Ok(mut state) = self.0.metrics.0.state.lock() else { - return; - }; - let at = DirectTurnMetrics::elapsed(&state); - if completed { - if state.items.remove(&id).is_some() { - state.coverage.entry(kind).or_default().end(&id, at); - state.all.end(&id, at); - } else { - state.unknown_item_starts += 1; - } - } else { - state.items.insert(id.clone(), kind); - state.coverage.entry(kind).or_default().start(&id, at); - state.all.start(&id, at); - } - } -} - -pub(crate) struct DirectMetricSpan { - attempt: DirectMetricAttempt, - id: String, - kind: &'static str, - start: u64, - started_at_ms: u64, - ended: bool, -} -impl DirectMetricSpan { - pub(crate) fn finish(&mut self, status: &'static str) { - self.finish_record("direct.codex.timing_span", json!({"status": status})); - } - fn finish_record(&mut self, record_type: &'static str, mut record: Value) { - if self.ended { - return; - } - self.ended = true; - let metrics = &self.attempt.0.metrics; - let Ok(mut state) = metrics.0.state.lock() else { - return; - }; - let at = DirectTurnMetrics::elapsed(&state); - state - .coverage - .entry(self.kind) - .or_default() - .end(&self.id, at); - state.all.end(&self.id, at); - record["recordType"] = json!(record_type); - record["attemptId"] = json!(self.attempt.id()); - record["requestId"] = json!(self.id); - record["category"] = json!(self.kind); - record["startedAtMs"] = json!(self.started_at_ms); - record["endedAtMs"] = json!(now_ms()); - record["durationMs"] = json!(at.saturating_sub(self.start)); - metrics.record_locked(&mut state, record); - if state.turn_end.is_some() { - let summary = metrics.summary_locked(&state); - metrics.record_locked( - &mut state, - json!({"recordType": "direct.codex.timing_summary", "timing": summary}), - ); - } - } -} -impl Drop for DirectMetricSpan { - fn drop(&mut self) { - self.finish("dropped"); - } -} - -/// Retains one bounded SSE event, and never writes its text. -#[derive(Default)] -struct SseObservation { - line: Vec, - data: Vec, - skip_event: bool, - first_event_ms: Option, - first_content_ms: Option, - reported_model: Option, - terminal: Option<&'static str>, -} -impl SseObservation { - fn chunk(&mut self, chunk: &[u8], at: u64) { - for &byte in chunk { - if byte != b'\n' { - if self.line.len() < MAX_SSE_EVENT_BYTES { - self.line.push(byte); - } else { - self.skip_event = true; - } - continue; - } - if self.line.last() == Some(&b'\r') { - self.line.pop(); - } - if self.line.is_empty() { - if !self.skip_event && !self.data.is_empty() { - self.event(at); - } - self.data.clear(); - self.skip_event = false; - } else if !self.skip_event { - if let Some(data) = self.line.strip_prefix(b"data:") { - let data = data.strip_prefix(b" ").unwrap_or(data); - if self.data.len() + data.len() + 1 <= MAX_SSE_EVENT_BYTES { - if !self.data.is_empty() { - self.data.push(b'\n'); - } - self.data.extend_from_slice(data); - } else { - self.skip_event = true; - self.data.clear(); - } - } - } - self.line.clear(); - } - } - fn event(&mut self, at: u64) { - if self.data == b"[DONE]" { - self.first_event_ms.get_or_insert(at); - return; - } - let Ok(value) = serde_json::from_slice::(&self.data) else { - return; - }; - self.first_event_ms.get_or_insert(at); - if self.reported_model.is_none() { - self.reported_model = value - .pointer("/response/model") - .and_then(Value::as_str) - .and_then(direct_safe_model_identifier); - } - match value.get("type").and_then(Value::as_str) { - Some( - "response.output_text.delta" - | "response.refusal.delta" - | "response.function_call_arguments.delta", - ) => { - if value - .get("delta") - .and_then(Value::as_str) - .is_some_and(|delta| !delta.is_empty()) - { - self.first_content_ms.get_or_insert(at); - } - } - Some("response.completed") => self.terminal = Some("completed"), - Some("response.failed" | "error") => self.terminal = Some("failed"), - Some("response.incomplete") => self.terminal = Some("incomplete"), - _ => {} - } - } -} - -pub(crate) struct DirectRequestTiming { - span: DirectMetricSpan, - origin: Instant, - dispatched_ms: Option, - headers_ms: Option, - first_chunk_ms: Option, - status: Option, - requested_model: Option, - reasoning_effort: Option, - sse: bool, - parser: SseObservation, - bytes: u64, -} -impl DirectRequestTiming { - pub(crate) fn new(attempt: DirectMetricAttempt) -> Self { - let origin = Instant::now(); - let span = attempt.span("http-request"); - if let Ok(mut state) = attempt.0.metrics.0.state.lock() { - attempt.0.metrics.record_locked( - &mut state, - json!({ - "recordType": "direct.codex.request_started", - "attemptId": attempt.id(), "requestId": span.id, - "startedAtMs": span.started_at_ms, - }), - ); - } - Self { - span, - origin, - dispatched_ms: None, - headers_ms: None, - first_chunk_ms: None, - status: None, - requested_model: None, - reasoning_effort: None, - sse: false, - parser: SseObservation::default(), - bytes: 0, - } - } - fn elapsed(&self) -> u64 { - u64::try_from(self.origin.elapsed().as_millis()).unwrap_or(u64::MAX) - } - pub(crate) fn request_body(&mut self, body: &[u8]) { - #[derive(serde::Deserialize)] - struct Metadata { - model: Option, - reasoning: Option, - } - #[derive(serde::Deserialize)] - struct Reasoning { - effort: Option, - } - if let Ok(metadata) = serde_json::from_slice::(body) { - self.requested_model = metadata - .model - .as_deref() - .and_then(direct_safe_model_identifier); - self.reasoning_effort = metadata - .reasoning - .and_then(|r| r.effort) - .and_then(|v| safe_effort(&v).map(str::to_string)); - } - } - pub(crate) fn dispatched(&mut self) { - self.dispatched_ms = Some(self.elapsed()); - } - pub(crate) fn headers(&mut self, status: u16, sse: bool) { - self.headers_ms = Some(self.elapsed()); - self.status = Some(status); - self.sse = sse; - } - pub(crate) fn chunk(&mut self, bytes: &[u8]) { - if bytes.is_empty() { - return; - } - let elapsed = self.elapsed(); - self.first_chunk_ms.get_or_insert(elapsed); - self.bytes = self.bytes.saturating_add(bytes.len() as u64); - if self.sse { - self.parser.chunk(bytes, elapsed); - } - } - pub(crate) fn finish(&mut self, transport_status: &'static str) { - if self.span.ended { - return; - } - let elapsed = self.elapsed(); - let phases = [ - ("requestDurationMs", Some(elapsed)), - ("upstreamDispatchOffsetMs", self.dispatched_ms), - ( - "dispatchToHeadersMs", - self.headers_ms - .zip(self.dispatched_ms) - .map(|(end, start)| end.saturating_sub(start)), - ), - ("firstBodyChunkOffsetMs", self.first_chunk_ms), - ("firstSseEventOffsetMs", self.parser.first_event_ms), - ("firstContentDeltaOffsetMs", self.parser.first_content_ms), - ( - "streamDurationMs", - self.headers_ms.map(|start| elapsed.saturating_sub(start)), - ), - ]; - // Aggregate before the detail-record cap. Long turns retain every observed - // phase's count/sum even after request_timing details have been truncated. - if let Ok(mut state) = self.span.attempt.0.metrics.0.state.lock() { - for (name, value) in phases { - state.http_phases.entry(name).or_default().observe(value); - } - } - self.span.finish_record("direct.codex.request_timing", json!({ - "transportStatus": transport_status, "responseStatus": self.parser.terminal, - "httpStatus": self.status, "requestedModel": self.requested_model, - "reasoningEffort": self.reasoning_effort, "responseReportedModel": self.parser.reported_model, - "upstreamDispatchOffsetMs": self.dispatched_ms, "responseHeadersOffsetMs": self.headers_ms, - "dispatchToHeadersMs": self.headers_ms.zip(self.dispatched_ms).map(|(end,start)| end.saturating_sub(start)), - "firstBodyChunkOffsetMs": self.first_chunk_ms, "firstSseEventOffsetMs": self.parser.first_event_ms, - "firstContentDeltaOffsetMs": self.parser.first_content_ms, - "streamDurationMs": self.headers_ms.map(|start| self.elapsed().saturating_sub(start)), - "responseBytes": self.bytes, - })); - } -} -impl Drop for DirectRequestTiming { - fn drop(&mut self) { - self.finish("dropped"); - } -} - -#[cfg(test)] -mod tests { - use super::*; - fn timing_log_path(root: &std::path::Path) -> PathBuf { - root.join(".agent/runtime/direct-codex/turns/turn.jsonl") - } - #[test] - fn writer_batches_respect_record_and_byte_limits_without_reordering() { - let mut queue = WriterQueue::default(); - for sequence in 1..=130 { - queue.records.push_back(PendingRecord { - sequence, - line: "{}".into(), - priority: RecordPriority::Detail, - }); - } - let first = take_writer_batch(&mut queue); - let second = take_writer_batch(&mut queue); - let third = take_writer_batch(&mut queue); - assert_eq!((first.len(), second.len(), third.len()), (64, 64, 2)); - assert_eq!( - first - .into_iter() - .chain(second) - .chain(third) - .map(|record| record.sequence) - .collect::>(), - (1..=130).collect::>() - ); - for sequence in 1..=3 { - queue.records.push_back(PendingRecord { - sequence, - line: "x".repeat(WRITER_BATCH_BYTES / 2), - priority: RecordPriority::Detail, - }); - } - // The record terminator also counts toward the byte bound. - assert_eq!(take_writer_batch(&mut queue).len(), 1); - assert_eq!(queue.records.front().unwrap().sequence, 2); - } - - #[test] - fn batch_append_preserves_record_bytes_and_existing_tail_repair() { - use std::io::Write; - let root = tempfile::tempdir().unwrap(); - let path = timing_log_path(root.path()); - crate::append_jsonl_line(&path, r#"{"id":0}"#, "批量计时测试").unwrap(); - std::fs::OpenOptions::new() - .append(true) - .open(&path) - .unwrap() - .write_all(br#"{"interrupted":"#) - .unwrap(); - let records = [ - r#"{"id":1,"text":"中文\n第二行","argsHash":"0123456789abcdef"}"#, - r#"{"id":2,"value":"quoted \"value\""}"#, - ]; - crate::append_jsonl_lines(&path, &records, "批量计时测试").unwrap(); - let expected = format!("{{\"id\":0}}\n{}\n{}\n", records[0], records[1]); - assert_eq!(std::fs::read(&path).unwrap(), expected.as_bytes()); - assert!(crate::append_jsonl_lines(&path, &["{}\n{}"], "批量计时测试").is_err()); - assert_eq!(std::fs::read(&path).unwrap(), expected.as_bytes()); - } - #[test] - fn blocked_writer_queue_is_bounded_and_keeps_terminal_order() { - // No worker consumes this queue: model a disk operation that is still blocked. - let shared = Arc::new(WriterShared { - queue: Mutex::new(WriterQueue::default()), - changed: Condvar::new(), - failed: AtomicBool::new(false), - dropped: AtomicU64::new(0), - }); - let writer = TimingWriter(Arc::new(WriterOwner { - shared: Arc::clone(&shared), - })); - for _ in 0..WRITER_CAPACITY { - assert!(writer.enqueue("{}".into(), RecordPriority::Detail)); - } - assert!(!writer.enqueue("{}".into(), RecordPriority::Detail)); - assert!(writer.enqueue("turn-end".into(), RecordPriority::TurnEnd)); - assert!(writer.enqueue("summary-1".into(), RecordPriority::LatestSummary)); - assert!(writer.enqueue("summary-2".into(), RecordPriority::LatestSummary)); - assert!(!writer.flush(Duration::ZERO)); - assert!(writer.failed()); - assert!(writer.dropped() >= 3); - let mut queue = shared.queue.lock().unwrap(); - assert!(queue.records.len() <= WRITER_CAPACITY); - let end = queue - .records - .iter() - .position(|record| record.line == "turn-end") - .unwrap(); - let summary = queue - .records - .iter() - .position(|record| record.line == "summary-2") - .unwrap(); - assert!(summary > end); - assert!(!queue - .records - .iter() - .any(|record| record.line == "summary-1")); - assert!(queue - .records - .iter() - .zip(queue.records.iter().skip(1)) - .all(|(a, b)| a.sequence < b.sequence)); - queue.records.clear(); - queue.completed = queue.submitted; - queue.closed = true; - } - #[test] - fn coverage_tracks_parallel_out_of_order_completion_without_double_counting() { - let mut coverage = Coverage::default(); - coverage.start("a", 10); - coverage.start("b", 20); - coverage.start("c", 30); - coverage.start("b", 40); - coverage.end("b", 70); - coverage.end("a", 90); - coverage.end("c", 100); - coverage.start("d", 150); - coverage.end("d", 170); - let result = coverage.snapshot(200); - assert_eq!(result["observedUnionMs"], 110); - assert_eq!(result["completedSumMs"], 220); - assert_eq!(result["startedCount"], 4); - assert_eq!(result["activeCount"], 0); - } - #[test] - fn safe_metadata_rejects_urls_keys_and_token_shaped_identifiers() { - assert_eq!( - direct_safe_model_identifier("gpt-5.6-sol"), - Some("gpt-5.6-sol".into()) - ); - for value in [ - "https://upstream/model?key=secret", - "sk-proj-secret", - "Bearer abc", - "eyJhbGciOiJIUzI1NiJ9.abc.def", - "abcdefghijklmnopqrstuvwxyz0123456789", - ] { - assert!(direct_safe_model_identifier(value).is_none(), "{value}"); - } - } - #[test] - fn fragmented_sse_separates_first_event_content_and_reported_model() { - let mut parser = SseObservation::default(); - parser.chunk(b": ping\n\ndata: {\"type\":\"response.created\",\"response\":{\"model\":\"gpt-5.6-sol\"}}\r\n\r\n", 10); - parser.chunk( - b"data: {\"type\":\"response.output_text.delta\",\"delta\":\"sec", - 30, - ); - assert_eq!(parser.first_event_ms, Some(10)); - assert_eq!(parser.first_content_ms, None); - parser.chunk( - b"ret text\"}\n\ndata: {\"type\":\"response.completed\"}\n\n", - 40, - ); - assert_eq!(parser.first_content_ms, Some(40)); - assert_eq!(parser.reported_model.as_deref(), Some("gpt-5.6-sol")); - assert_eq!(parser.terminal, Some("completed")); - assert!(parser.data.is_empty()); - } - #[test] - fn oversized_sse_recovers_at_next_event_and_reasoning_is_not_output() { - let mut parser = SseObservation::default(); - parser.chunk(b"data: ", 1); - parser.chunk(&vec![b'x'; MAX_SSE_EVENT_BYTES + 10], 2); - parser.chunk( - b"\n\ndata: {\"type\":\"response.reasoning_text.delta\",\"delta\":\"private\"}\n\n", - 3, - ); - assert_eq!(parser.first_content_ms, None); - parser.chunk(b"data: {\"type\":\"response.failed\"}\n\n", 4); - assert_eq!(parser.terminal, Some("failed")); - assert!(parser.line.capacity() <= MAX_SSE_EVENT_BYTES * 2); - } - #[test] - fn summary_continues_after_detail_cap_and_missing_boundaries_stay_unknown() { - let root = tempfile::tempdir().unwrap(); - let metrics = DirectTurnMetrics::new(timing_log_path(root.path()), "turn-1"); - metrics.0.state.lock().unwrap().records = MAX_TIMING_RECORDS; - let attempt = metrics.attempt("gpt-5.6-sol", "gpt-5.6-sol", "high"); - for i in 0..600 { - let item = json!({"id": format!("tool-{i}"), "type": "mcpToolCall"}); - attempt.observe_item(&item, false); - attempt.observe_item(&item, true); - } - attempt.observe_item(&json!({"id":"missing", "type":"mcpToolCall"}), true); - attempt.finish("completed"); - let summary = metrics.finish(); - assert_eq!( - summary["categories"]["tool-execution"]["completedCount"], - 600 - ); - assert_eq!(summary["unknownItemStartCount"], 1); - assert_eq!(summary["detailsTruncated"], true); - assert!(summary["modelInferenceMs"].is_null()); - assert!(summary["categories"].get("http-request").is_none()); - } - #[test] - fn dropped_request_keeps_original_attempt_and_never_persists_content() { - let root = tempfile::tempdir().unwrap(); - let path = timing_log_path(root.path()); - let metrics = DirectTurnMetrics::new(path.clone(), "turn-1"); - let first = metrics.attempt("gpt-5.6-sol", "gpt-5.6-sol", "high"); - let mut request = DirectRequestTiming::new(first.clone()); - request.request_body(br#"{"model":"sk-secret","input":"TOP SECRET"}"#); - request.headers(200, true); - request - .chunk(b"data: {\"type\":\"response.output_text.delta\",\"delta\":\"TOP SECRET\"}\n\n"); - let second = metrics.attempt("gpt-5.6-sol", "gpt-5.6-sol", "high"); - drop(request); - assert!( - metrics.flush_for_test(), - "writer failed: {}", - metrics.snapshot() - ); - let text = std::fs::read_to_string(path).unwrap(); - assert!(!text.contains("TOP SECRET")); - assert!(!text.contains("sk-secret")); - let record: Value = text - .lines() - .map(|s| serde_json::from_str::(s).unwrap()) - .find(|v| v["recordType"] == "direct.codex.request_timing") - .unwrap(); - assert_eq!(record["attemptId"], first.id()); - assert_ne!(record["attemptId"], second.id()); - assert_eq!(record["transportStatus"], "dropped"); - assert!(record["upstreamDispatchOffsetMs"].is_null()); - } - - #[test] - fn raw_calls_track_dispatch_wait_and_unmatched_execution_is_not_fabricated() { - let root = tempfile::tempdir().unwrap(); - let metrics = DirectTurnMetrics::new(timing_log_path(root.path()), "turn-queue"); - let attempt = metrics.attempt("gpt-5.6-sol", "gpt-5.6-sol", "high"); - attempt.observe_raw_item(&json!({"type":"function_call", "call_id":"a"})); - attempt.observe_raw_item(&json!({"type":"function_call", "call_id":"b"})); - attempt.observe_item(&json!({"type":"mcpToolCall", "id":"a"}), false); - attempt.observe_item(&json!({"type":"mcpToolCall", "id":"a"}), true); - attempt.observe_raw_item(&json!({"type":"function_call_output", "call_id":"a"})); - attempt.observe_raw_item(&json!({"type":"function_call_output", "call_id":"b"})); - attempt.finish("completed"); - let summary = metrics.finish(); - assert_eq!(summary["categories"]["tool-pending"]["completedCount"], 2); - assert_eq!( - summary["categories"]["tool-dispatch-wait"]["completedCount"], - 1 - ); - assert_eq!( - summary["categories"]["tool-dispatch-wait"]["closedWithoutCompletionCount"], - 1 - ); - assert_eq!(summary["categories"]["tool-execution"]["completedCount"], 1); - assert_eq!(summary["observed"]["activeCount"], 0); - assert_eq!(summary["complete"], false); - } - - #[test] - fn interrupted_items_and_failed_persistence_remain_diagnostic_only() { - let root = tempfile::tempdir().unwrap(); - let blocked = root.path().join("not-a-directory"); - std::fs::write(&blocked, "fixture").unwrap(); - let metrics = DirectTurnMetrics::new(timing_log_path(&blocked), "turn-interrupted"); - let attempt = metrics.attempt("gpt-5.6-sol", "gpt-5.6-sol", "high"); - attempt.observe_item(&json!({"type":"contextCompaction", "id":"compact"}), false); - attempt.observe_item(&json!({"type":"mcpToolCall", "id":"unfinished"}), false); - attempt.finish("interrupted"); - metrics.flush_for_test(); - let summary = metrics.finish(); - assert_eq!(summary["writeFailed"], true); - assert_eq!( - summary["categories"]["context-compaction"]["completedCount"], - 0 - ); - assert_eq!( - summary["categories"]["context-compaction"]["closedWithoutCompletionCount"], - 1 - ); - assert_eq!(summary["categories"]["tool-execution"]["completedCount"], 0); - assert_eq!(summary["observed"]["activeCount"], 0); - assert_eq!(summary["complete"], false); - } - - #[test] - fn http_phase_aggregates_survive_more_than_the_detail_cap() { - let root = tempfile::tempdir().unwrap(); - let metrics = DirectTurnMetrics::new(timing_log_path(root.path()), "turn-many-requests"); - let attempt = metrics.attempt("gpt-5.6-sol", "gpt-5.6-sol", "high"); - metrics.0.state.lock().unwrap().records = MAX_TIMING_RECORDS; - for _ in 0..600 { - let mut request = DirectRequestTiming::new(attempt.clone()); - request.dispatched(); - request.headers(200, true); - request.chunk(b"data: {\"type\":\"response.created\"}\n\n"); - request.finish("eof"); - } - attempt.finish("completed"); - let summary = metrics.finish(); - assert_eq!(summary["categories"]["http-request"]["completedCount"], 600); - for phase in [ - "requestDurationMs", - "dispatchToHeadersMs", - "firstBodyChunkOffsetMs", - "firstSseEventOffsetMs", - "streamDurationMs", - ] { - assert_eq!(summary["httpPhases"][phase]["observedCount"], 600); - assert!(summary["httpPhases"][phase]["totalMs"].is_number()); - } - assert_eq!( - summary["httpPhases"]["firstContentDeltaOffsetMs"]["observedCount"], - 0 - ); - assert!(summary["httpPhases"]["firstContentDeltaOffsetMs"]["totalMs"].is_null()); - assert_eq!(summary["detailsTruncated"], true); - } - - #[test] - fn audit_terminal_precedes_late_body_summary_in_single_writer_order() { - let root = tempfile::tempdir().unwrap(); - let path = timing_log_path(root.path()); - let metrics = DirectTurnMetrics::new(path.clone(), "turn-order"); - metrics.append_audit_record(json!({"recordType":"direct.codex.turn_start"})); - let attempt = metrics.attempt("gpt-5.6-sol", "gpt-5.6-sol", "high"); - let mut request = DirectRequestTiming::new(attempt.clone()); - metrics.append_audit_record(json!({"recordType":"direct.codex.turn_end"})); - request.finish("eof"); - attempt.finish("completed"); - assert!( - metrics.flush_for_test(), - "writer failed: {}", - metrics.snapshot() - ); - let records: Vec = std::fs::read_to_string(path) - .unwrap() - .lines() - .map(|line| serde_json::from_str(line).unwrap()) - .collect(); - let end = records - .iter() - .position(|record| record["recordType"] == "direct.codex.turn_end") - .unwrap(); - let summary = records - .iter() - .rposition(|record| record["recordType"] == "direct.codex.timing_summary") - .unwrap(); - assert!(summary > end); - assert_eq!(records[end]["timing"]["complete"], false); - assert_eq!(records[summary]["timing"]["complete"], true); - assert_eq!( - records[end]["timing"]["withinTurn"], - records[summary]["timing"]["withinTurn"] - ); - } -} diff --git a/apps/ai-game-creator-shell/src-tauri/src/project/agent_db.rs b/apps/ai-game-creator-shell/src-tauri/src/project/agent_db.rs index aaa98aeae..bcd6d99d4 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/project/agent_db.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/project/agent_db.rs @@ -4780,27 +4780,6 @@ pub(crate) fn append_jsonl_line(path: &Path, line: &str, error_label: &str) -> R append_jsonl_line_unlocked(path, line, error_label) } -/// Append already serialized JSON lines under one existing append lock and fsync. -/// Keep each record byte-for-byte intact; physical newlines belong to this framing layer. -pub(crate) fn append_jsonl_lines( - path: &Path, - lines: &[&str], - error_label: &str, -) -> Result<(), String> { - if lines.is_empty() { - return Ok(()); - } - if lines - .iter() - .any(|line| line.is_empty() || line.contains('\n') || line.contains('\r')) - { - return Err(format!("{error_label}批量记录必须是非空单行 JSON")); - } - // Reuse all secure-open, path/handle verification, tail repair, and durability - // checks. append_jsonl_line adds the final newline for the last record. - append_jsonl_line(path, &lines.join("\n"), error_label) -} - fn agent_db_has_conversation_message_audit_unlocked( file: &mut File, path: &Path, diff --git a/docs/project-memory/shared-memory/decision-log.md b/docs/project-memory/shared-memory/decision-log.md index 02746a98d..5f79b893d 100644 --- a/docs/project-memory/shared-memory/decision-log.md +++ b/docs/project-memory/shared-memory/decision-log.md @@ -913,9 +913,9 @@ Godot 编辑器操控复用既有 AGC 插件宿主、EditorAdapter、Runner 和 ## DirectProject 平行审计与请求分段计时退役边界(2026-09-23 核准) - 当前合同:完整用户与 Codex 完成 item 保存在 `.agent/conversations/project.jsonl`;GUI 回合不再创建 `runtime/direct-codex/turns/.jsonl` 或对应 `agent.db` 的 `direct.codex.turn` 摘要。旧 `offeredRead` / `firstDesign` 和请求分段计时不再属于生产保证,旧审计专题仅作历史追溯。 -- 代码依据:`agent/direct_runtime/user_input.rs` 明确传入 `audit: None`;`DirectCodexTurnAudit::start` 仅有测试调用,`DirectTurnMetrics` 和 Provider proxy 的 timing scope 随之不在生产构造。Provider proxy 本体及独立 model-usage observer 仍有现役用途,不能随旧计时链退役。 +- 实现边界:旧 `DirectCodexTurnAudit`、`DirectTurnMetrics`、可选审计 / 计时参数及专用批量 JSONL 追加包装已删除。Provider proxy 本体及独立 model-usage observer 仍有现役用途,不能随旧计时链退役;字节流透传和上游错误传递继续由现有测试验证。 - 保留边界:运行中的界面对话 / 工具耗时、`.agent/model-usage.jsonl`、产品埋点及 Runtime Agent 审计保持各自合同。`project.jsonl` 的完成 item 写入时间不等于 turn 起止或请求阶段计时;不据此补造旧历史耗时。本次不清理或迁移用户项目内的旧审计文件。 -- 维护依据:AGC 实施计划“Direct 历史、审计与耗时的现行边界”和“DirectProject Codex 原始历史与异常恢复”。残留旧 writer 与测试按未使用代码范围清理,不能以原审计方案为由永久保留或恢复。 +- 维护依据:AGC 实施计划“Direct 历史、审计与耗时的现行边界”和“DirectProject Codex 原始历史与异常恢复”。不能以原审计方案或已退役测试为由恢复旧 writer。 ## 2026-08-31 Direct 本轮附件只映射路径,不灌正文、不区别 GDD diff --git a/docs/project-memory/todos/【待解决】主要构建入口编译警告清单-2026-09-23.md b/docs/project-memory/todos/【待解决】主要构建入口编译警告清单-2026-09-23.md index 3058ea238..86cb40580 100644 --- a/docs/project-memory/todos/【待解决】主要构建入口编译警告清单-2026-09-23.md +++ b/docs/project-memory/todos/【待解决】主要构建入口编译警告清单-2026-09-23.md @@ -7,7 +7,7 @@ ## 1. 基线与范围 -首次诊断的基线提交为 `016356e509f11a1a638ce45ed51b9e40ef3e36a2`,诊断时工作树干净。环境为 Windows x64、Rust `1.98.1`、Node `v24.15.0`、npm `12.0.2`。第 2~6 节和附录保留首次诊断快照;后续处理状态及契约核查见第 7~8 节,不能将附录的全部条目视为仍未解决。 +首次诊断的基线提交为 `016356e509f11a1a638ce45ed51b9e40ef3e36a2`,诊断时工作树干净。环境为 Windows x64、Rust `1.98.1`、Node `v24.15.0`、npm `12.0.2`。第 2~6 节和附录保留首次诊断快照;后续处理状态及契约核查见第 7~9 节,不能将附录的全部条目视为仍未解决。 AGC Rust 使用默认 features、dev profile;237 是本轮编译器诊断数,不代表 237 个独立根因,也不是所有平台、features 和 test targets 的总数。另有 5 条不带常规源码 span 的 ts-rs 宏提示,不计入 237。 @@ -143,7 +143,8 @@ three、FBX/GLTF loader、OrbitControls 已动态 import。不能把“改为动 ### 剩余工作 - [x] 第一批:清理确认不改变行为的低风险项,修复预览部署器嵌套 npm 及移动 smoke 的 Windows 调用。 -- [x] 对当前 158 条 dead_code 和 27 条 unused_imports 完成调用、构建条件及保留契约的静态初审,结论见第 8 节;尚未执行第二批删除。 +- [x] 对清理前 158 条 dead_code 和 27 条 unused_imports 完成调用、构建条件及保留契约的静态初审,结论快照见第 8 节。 +- [x] 第二批首组:清理已明确退役的 Direct 审计 / 计时链 12 条诊断,结果见第 9 节。 - [ ] 单独核查保留的兼容重导出、身份/锁/门禁参数及 Direct 重试回合初始化,确认合同后再修改。 - [ ] 第二批:对当前 158 条 dead_code(首次快照为 165 条)核对测试、正式 features、平台及退役合同,逐项确定保留、条件编译或删除;涉及退役范围扩大时先补方案。 - [ ] 第三批:核实 5 条 ts-rs 提示的源类型及 TS 输出,保留现有反序列化约束。 @@ -199,6 +200,26 @@ three、FBX/GLTF loader、OrbitControls 已动态 import。不能把“改为动 本次仅补充核查文档,未删除或修改业务代码。实际执行 AGC 默认 Windows dev 编译;未重跑 Rust 单元测试、正式 editor features、Linux/macOS 编译或 GUI 端到端验收。后续实施按改动范围补定向验证。 +## 9. Direct 审计 / 计时链清理(2026-09-23) + +契约文档提交 `5a86a64e4` 后,删除 `direct_codex_audit.rs`、`direct_turn_metrics.rs`、专属批量 JSONL 追加包装,以及 Direct Runtime / app-server / Provider proxy 中的审计和计时参数、作用域与观察分支。退役模块的专属测试一并删除;原代理测试保留分块 SSE 字节透传和上游流错误传递断言,独立模型使用记录测试继续维护。 + +未改动现役 Provider 代理的鉴权 / 路由、模型使用记录、`project.jsonl`、GUI 回合与工具事件、付费 / 资产审计或用户项目内已有文件。 + +默认 Windows dev `cargo check --locked --offline` 通过,**195 → 183 条**,没有新增诊断: + +| 类别 | 清理前 | 清理后 | +| --- | ---: | ---: | +| `dead_code` | 158 | 146 | +| `unused_imports` | 27 | 27 | +| `unused_variables` | 9 | 9 | +| `unused_assignments` | 1 | 1 | +| **合计** | **195** | **183** | + +消除的是首次清单 W023~W025、W041~W048、W134,共 12 条;5 条 ts-rs 提示不在本次范围。第 8 节保留清理前的静态分类,其中 Direct 审计 / 计时 12 条已完成代码清理,其余候选仍按实际调用与测试范围推进。 + +验证:默认编译、修改文件 rustfmt、编码、文档索引和 diff 检查通过;独立只读审查未发现参数错位或现役链路回归。**41 个定向 Rust 测试通过**:`codex_provider_proxy::` 14 个(含鉴权、主站路由标记、并行、透传与模型记录),`direct_project_history::` 22 个,app-server 的 Direct 输入 / 宿主继续和事件投影 5 个。先通过 `cargo test --locked --offline` 构建并执行代理测试,再直接复用同一测试二进制执行后两组,均单线程运行。Unix 条件下的 fake app-server 协议用例未在 Windows 执行;未运行正式 editor features、Linux/macOS 编译、真实 Provider、GUI 或安装包验收。 + ## 附录:237 条编译器诊断位置 以下是诊断的人工可读整理,不保存原始构建日志、本机路径或 target 缓存路径。位置相对 `apps/ai-game-creator-shell/src-tauri/`;行号只对应本节基线,修改后以符号搜索及重新编译为准。每条保留一个主要 span,编号用于本清单内跟踪,不表示独立业务缺陷。生成产物项须回到 `build_support/runtime_prompt_bundle.rs` 处理。 diff --git a/docs/technical/【技术方案】AI游戏创作智能体App实施计划-2026-06-24.md b/docs/technical/【技术方案】AI游戏创作智能体App实施计划-2026-06-24.md index 88e0adaa2..15ec570cb 100644 --- a/docs/technical/【技术方案】AI游戏创作智能体App实施计划-2026-06-24.md +++ b/docs/technical/【技术方案】AI游戏创作智能体App实施计划-2026-06-24.md @@ -109,7 +109,7 @@ ### Direct 历史、审计与耗时的现行边界(2026-09-23 核准) - DirectProject 的完整回合条目只写入 `.agent/conversations/project.jsonl`,按 [原始历史与异常恢复](./【技术方案】DirectProject%20Codex原始历史与异常恢复-2026-09-04.md) 保存 canonical 用户条目和 Codex 完成的原始 item。工具调用与结果从同一事实源读取,前端继续使用安全投影;完整私有历史不能直接作为埋点或上传报告。 -- GUI 回合不再创建 `.agent/runtime/direct-codex/turns/.jsonl` 平行审计日志,也不再由该链写入 `agent.db` 的 `direct.codex.turn` 摘要。`direct_runtime/user_input.rs` 明确向回合执行器传入 `audit: None`;旧 `DirectCodexTurnAudit`、`DirectTurnMetrics` 构造和配套测试不构成现役行为承诺。旧审计专题归入历史,不要求恢复其 writer、`offeredRead` / `firstDesign` 投影或分段计时落盘。 +- GUI 回合不再创建 `.agent/runtime/direct-codex/turns/.jsonl` 平行审计日志,也不再由该链写入 `agent.db` 的 `direct.codex.turn` 摘要。旧 `DirectCodexTurnAudit`、`DirectTurnMetrics`、沿途审计 / 计时参数、专用批量 JSONL 追加包装和退役测试已删除。旧审计专题归入历史,不要求恢复其 writer、`offeredRead` / `firstDesign` 投影或分段计时落盘;代理的字节流透传、错误传递和独立模型记录测试继续维护。 - 会话运行中的对话和工具界面耗时仍按生命周期事件显示;模型请求 / 响应身份仍由 `.agent/model-usage.jsonl` 独立记录,Provider proxy 本体继续承担路由与响应型号观察。两者都不能证明 HTTP 首包、首 SSE、首内容 delta 或各阶段占用已形成生产分段计时记录;旧计时 fixture 通过也不能作为生产接入证据。`project.jsonl` 的 `recordedAt` 只是完成 item 的写入时间,不是 turn 起止时间;重进历史缺终态边界时不能承诺精确耗时,也不得补造请求阶段统计。 - 历史项目可能留有旧审计文件或摘要;本次契约收敛不删除、迁移或重写用户数据,也不要求新回合继续追加。现役模型使用记录、产品埋点、Runtime Agent 审计和付费 / 恢复凭证各守原有合同,不因 Direct 平行日志停用而退役。 diff --git a/docs/technical/【技术方案】DirectProject Codex原始历史与异常恢复-2026-09-04.md b/docs/technical/【技术方案】DirectProject Codex原始历史与异常恢复-2026-09-04.md index 635af7f5e..52918637c 100644 --- a/docs/technical/【技术方案】DirectProject Codex原始历史与异常恢复-2026-09-04.md +++ b/docs/technical/【技术方案】DirectProject Codex原始历史与异常恢复-2026-09-04.md @@ -6,7 +6,7 @@ DirectProject 只使用 `.agent/conversations/project.jsonl` 作为对话历史。历史保存 Codex Responses API 的完整 item,使聊天展示与新线程恢复使用同一份事实来源;两者只是不同读取动作。 -本方案只适用于 DirectProject,不改变 Agent session 历史。DirectProject 已停用 `runtime/direct-codex/turns` 平行审计账本,GUI 入口不再构造旧审计对象;完整回合条目统一来自本方案的 `project.jsonl`。旧文件不作为新回合必需产物,也不因本次契约更新被删除或迁移。 +本方案只适用于 DirectProject,不改变 Agent session 历史。DirectProject 已退役 `runtime/direct-codex/turns` 平行审计账本并删除旧审计 / 计时实现;完整回合条目统一来自本方案的 `project.jsonl`。用户项目中的旧日志不作为新回合必需产物,也不因代码清理被删除或迁移。 ## 文件格式 diff --git a/docs/technical/【技术方案】Direct回合行为审计账本-2026-08-31.md b/docs/technical/【技术方案】Direct回合行为审计账本-2026-08-31.md index ac547c5da..d01893b3a 100644 --- a/docs/technical/【技术方案】Direct回合行为审计账本-2026-08-31.md +++ b/docs/technical/【技术方案】Direct回合行为审计账本-2026-08-31.md @@ -323,7 +323,7 @@ chat_with_game_creator_direct_codex ## 9. 代码落地 -新增 [`apps/ai-game-creator-shell/src-tauri/src/agent/direct_codex_audit.rs`](../../apps/ai-game-creator-shell/src-tauri/src/agent/direct_codex_audit.rs): +原方案新增 `apps/ai-game-creator-shell/src-tauri/src/agent/direct_codex_audit.rs`(现已删除,以下仅作历史追溯): - `DirectCodexTurnAudit` - `start` / `observe_item` / `finish` diff --git a/docs/technical/【技术方案】客户端本地埋点与主站入库契约-2026-09-21.md b/docs/technical/【技术方案】客户端本地埋点与主站入库契约-2026-09-21.md index 20198c0d6..6cc64acbd 100644 --- a/docs/technical/【技术方案】客户端本地埋点与主站入库契约-2026-09-21.md +++ b/docs/technical/【技术方案】客户端本地埋点与主站入库契约-2026-09-21.md @@ -473,7 +473,7 @@ session.json 是本地恢复元数据,不是待上传事件;事件文件不 - 生命周期与配置:`apps/ai-game-creator-shell/src-tauri/src/main.rs`、`config.rs`、`platform_session.rs`。 - 项目创建与打开:`apps/ai-game-creator-shell/src-tauri/src/commands.rs`、`src/features/app-shell/useHomeProjectCreation.ts`;离开登记由 `WorkspaceLauncher.tsx` 保持。 - Direct 前端尝试与终态确认:`apps/ai-game-creator-shell/src/view/project-development/chat/controller/useDirectProjectChatController.ts`;沿用 `services/clientAnalytics.ts` 冻结账号代次、每次原生重试生成 attempt ID、只确认最后一次尝试。不在已退役的 App 聊天状态链恢复接线。 -- Direct 执行与审计:`apps/ai-game-creator-shell/src-tauri/src/agent/direct_runtime/user_input.rs`、`agent/direct_runtime/mod.rs`、`agent/direct_codex_audit.rs`。 +- Direct 执行与现役历史:`apps/ai-game-creator-shell/src-tauri/src/agent/direct_runtime/user_input.rs`、`agent/direct_runtime/mod.rs`、`agent/direct_project_history.rs`;旧 `direct_codex_audit.rs` 已删除,不作为埋点接入点。 - Design 执行与持久化:`apps/ai-game-creator-shell/src-tauri/src/agent/design_runtime.rs`、`agent/runtime_protocol/design_session.rs`、`agent/design_tools.rs`。 - revision:`apps/ai-game-creator-shell/src-tauri/src/agent/runtime_actions/project_gates.rs`,结合各实际写入调用方。 - 预览与保存:`apps/ai-game-creator-shell/src-tauri/src/preview.rs`、`ui_editor/persistence.rs`、`project/checkpoint.rs`。