diff --git a/crates/app/ironclaw_composition/src/observability/budget.rs b/crates/app/ironclaw_composition/src/observability/budget.rs index d073183738a..48ac9e139ac 100644 --- a/crates/app/ironclaw_composition/src/observability/budget.rs +++ b/crates/app/ironclaw_composition/src/observability/budget.rs @@ -122,6 +122,7 @@ mod tests { surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }; let _ = governor diff --git a/crates/contracts/ironclaw_loop_contracts/src/host/model.rs b/crates/contracts/ironclaw_loop_contracts/src/host/model.rs index a0f18c03e42..b8a9b583b2e 100644 --- a/crates/contracts/ironclaw_loop_contracts/src/host/model.rs +++ b/crates/contracts/ironclaw_loop_contracts/src/host/model.rs @@ -41,6 +41,9 @@ pub struct LoopModelRequest { /// Zero-based index into the host-resolved ordered fallback chain. #[serde(default)] pub fallback_index: u32, + /// Zero-based agent-loop iteration that issued this provider call. + #[serde(default)] + pub iteration: u32, #[serde(default, skip_serializing_if = "Option::is_none")] pub capability_view: Option, } diff --git a/crates/kernel/ironclaw_turns/tests/agent_loop_host_contract.rs b/crates/kernel/ironclaw_turns/tests/agent_loop_host_contract.rs index 8a42aa05142..b02308900e8 100644 --- a/crates/kernel/ironclaw_turns/tests/agent_loop_host_contract.rs +++ b/crates/kernel/ironclaw_turns/tests/agent_loop_host_contract.rs @@ -286,6 +286,7 @@ async fn host_managed_model_port_routes_gateway_and_emits_model_milestones() { surface_version: Some(CapabilitySurfaceVersion::new("surface-v1").unwrap()), model_preference: Some(context.resolved_run_profile.model_profile_id.clone()), fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -345,6 +346,7 @@ async fn host_managed_model_port_returns_response_when_model_started_milestone_f surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -390,6 +392,7 @@ async fn host_managed_model_port_returns_response_when_model_completed_milestone surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -432,6 +435,7 @@ async fn host_managed_model_port_sanitizes_gateway_errors() { surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -2913,6 +2917,7 @@ impl AgentLoopDriver for ReplyDriver { .clone(), ), fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -3205,6 +3210,7 @@ async fn host_managed_model_port_times_out_a_hung_gateway() { surface_version: Some(CapabilitySurfaceVersion::new("surface-v1").unwrap()), model_preference: Some(context.resolved_run_profile.model_profile_id.clone()), fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -3243,6 +3249,7 @@ async fn host_managed_model_port_allows_long_calls_that_keep_streaming_progress( surface_version: Some(CapabilitySurfaceVersion::new("surface-v1").unwrap()), model_preference: Some(context.resolved_run_profile.model_profile_id.clone()), fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -3860,6 +3867,7 @@ fn simple_model_request(context: &LoopRunContext) -> LoopModelRequest { surface_version: None, model_preference: Some(context.resolved_run_profile.model_profile_id.clone()), fallback_index: 0, + iteration: 0, capability_view: None, } } diff --git a/crates/loop/ironclaw_agent_loop/src/executor/failure_explanation.rs b/crates/loop/ironclaw_agent_loop/src/executor/failure_explanation.rs index fa2e5c96bc0..cfdd8edbc55 100644 --- a/crates/loop/ironclaw_agent_loop/src/executor/failure_explanation.rs +++ b/crates/loop/ironclaw_agent_loop/src/executor/failure_explanation.rs @@ -74,6 +74,7 @@ pub(super) async fn explain_failure( surface_version: None, model_preference: None, fallback_index: state.model_state.fallback_index, + iteration: state.iteration, capability_view: Some(LoopModelCapabilityView { visible_capability_ids: Vec::new(), }), diff --git a/crates/loop/ironclaw_agent_loop/src/executor/model.rs b/crates/loop/ironclaw_agent_loop/src/executor/model.rs index 095910bf8f1..686717d484c 100644 --- a/crates/loop/ironclaw_agent_loop/src/executor/model.rs +++ b/crates/loop/ironclaw_agent_loop/src/executor/model.rs @@ -79,6 +79,7 @@ impl ExecutorStage for ModelStage { surface_version: Some(surface_version.clone()), model_preference, fallback_index, + iteration: state.iteration, capability_view: Some(capability_view.clone()), }; let visible_capability_count = capability_view.visible_capability_ids.len(); diff --git a/crates/loop/ironclaw_agent_loop/src/executor/tests/cancellation.rs b/crates/loop/ironclaw_agent_loop/src/executor/tests/cancellation.rs index f59223ba28b..50a9cfe043d 100644 --- a/crates/loop/ironclaw_agent_loop/src/executor/tests/cancellation.rs +++ b/crates/loop/ironclaw_agent_loop/src/executor/tests/cancellation.rs @@ -488,7 +488,8 @@ async fn model_cancelled_returns_cancelled_without_retry() { "model cancelled", )]); let executor = CanonicalAgentLoopExecutor; - let state = LoopExecutionState::initial_for_run(host.run_context()); + let mut state = LoopExecutionState::initial_for_run(host.run_context()); + state.iteration = 3; let result = executor .execute_family(&crate::families::default(), &host, state) @@ -496,6 +497,7 @@ async fn model_cancelled_returns_cancelled_without_retry() { assert!(matches!(result, Err(AgentLoopExecutorError::Cancelled))); assert_eq!(host.model_requests().len(), 1); + assert_eq!(host.model_requests()[0].iteration, 3); } #[tokio::test] diff --git a/crates/loop/ironclaw_hooks/src/middleware/model_port.rs b/crates/loop/ironclaw_hooks/src/middleware/model_port.rs index fa9d9c3e64b..8a36384c50f 100644 --- a/crates/loop/ironclaw_hooks/src/middleware/model_port.rs +++ b/crates/loop/ironclaw_hooks/src/middleware/model_port.rs @@ -165,6 +165,7 @@ mod tests { surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, } } diff --git a/crates/loop/ironclaw_loop_host/src/budget_accountant.rs b/crates/loop/ironclaw_loop_host/src/budget_accountant.rs index 4d41daac880..725122a0dd3 100644 --- a/crates/loop/ironclaw_loop_host/src/budget_accountant.rs +++ b/crates/loop/ironclaw_loop_host/src/budget_accountant.rs @@ -899,6 +899,7 @@ mod tests { surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, } } diff --git a/crates/loop/ironclaw_loop_host/src/lib.rs b/crates/loop/ironclaw_loop_host/src/lib.rs index 9649247e18f..406a5a4bb47 100644 --- a/crates/loop/ironclaw_loop_host/src/lib.rs +++ b/crates/loop/ironclaw_loop_host/src/lib.rs @@ -195,6 +195,7 @@ pub use token_estimator::{ use tokio::sync::{Mutex, OnceCell}; use async_trait::async_trait; +use chrono::{DateTime, Utc}; use ironclaw_host_api::ids::{CapabilityId, RunId}; use ironclaw_loop_contracts::{ AgentLoopHostError, AgentLoopHostErrorKind, AgentLoopHostErrorReasonKind, @@ -1412,6 +1413,13 @@ where // with_budget_accountant"). let resolved_messages = self.resolve_model_messages(prompt_grant.messages).await?; + let diagnostic_requested_model = self.prompt_diagnostic_sink.as_ref().map(|_| { + requested_model_profile_id + .as_ref() + .unwrap_or(&model_profile_id) + .as_str() + .to_string() + }); if let Some(sink) = self.prompt_diagnostic_sink.as_ref() { let effective_model = self.gateway.diagnostic_effective_model( &model_profile_id, @@ -1460,6 +1468,8 @@ where } self.emit_model_started(requested_model_profile_id).await; + let diagnostic_started_at = Utc::now(); + let diagnostic_timer = Instant::now(); let host_request = HostManagedModelRequest { model_profile_id: model_profile_id.clone(), fallback_index: request.fallback_index, @@ -1500,6 +1510,16 @@ where self.gateway.stream_model(host_request).await }; + let diagnostic_effective_model = match &gateway_result { + Ok(response) => response + .diagnostic_effective_model + .as_ref() + .map(|model| model.as_str().to_string()), + Err(error) => error + .diagnostic_effective_model + .as_ref() + .map(|model| model.as_str().to_string()), + }; let host_response_result = match gateway_result { Ok(response) => { let HostManagedModelResponse { @@ -1508,12 +1528,17 @@ where output, usage, effective_fallback_index, + diagnostic_effective_model: _, } = response; if effective_fallback_index != Some(request.fallback_index) { - Err(AgentLoopHostError::new( + let error = AgentLoopHostError::new( AgentLoopHostErrorKind::Internal, "model gateway returned mismatched fallback route evidence", - )) + ); + Err(match usage { + Some(usage) => error.with_usage(usage), + None => error, + }) } else { let chunks = safe_text_deltas .into_iter() @@ -1534,6 +1559,46 @@ where Err(error) => Err(model_gateway_error(error)), }; + if let (Some(sink), Some(requested_model)) = ( + self.prompt_diagnostic_sink.as_ref(), + diagnostic_requested_model.as_ref(), + ) { + let (status, usage, failure_summary) = match &host_response_result { + Ok(response) => ( + HostManagedModelCallDiagnosticStatus::Succeeded, + diagnostic_usage(response.usage), + None, + ), + Err(error) => ( + HostManagedModelCallDiagnosticStatus::Failed, + diagnostic_usage(error.usage), + Some(error.safe_summary.as_str().to_string()), + ), + }; + let effective_model = diagnostic_effective_model.or_else(|| { + self.gateway + .diagnostic_effective_model( + &model_profile_id, + request.fallback_index, + self.run_context.resolved_model_route.as_ref(), + ) + .map(ProviderModelId::into_inner) + }); + sink.record_model_call(HostManagedModelCallDiagnosticCapture { + context: self.run_context.clone(), + iteration: request.iteration, + requested_model: requested_model.clone(), + effective_model, + started_at: diagnostic_started_at, + completed_at: Utc::now(), + duration_ms: u64::try_from(diagnostic_timer.elapsed().as_millis()) + .unwrap_or(u64::MAX), + status, + usage, + failure_summary, + }); + } + match host_response_result { Ok(response) => { self.emit_model_completed(model_profile_id).await; @@ -1942,6 +2007,8 @@ pub trait HostManagedModelStreamSink: Send + Sync { /// product events and apply their own authorization, redaction, and bounds. pub trait HostManagedPromptDiagnosticSink: Send + Sync { fn record_prompt(&self, capture: HostManagedPromptDiagnosticCapture); + + fn record_model_call(&self, _capture: HostManagedModelCallDiagnosticCapture) {} } /// Validated concrete provider model identifier used only for diagnostics. @@ -1986,14 +2053,20 @@ impl fmt::Display for ProviderModelId { /// Bounded, non-blocking decorator for best-effort prompt diagnostics. /// -/// Captures are dropped when the worker cannot keep up so diagnostic work -/// never adds backpressure to the provider request path. +/// Prompt and model-call captures share one ordered queue. Captures are dropped +/// when the worker cannot keep up so diagnostic work never adds backpressure to +/// the provider request path. pub struct BufferedPromptDiagnosticSink { - sender: tokio::sync::mpsc::Sender, + sender: tokio::sync::mpsc::Sender, } pub const DEFAULT_PROMPT_DIAGNOSTIC_QUEUE_CAPACITY: usize = 8; +enum BufferedDiagnosticCapture { + Prompt(HostManagedPromptDiagnosticCapture), + ModelCall(HostManagedModelCallDiagnosticCapture), +} + impl BufferedPromptDiagnosticSink { pub fn new( inner: Arc, @@ -2008,20 +2081,22 @@ impl BufferedPromptDiagnosticSink { runtime.spawn(async move { while let Some(capture) = receiver.recv().await { let sink = Arc::clone(&inner); - if let Err(error) = - tokio::task::spawn_blocking(move || sink.record_prompt(capture)).await + if let Err(error) = tokio::task::spawn_blocking(move || match capture { + BufferedDiagnosticCapture::Prompt(capture) => sink.record_prompt(capture), + BufferedDiagnosticCapture::ModelCall(capture) => { + sink.record_model_call(capture); + } + }) + .await { - tracing::debug!(%error, "prompt diagnostic worker failed"); + tracing::debug!(%error, "diagnostic worker failed"); } } }); Ok(Self { sender }) } -} -impl HostManagedPromptDiagnosticSink for BufferedPromptDiagnosticSink { - fn record_prompt(&self, capture: HostManagedPromptDiagnosticCapture) { - let run_id = capture.context.run_id; + fn enqueue(&self, run_id: TurnRunId, capture: BufferedDiagnosticCapture) { if let Err(error) = self.sender.try_send(capture) { let queue_state = match error { tokio::sync::mpsc::error::TrySendError::Full(_) => "full", @@ -2030,12 +2105,24 @@ impl HostManagedPromptDiagnosticSink for BufferedPromptDiagnosticSink { tracing::debug!( %run_id, queue_state, - "dropping best-effort prompt diagnostic capture" + "dropping best-effort diagnostic capture" ); } } } +impl HostManagedPromptDiagnosticSink for BufferedPromptDiagnosticSink { + fn record_prompt(&self, capture: HostManagedPromptDiagnosticCapture) { + let run_id = capture.context.run_id; + self.enqueue(run_id, BufferedDiagnosticCapture::Prompt(capture)); + } + + fn record_model_call(&self, capture: HostManagedModelCallDiagnosticCapture) { + let run_id = capture.context.run_id; + self.enqueue(run_id, BufferedDiagnosticCapture::ModelCall(capture)); + } +} + #[derive(Debug, Clone)] pub struct HostManagedPromptDiagnosticCapture { pub context: LoopRunContext, @@ -2056,6 +2143,35 @@ pub struct HostManagedPromptDiagnosticMessage { pub content: String, } +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum HostManagedModelCallDiagnosticStatus { + Succeeded, + Failed, +} + +#[derive(Debug, Clone)] +pub struct HostManagedModelCallDiagnosticCapture { + pub context: LoopRunContext, + pub iteration: u32, + pub requested_model: String, + pub effective_model: Option, + pub started_at: DateTime, + pub completed_at: DateTime, + pub duration_ms: u64, + pub status: HostManagedModelCallDiagnosticStatus, + pub usage: Option, + pub failure_summary: Option, +} + +fn diagnostic_usage(usage: Option) -> Option { + usage.filter(|usage| { + usage.input_tokens > 0 + || usage.output_tokens > 0 + || usage.cache_read_input_tokens > 0 + || usage.cache_creation_input_tokens > 0 + }) +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub struct HostManagedModelRequest { pub model_profile_id: ModelProfileId, @@ -2188,6 +2304,10 @@ pub struct HostManagedModelResponse { /// Authoritative ordered-chain index used for this successful call. #[serde(default, skip_serializing_if = "Option::is_none")] pub effective_fallback_index: Option, + /// Concrete provider model that handled this call. This is runtime + /// diagnostic evidence, not an authority or routing input. + #[serde(skip)] + pub diagnostic_effective_model: Option>, } impl HostManagedModelResponse { @@ -2202,6 +2322,7 @@ impl HostManagedModelResponse { }), usage: None, effective_fallback_index: Some(0), + diagnostic_effective_model: None, } } @@ -2229,6 +2350,7 @@ impl HostManagedModelResponse { output: ParentLoopOutput::CapabilityCalls(calls), usage: None, effective_fallback_index: Some(0), + diagnostic_effective_model: None, } } @@ -2253,6 +2375,11 @@ impl HostManagedModelResponse { self.effective_fallback_index = Some(fallback_index); self } + + pub fn with_diagnostic_effective_model(mut self, model: impl Into) -> Self { + self.diagnostic_effective_model = Some(Arc::new(model.into())); + self + } } fn sanitized_reasoning_deltas(reasoning: Option) -> Vec { @@ -2309,7 +2436,7 @@ pub struct HostManagedModelError { pub kind: HostManagedModelErrorKind, pub safe_summary: String, pub reason_kind: Option, - pub gate_ref: Option, + pub gate_ref: Option>, /// Provider-supplied retry delay. Typed so the recovery strategy does not /// have to parse model-visible detail text. pub retry_after_ms: Option, @@ -2318,6 +2445,9 @@ pub struct HostManagedModelError { pub next_fallback_index: Option, /// Provider-reported usage for a call that consumed tokens before failing. pub usage: Option, + /// Concrete provider model that handled this failed call. This is runtime + /// diagnostic evidence, not an authority or routing input. + pub diagnostic_effective_model: Option>, /// Model-visible, secret-scrubbed raw cause (status line, provider body /// snippet). Unlike `safe_summary`, this carries the original message so the /// failure explainer can describe the real fault. Secret VALUES must be @@ -2337,6 +2467,7 @@ impl HostManagedModelError { retry_after_ms: None, next_fallback_index: None, usage: None, + diagnostic_effective_model: None, detail: None, } } @@ -2350,6 +2481,7 @@ impl HostManagedModelError { retry_after_ms: None, next_fallback_index: None, usage: None, + diagnostic_effective_model: None, detail: None, } } @@ -2378,7 +2510,7 @@ impl HostManagedModelError { } pub fn with_gate_ref(mut self, gate_ref: LoopGateRef) -> Self { - self.gate_ref = Some(gate_ref); + self.gate_ref = Some(Box::new(gate_ref)); self } @@ -2402,6 +2534,11 @@ impl HostManagedModelError { self.usage = Some(usage); self } + + pub fn with_diagnostic_effective_model(mut self, model: impl Into) -> Self { + self.diagnostic_effective_model = Some(Arc::new(model.into())); + self + } } fn validate_thread_scope_for_run( @@ -2791,7 +2928,7 @@ fn model_gateway_error(error: HostManagedModelError) -> AgentLoopHostError { host_error = host_error.with_reason_kind(reason_kind); } if let Some(gate_ref) = error.gate_ref { - host_error = host_error.with_gate_ref(gate_ref); + host_error = host_error.with_gate_ref(*gate_ref); } if let Some(retry_after_ms) = error.retry_after_ms { host_error = host_error.with_retry_after_ms(retry_after_ms); @@ -2880,6 +3017,7 @@ mod tests { struct BlockingPromptDiagnosticSink { calls: std::sync::atomic::AtomicUsize, + model_calls: std::sync::atomic::AtomicUsize, first_started: tokio::sync::Notify, release_first: (std::sync::Mutex, std::sync::Condvar), } @@ -2888,6 +3026,7 @@ mod tests { fn new() -> Self { Self { calls: std::sync::atomic::AtomicUsize::new(0), + model_calls: std::sync::atomic::AtomicUsize::new(0), first_started: tokio::sync::Notify::new(), release_first: (std::sync::Mutex::new(false), std::sync::Condvar::new()), } @@ -2912,6 +3051,10 @@ mod tests { } } } + + fn record_model_call(&self, _capture: HostManagedModelCallDiagnosticCapture) { + self.model_calls.fetch_add(1, Ordering::SeqCst); + } } fn prompt_diagnostic_capture_for_test() -> HostManagedPromptDiagnosticCapture { @@ -2980,6 +3123,39 @@ mod tests { assert_eq!(inner.calls.load(Ordering::SeqCst), 2); } + #[tokio::test] + async fn buffered_prompt_diagnostics_forward_model_calls() { + let inner = Arc::new(BlockingPromptDiagnosticSink::new()); + let buffered = BufferedPromptDiagnosticSink::new( + inner.clone() as Arc, + 1, + ) + .expect("buffered sink"); + let context = prompt_diagnostic_capture_for_test().context; + let now = Utc::now(); + + buffered.record_model_call(HostManagedModelCallDiagnosticCapture { + context, + iteration: 1, + requested_model: "interactive_model".to_string(), + effective_model: Some("provider-model".to_string()), + started_at: now, + completed_at: now, + duration_ms: 1, + status: HostManagedModelCallDiagnosticStatus::Succeeded, + usage: None, + failure_summary: None, + }); + + tokio::time::timeout(Duration::from_secs(1), async { + while inner.model_calls.load(Ordering::SeqCst) < 1 { + tokio::task::yield_now().await; + } + }) + .await + .expect("model-call capture drains"); + } + #[test] fn provider_model_id_enforces_bounded_route_grammar() { assert_eq!( diff --git a/crates/loop/ironclaw_loop_host/src/model_gateway.rs b/crates/loop/ironclaw_loop_host/src/model_gateway.rs index eb5143172fb..09fcefe9d35 100644 --- a/crates/loop/ironclaw_loop_host/src/model_gateway.rs +++ b/crates/loop/ironclaw_loop_host/src/model_gateway.rs @@ -455,7 +455,8 @@ where )?; add_request_metadata(&mut completion, &model_profile_id, run_id, turn_id); - complete_model_request( + let diagnostic_effective_model = replay_identity.provider_model_id.clone(); + let result = complete_model_request( self.provider.as_ref(), completion, None, @@ -464,8 +465,8 @@ where ProviderRequestContext::new(replay_identity, next_fallback_index), Some(self.prompt_cache_scope(run_id)), ) - .await - .map(|response| response.with_effective_fallback_index(effective_fallback_index)) + .await; + with_model_diagnostic_evidence(result, effective_fallback_index, diagnostic_effective_model) } async fn stream_model_with_progress( @@ -503,7 +504,8 @@ where )?; add_request_metadata(&mut completion, &model_profile_id, run_id, turn_id); - complete_model_request( + let diagnostic_effective_model = replay_identity.provider_model_id.clone(); + let result = complete_model_request( self.provider.as_ref(), completion, None, @@ -512,8 +514,8 @@ where ProviderRequestContext::new(replay_identity, next_fallback_index), Some(self.prompt_cache_scope(run_id)), ) - .await - .map(|response| response.with_effective_fallback_index(effective_fallback_index)) + .await; + with_model_diagnostic_evidence(result, effective_fallback_index, diagnostic_effective_model) } async fn stream_model_with_capabilities( @@ -555,7 +557,8 @@ where "run={run_id}\nturn={turn_id}\nmodel_call={}", self.provider_turn_sequence.fetch_add(1, Ordering::Relaxed) ); - complete_model_request( + let diagnostic_effective_model = replay_identity.provider_model_id.clone(); + let result = complete_model_request( self.provider.as_ref(), completion, Some(capabilities), @@ -564,8 +567,8 @@ where ProviderRequestContext::new(replay_identity, next_fallback_index), Some(self.prompt_cache_scope(run_id)), ) - .await - .map(|response| response.with_effective_fallback_index(effective_fallback_index)) + .await; + with_model_diagnostic_evidence(result, effective_fallback_index, diagnostic_effective_model) } async fn stream_model_with_capabilities_and_progress( @@ -608,7 +611,8 @@ where "run={run_id}\nturn={turn_id}\nmodel_call={}", self.provider_turn_sequence.fetch_add(1, Ordering::Relaxed) ); - complete_model_request( + let diagnostic_effective_model = replay_identity.provider_model_id.clone(); + let result = complete_model_request( self.provider.as_ref(), completion, Some(capabilities), @@ -617,8 +621,8 @@ where ProviderRequestContext::new(replay_identity, next_fallback_index), Some(self.prompt_cache_scope(run_id)), ) - .await - .map(|response| response.with_effective_fallback_index(effective_fallback_index)) + .await; + with_model_diagnostic_evidence(result, effective_fallback_index, diagnostic_effective_model) } } @@ -783,7 +787,8 @@ where add_request_metadata(&mut completion, &model_profile_id, run_id, turn_id); add_route_metadata(&mut completion, &snapshot); - complete_model_request( + let diagnostic_effective_model = replay_identity.provider_model_id.clone(); + let result = complete_model_request( provider.as_ref(), completion, None, @@ -792,8 +797,8 @@ where ProviderRequestContext::new(replay_identity, next_fallback_index), Some(self.prompt_cache_scope(run_id)), ) - .await - .map(|response| response.with_effective_fallback_index(effective_fallback_index)) + .await; + with_model_diagnostic_evidence(result, effective_fallback_index, diagnostic_effective_model) } async fn stream_model_with_progress( @@ -824,7 +829,8 @@ where add_request_metadata(&mut completion, &model_profile_id, run_id, turn_id); add_route_metadata(&mut completion, &snapshot); - complete_model_request( + let diagnostic_effective_model = replay_identity.provider_model_id.clone(); + let result = complete_model_request( provider.as_ref(), completion, None, @@ -833,8 +839,8 @@ where ProviderRequestContext::new(replay_identity, next_fallback_index), Some(self.prompt_cache_scope(run_id)), ) - .await - .map(|response| response.with_effective_fallback_index(effective_fallback_index)) + .await; + with_model_diagnostic_evidence(result, effective_fallback_index, diagnostic_effective_model) } async fn stream_model_with_capabilities( @@ -869,7 +875,8 @@ where "run={run_id}\nturn={turn_id}\nmodel_call={}", self.provider_turn_sequence.fetch_add(1, Ordering::Relaxed) ); - complete_model_request( + let diagnostic_effective_model = replay_identity.provider_model_id.clone(); + let result = complete_model_request( provider.as_ref(), completion, Some(capabilities), @@ -878,8 +885,8 @@ where ProviderRequestContext::new(replay_identity, next_fallback_index), Some(self.prompt_cache_scope(run_id)), ) - .await - .map(|response| response.with_effective_fallback_index(effective_fallback_index)) + .await; + with_model_diagnostic_evidence(result, effective_fallback_index, diagnostic_effective_model) } async fn stream_model_with_capabilities_and_progress( @@ -915,7 +922,8 @@ where "run={run_id}\nturn={turn_id}\nmodel_call={}", self.provider_turn_sequence.fetch_add(1, Ordering::Relaxed) ); - complete_model_request( + let diagnostic_effective_model = replay_identity.provider_model_id.clone(); + let result = complete_model_request( provider.as_ref(), completion, Some(capabilities), @@ -924,8 +932,8 @@ where ProviderRequestContext::new(replay_identity, next_fallback_index), Some(self.prompt_cache_scope(run_id)), ) - .await - .map(|response| response.with_effective_fallback_index(effective_fallback_index)) + .await; + with_model_diagnostic_evidence(result, effective_fallback_index, diagnostic_effective_model) } } @@ -967,6 +975,20 @@ fn add_request_metadata( .insert("run_id".to_string(), run_id.to_string()); } +fn with_model_diagnostic_evidence( + result: Result, + effective_fallback_index: u32, + effective_model: String, +) -> Result { + result + .map(|response| { + response + .with_effective_fallback_index(effective_fallback_index) + .with_diagnostic_effective_model(effective_model.clone()) + }) + .map_err(|error| error.with_diagnostic_effective_model(effective_model)) +} + fn resolve_fallback_route

( provider: &P, fallback_index: u32, diff --git a/crates/loop/ironclaw_loop_host/tests/llm_gateway.rs b/crates/loop/ironclaw_loop_host/tests/llm_gateway.rs index ccc26136af9..1a12f2824c0 100644 --- a/crates/loop/ironclaw_loop_host/tests/llm_gateway.rs +++ b/crates/loop/ironclaw_loop_host/tests/llm_gateway.rs @@ -135,6 +135,13 @@ async fn gateway_calls_llm_provider_for_allowed_model_profile() { response.safe_text_deltas, vec!["assistant response".to_string()] ); + assert_eq!( + response + .diagnostic_effective_model + .as_ref() + .map(|model| model.as_str()), + Some("host-selected-model") + ); let requests = provider.requests.lock().unwrap(); assert_eq!(requests.len(), 1); assert_eq!(requests[0].model.as_deref(), Some("host-selected-model")); @@ -3010,6 +3017,7 @@ async fn production_loop_model_gateway_rejects_forged_context_summary_before_pro surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -3060,6 +3068,7 @@ async fn production_loop_model_gateway_rejects_unvalidated_surface_before_provid surface_version: Some(CapabilitySurfaceVersion::new("surface-stale").unwrap()), model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -3121,6 +3130,13 @@ async fn gateway_sanitizes_provider_errors() { .unwrap_err(); assert_eq!(error.kind, HostManagedModelErrorKind::Unavailable); + assert_eq!( + error + .diagnostic_effective_model + .as_ref() + .map(|model| model.as_str()), + Some("host-selected-model") + ); assert!(!error.safe_summary.contains("RAW_PROVIDER_SECRET")); assert!(!format!("{error:?}").contains("RAW_PROVIDER_SECRET")); } @@ -4075,6 +4091,7 @@ async fn production_loop_request_with_safety_and_inline_messages( surface_version: None, model_preference, fallback_index: 0, + iteration: 0, capability_view: None, } } diff --git a/crates/loop/ironclaw_loop_host/tests/thread_loop_host_contract.rs b/crates/loop/ironclaw_loop_host/tests/thread_loop_host_contract.rs index cbb5ecf6f9d..249bb46a52f 100644 --- a/crates/loop/ironclaw_loop_host/tests/thread_loop_host_contract.rs +++ b/crates/loop/ironclaw_loop_host/tests/thread_loop_host_contract.rs @@ -26,27 +26,29 @@ use ironclaw_loop_contracts::{ LoopContextBundle, LoopContextCompactionKind, LoopContextMessage, LoopContextPort, LoopContextRequest, LoopContextSnippet, LoopDriverNoteKind, LoopHostMilestoneKind, LoopHostMilestoneSink, LoopInputCursor, LoopInputCursorToken, LoopModelCapabilityView, - LoopModelMessage, LoopModelPort, LoopModelRequest, LoopModelRouteSnapshot, LoopPromptBundle, - LoopPromptBundleAuthority, LoopPromptBundleRef, LoopPromptBundleRequest, LoopPromptPort, - LoopRequest, LoopRequestBatch, LoopRunContext, LoopTranscriptPort, ModelProfileId, - ModelVisibleToolObservation, ObservationTrust, ParentLoopOutput, PersonalContextPolicy, - PromptMode, PromptSkillContextMetadata, ProviderToolCallReference, ProviderToolCallReplay, - ProviderToolDefinition, RunProfileResolutionRequest, RunProfileResolver, SkillName, - SkillTrustLevel, SkillVisibility, ToolObservationDetail, ToolObservationStatus, - UpdateAssistantDraft, VisibleCapabilityRequest, VisibleCapabilitySurface, resolution, + LoopModelMessage, LoopModelPort, LoopModelRequest, LoopModelRouteSnapshot, LoopModelUsage, + LoopPromptBundle, LoopPromptBundleAuthority, LoopPromptBundleRef, LoopPromptBundleRequest, + LoopPromptPort, LoopRequest, LoopRequestBatch, LoopRunContext, LoopTranscriptPort, + ModelProfileId, ModelVisibleToolObservation, ObservationTrust, ParentLoopOutput, + PersonalContextPolicy, PromptMode, PromptSkillContextMetadata, ProviderToolCallReference, + ProviderToolCallReplay, ProviderToolDefinition, RunProfileResolutionRequest, + RunProfileResolver, SkillName, SkillTrustLevel, SkillVisibility, ToolObservationDetail, + ToolObservationStatus, UpdateAssistantDraft, VisibleCapabilityRequest, + VisibleCapabilitySurface, resolution, }; use ironclaw_loop_host::{ EmptyLoopCapabilityPort, HostIdentityContextBuildError, HostIdentityContextCandidate, - HostIdentityContextSource, HostIdentityMessageContent, HostManagedModelError, - HostManagedModelErrorKind, HostManagedModelGateway, HostManagedModelMessageRole, - HostManagedModelRequest, HostManagedModelResponse, HostManagedPromptDiagnosticCapture, - HostManagedPromptDiagnosticSink, HostManagedToolResultContent, HostSkillContextBuildError, - HostSkillContextCandidate, HostSkillContextSource, IdentityApplicability, IdentityBudget, - IdentityFileName, LoopAttachmentReadError, LoopAttachmentReadPort, PromptContextTokenBudget, - ProviderModelId, SkillBundleContextSource, SkillBundleDescriptor, SkillBundleId, - SkillBundleSource, SkillBundleSourceError, SkillFilePath, SkillSourceKind, - ThreadBackedLoopContextPort, ThreadBackedLoopModelPort, ThreadBackedLoopTranscriptPort, - ThreadContextWindowCache, build_skill_run_snapshot, identity_message_ref, + HostIdentityContextSource, HostIdentityMessageContent, HostManagedModelCallDiagnosticCapture, + HostManagedModelError, HostManagedModelErrorKind, HostManagedModelGateway, + HostManagedModelMessageRole, HostManagedModelRequest, HostManagedModelResponse, + HostManagedPromptDiagnosticCapture, HostManagedPromptDiagnosticSink, + HostManagedToolResultContent, HostSkillContextBuildError, HostSkillContextCandidate, + HostSkillContextSource, IdentityApplicability, IdentityBudget, IdentityFileName, + LoopAttachmentReadError, LoopAttachmentReadPort, PromptContextTokenBudget, ProviderModelId, + SkillBundleContextSource, SkillBundleDescriptor, SkillBundleId, SkillBundleSource, + SkillBundleSourceError, SkillFilePath, SkillSourceKind, ThreadBackedLoopContextPort, + ThreadBackedLoopModelPort, ThreadBackedLoopTranscriptPort, ThreadContextWindowCache, + build_skill_run_snapshot, identity_message_ref, }; use ironclaw_outbound::{ OutboundError, OutboundStateStore, ReplyAttachmentHandle, ReplyAttachmentIntent, @@ -213,6 +215,7 @@ async fn model_port_empty_request_applies_prompt_token_budget_to_context_fallbac surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -227,7 +230,16 @@ async fn model_port_empty_request_applies_prompt_token_budget_to_context_fallbac #[tokio::test] async fn model_port_records_resolved_prompt_with_fallback_model_at_the_host_boundary() { let fixture = ThreadFixture::new_with_user_content("diagnostic prompt body").await; - let gateway = Arc::new(RecordingGateway::reply_with_fallback("model says hi", 2)); + let gateway = Arc::new(RecordingGateway::reply_with_usage_and_fallback( + "model says hi", + LoopModelUsage { + input_tokens: 21, + output_tokens: 8, + cache_read_input_tokens: 5, + cache_creation_input_tokens: 3, + }, + 2, + )); let sink = Arc::new(RecordingPromptDiagnosticSink::default()); let messages = user_model_messages(&fixture); let bundle = LoopPromptBundle { @@ -267,6 +279,7 @@ async fn model_port_records_resolved_prompt_with_fallback_model_at_the_host_boun surface_version: None, model_preference: None, fallback_index: 2, + iteration: 7, capability_view: Some(LoopModelCapabilityView { visible_capability_ids: vec![CapabilityId::new("filesystem.read").expect("capability")], }), @@ -288,6 +301,142 @@ async fn model_port_records_resolved_prompt_with_fallback_model_at_the_host_boun Some("fallback-provider-model") ); assert_eq!(captures[0].context_limit, 64_000); + drop(captures); + let model_calls = sink.model_calls.lock().expect("model calls"); + assert_eq!(model_calls.len(), 1); + assert_eq!(model_calls[0].iteration, 7); + assert_eq!(model_calls[0].requested_model, "interactive_model"); + assert_eq!( + model_calls[0].effective_model.as_deref(), + Some("provider-model-from-response") + ); + assert_eq!( + model_calls[0].usage.map(|usage| usage.input_tokens), + Some(21) + ); +} + +#[tokio::test] +async fn model_port_keeps_effective_model_unavailable_without_provider_evidence() { + let fixture = ThreadFixture::new_with_user_content("diagnostic prompt body").await; + let gateway = Arc::new(MissingDiagnosticModelGateway); + let sink = Arc::new(RecordingPromptDiagnosticSink::default()); + let messages = user_model_messages(&fixture); + issue_prompt_grant(&fixture.run_context, &messages); + + ThreadBackedLoopModelPort::new( + Arc::clone(&fixture.thread_service), + fixture.thread_scope.clone(), + fixture.run_context.clone(), + gateway, + 16, + ) + .with_prompt_diagnostic_sink(sink.clone()) + .stream_model(LoopModelRequest { + inline_messages: Vec::new(), + messages, + surface_version: None, + model_preference: None, + fallback_index: 2, + iteration: 7, + capability_view: None, + }) + .await + .expect("model response"); + + let model_calls = sink.model_calls.lock().expect("model calls"); + assert_eq!(model_calls.len(), 1); + assert_eq!(model_calls[0].requested_model, "interactive_model"); + assert_eq!(model_calls[0].effective_model, None); +} + +#[tokio::test] +async fn model_port_retains_usage_reported_by_failed_calls() { + let fixture = ThreadFixture::new_with_user_content("failed diagnostic prompt").await; + let usage = LoopModelUsage { + input_tokens: 34, + output_tokens: 2, + cache_read_input_tokens: 8, + cache_creation_input_tokens: 1, + }; + let gateway = Arc::new(RecordingGateway::model_error_with_usage( + HostManagedModelErrorKind::ProviderUnavailable, + "model provider unavailable", + usage, + )); + let sink = Arc::new(RecordingPromptDiagnosticSink::default()); + let messages = user_model_messages(&fixture); + issue_prompt_grant(&fixture.run_context, &messages); + + let error = ThreadBackedLoopModelPort::new( + Arc::clone(&fixture.thread_service), + fixture.thread_scope.clone(), + fixture.run_context.clone(), + gateway, + 16, + ) + .with_prompt_diagnostic_sink(sink.clone()) + .stream_model(LoopModelRequest { + inline_messages: Vec::new(), + messages, + surface_version: None, + model_preference: None, + fallback_index: 0, + iteration: 0, + capability_view: None, + }) + .await + .expect_err("model call fails"); + assert_eq!(error.kind, AgentLoopHostErrorKind::Unavailable); + + let model_calls = sink.model_calls.lock().expect("model calls"); + assert_eq!(model_calls.len(), 1); + assert_eq!( + model_calls[0].status, + ironclaw_loop_host::HostManagedModelCallDiagnosticStatus::Failed + ); + assert_eq!(model_calls[0].usage, Some(usage)); + assert_eq!( + model_calls[0].effective_model.as_deref(), + Some("provider-model-from-error") + ); + assert_eq!( + model_calls[0].failure_summary.as_deref(), + Some("model provider unavailable") + ); +} + +#[tokio::test] +async fn model_port_keeps_omitted_usage_unavailable() { + let fixture = ThreadFixture::new_with_user_content("usage unavailable prompt").await; + let gateway = Arc::new(RecordingGateway::reply("model says hi")); + let sink = Arc::new(RecordingPromptDiagnosticSink::default()); + let messages = user_model_messages(&fixture); + issue_prompt_grant(&fixture.run_context, &messages); + + ThreadBackedLoopModelPort::new( + Arc::clone(&fixture.thread_service), + fixture.thread_scope.clone(), + fixture.run_context.clone(), + gateway, + 16, + ) + .with_prompt_diagnostic_sink(sink.clone()) + .stream_model(LoopModelRequest { + inline_messages: Vec::new(), + messages, + surface_version: None, + model_preference: None, + fallback_index: 0, + iteration: 0, + capability_view: None, + }) + .await + .expect("model call succeeds"); + + let model_calls = sink.model_calls.lock().expect("model calls"); + assert_eq!(model_calls.len(), 1); + assert_eq!(model_calls[0].usage, None); } #[tokio::test] @@ -316,6 +465,7 @@ async fn model_port_records_full_capability_surface_when_request_has_no_view() { surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -358,6 +508,7 @@ async fn model_port_continues_when_diagnostic_capability_lookup_fails() { surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -429,6 +580,7 @@ async fn prompt_and_model_ports_share_cached_context_window_for_one_request() { surface_version: prompt_bundle.surface_version, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -484,6 +636,7 @@ async fn model_port_reuses_smaller_prompt_context_window_for_explicit_prompt_ref surface_version: prompt_bundle.surface_version, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -578,6 +731,7 @@ async fn context_window_cache_does_not_cross_thread_scope_boundaries() { surface_version: prompt_bundle.surface_version, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -1500,6 +1654,7 @@ async fn prompt_and_model_ports_materialize_trusted_identity_content() { surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -1541,6 +1696,7 @@ async fn model_port_limits_provider_tool_definitions_to_model_visible_capability surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: Some(LoopModelCapabilityView { visible_capability_ids: vec![allowed_id], }), @@ -1584,6 +1740,7 @@ async fn model_port_maps_invalid_model_output_to_recoverable_model_error() { surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -1629,6 +1786,7 @@ async fn model_port_preserves_capability_info_for_filtered_capability_view() { surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: Some(LoopModelCapabilityView { visible_capability_ids: vec![allowed_id], }), @@ -1901,6 +2059,7 @@ async fn prompt_and_model_ports_send_selected_skill_context_to_gateway() { surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -2004,6 +2163,7 @@ async fn prompt_and_model_ports_resolve_skill_refs_after_prompt_sorting() { surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -2093,6 +2253,7 @@ async fn prompt_and_model_ports_resolve_instruction_memory_and_identity_refs() { surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -2141,6 +2302,7 @@ async fn model_port_rejects_policy_denied_identity_ref_before_gateway_call() { surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -2342,6 +2504,7 @@ async fn prompt_and_model_ports_keep_duplicate_skill_names_distinct() { surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -2411,6 +2574,7 @@ async fn model_port_rejects_skill_context_refs_when_source_changes_after_prompt_ surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -3788,6 +3952,7 @@ async fn model_port_resolves_thread_message_refs_and_delegates_to_gateway() { surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -3834,6 +3999,7 @@ async fn model_port_rejects_mismatched_fallback_route_evidence() { surface_version: None, model_preference: None, fallback_index: 2, + iteration: 0, capability_view: None, }) .await @@ -3880,6 +4046,7 @@ async fn model_port_rejects_missing_fallback_route_evidence() { surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -3916,6 +4083,7 @@ async fn model_port_accepts_matching_fallback_route_evidence() { surface_version: None, model_preference: None, fallback_index: 2, + iteration: 0, capability_view: None, }) .await @@ -4006,6 +4174,7 @@ async fn model_port_reads_image_attachment_bytes_into_model_image_parts() { surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -4055,6 +4224,7 @@ async fn model_port_merges_consecutive_text_user_messages_for_prompt() { model_preference: None, capability_view: None, fallback_index: 0, + iteration: 0, }) .await .unwrap(); @@ -4112,6 +4282,7 @@ async fn model_port_threads_resolved_model_route_snapshot_to_gateway() { surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -4157,6 +4328,7 @@ async fn model_port_resolves_explicit_refs_that_fall_outside_context_window() { surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -4229,6 +4401,7 @@ async fn model_port_preserves_provider_metadata_for_explicit_refs_outside_contex surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -4352,6 +4525,7 @@ async fn model_port_round_trips_tool_result_reference_context_as_typed_model_inp surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -4403,6 +4577,7 @@ async fn model_port_rejects_malformed_tool_result_reference_content() { surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -4447,6 +4622,7 @@ async fn model_port_rejects_missing_explicit_tool_result_reference_before_gatewa surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -4491,6 +4667,7 @@ async fn model_port_emits_model_milestones_without_prompt_or_output_payloads() { .clone(), ), fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -4550,6 +4727,7 @@ async fn model_port_emits_started_and_failed_milestones_when_gateway_fails() { surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -4609,6 +4787,7 @@ async fn model_port_logs_model_started_milestone_failure_without_losing_response surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -4651,6 +4830,7 @@ async fn model_port_logs_model_completed_milestone_failure_without_losing_respon surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -4691,6 +4871,7 @@ async fn model_port_rejects_message_role_that_disagrees_with_thread_record() { surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -4721,6 +4902,7 @@ async fn model_port_surfaces_fail_closed_gateway_policy_errors_without_raw_detai surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -4757,6 +4939,7 @@ async fn model_port_replaces_invalid_gateway_safe_summary_with_stable_summary() surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -4817,6 +5000,7 @@ async fn model_port_preserves_gateway_safe_reason_kind() { surface_version: None, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await @@ -5997,15 +6181,35 @@ struct RecordingGateway { response: Result, } +struct MissingDiagnosticModelGateway; + +#[async_trait] +impl HostManagedModelGateway for MissingDiagnosticModelGateway { + async fn stream_model( + &self, + request: HostManagedModelRequest, + ) -> Result { + Ok( + HostManagedModelResponse::assistant_reply("model says hi".to_string()) + .with_effective_fallback_index(request.fallback_index), + ) + } +} + #[derive(Default)] struct RecordingPromptDiagnosticSink { captures: Mutex>, + model_calls: Mutex>, } impl HostManagedPromptDiagnosticSink for RecordingPromptDiagnosticSink { fn record_prompt(&self, capture: HostManagedPromptDiagnosticCapture) { self.captures.lock().expect("captures").push(capture); } + + fn record_model_call(&self, capture: HostManagedModelCallDiagnosticCapture) { + self.model_calls.lock().expect("model calls").push(capture); + } } impl RecordingGateway { @@ -6019,6 +6223,23 @@ impl RecordingGateway { } } + fn reply_with_usage_and_fallback( + content: &str, + usage: LoopModelUsage, + fallback_index: u32, + ) -> Self { + Self { + calls: Mutex::new(Vec::new()), + tool_definition_calls: Mutex::new(Vec::new()), + response: Ok( + HostManagedModelResponse::assistant_reply(content.to_string()) + .with_usage(usage) + .with_effective_fallback_index(fallback_index) + .with_diagnostic_effective_model("provider-model-from-response"), + ), + } + } + fn reply_with_fallback(content: &str, fallback_index: u32) -> Self { Self { calls: Mutex::new(Vec::new()), @@ -6070,6 +6291,20 @@ impl RecordingGateway { } } + fn model_error_with_usage( + kind: HostManagedModelErrorKind, + safe_summary: &str, + usage: LoopModelUsage, + ) -> Self { + Self { + calls: Mutex::new(Vec::new()), + tool_definition_calls: Mutex::new(Vec::new()), + response: Err(HostManagedModelError::safe(kind, safe_summary) + .with_usage(usage) + .with_diagnostic_effective_model("provider-model-from-error")), + } + } + fn model_error_with_reason_kind( kind: HostManagedModelErrorKind, safe_summary: &str, diff --git a/crates/loop/ironclaw_turn_runner/src/text_loop_driver.rs b/crates/loop/ironclaw_turn_runner/src/text_loop_driver.rs index 8e468489ec9..a7140dfa1e3 100644 --- a/crates/loop/ironclaw_turn_runner/src/text_loop_driver.rs +++ b/crates/loop/ironclaw_turn_runner/src/text_loop_driver.rs @@ -82,6 +82,7 @@ impl AgentLoopDriver for TextOnlyModelReplyDriver { surface_version: prompt_bundle.surface_version, model_preference: None, fallback_index: 0, + iteration: 0, capability_view: None, }) .await diff --git a/crates/product/ironclaw_assistant/src/inspector_store.rs b/crates/product/ironclaw_assistant/src/inspector_store.rs index 10bccde7232..8549a21e862 100644 --- a/crates/product/ironclaw_assistant/src/inspector_store.rs +++ b/crates/product/ironclaw_assistant/src/inspector_store.rs @@ -4,6 +4,7 @@ //! content out of durable events and drops all state at process restart. use std::collections::{HashMap, VecDeque}; +use std::hash::Hash; use std::sync::{Arc, OnceLock, RwLock}; use chrono::Utc; @@ -12,6 +13,7 @@ use ironclaw_host_api::{ turn::TurnRunId, }; use ironclaw_loop_host::{ + HostManagedModelCallDiagnosticCapture, HostManagedModelCallDiagnosticStatus, HostManagedModelMessageRole, HostManagedPromptDiagnosticCapture, HostManagedPromptDiagnosticMessage, HostManagedPromptDiagnosticSink, estimate_tokens_from_chars, @@ -20,15 +22,21 @@ use ironclaw_product_contracts::inspector::{ DEFAULT_MAX_ACTIVITY_ENTRIES, DEFAULT_MAX_LIVE_UPDATE_SCOPES, DEFAULT_MAX_MODEL_CALLS_PER_RUN, DEFAULT_MAX_RETAINED_RUNS_PER_SESSION, DEFAULT_MAX_RETAINED_UPDATES_PER_RUN, DEFAULT_MAX_TOOL_EXECUTIONS_PER_RUN, DEFAULT_MAX_TRACKED_SESSIONS, DiagnosticActivityEntry, - DiagnosticActivityEvent, DiagnosticCursor, DiagnosticScope, DiagnosticSequence, - DiagnosticSnapshot, DiagnosticStreamId, DiagnosticUpdateBatch, DiagnosticUpdateEnvelope, - DiagnosticUpdateKind, ModelCallDiagnostic, PromptComponentDiagnostic, PromptComponentKind, - PromptDiagnostic, SessionDiagnosticStats, ToolExecutionDiagnostic, + DiagnosticActivityEvent, DiagnosticCursor, DiagnosticMetricTotal, DiagnosticModelCallId, + DiagnosticModelCount, DiagnosticScope, DiagnosticSequence, DiagnosticSnapshot, + DiagnosticStreamId, DiagnosticUpdateBatch, DiagnosticUpdateEnvelope, DiagnosticUpdateKind, + InspectorModelCallStatus, MAX_MODELS_IN_STATS, ModelCallDiagnostic, ModelTokenUsage, + PromptComponentDiagnostic, PromptComponentKind, PromptDiagnostic, SessionDiagnosticStats, + ToolExecutionDiagnostic, ToolExecutionStatus, }; use ironclaw_safety::LeakDetector; use thiserror::Error; use tokio::sync::broadcast; +// Stats replacement identity must outlive the smaller operator-display queues. +// Bound it to the maximum retained diagnostic update horizon instead. +const MAX_STATS_RECORDS_PER_RUN: usize = DEFAULT_MAX_RETAINED_UPDATES_PER_RUN; + #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub struct DiagnosticStoreLimits { pub max_sessions: usize, @@ -163,6 +171,11 @@ struct DiagnosticRunState { prompt: Option, model_calls: VecDeque, tool_executions: VecDeque, + model_stats_records: HashMap, + model_stats_order: VecDeque, + tool_stats_records: + HashMap, + tool_stats_order: VecDeque, activity: VecDeque, stats: SessionDiagnosticStats, updates: VecDeque, @@ -176,6 +189,10 @@ impl Default for DiagnosticRunState { prompt: None, model_calls: VecDeque::new(), tool_executions: VecDeque::new(), + model_stats_records: HashMap::new(), + model_stats_order: VecDeque::new(), + tool_stats_records: HashMap::new(), + tool_stats_order: VecDeque::new(), activity: VecDeque::new(), stats: SessionDiagnosticStats::default(), updates: VecDeque::new(), @@ -361,13 +378,14 @@ impl InMemoryDiagnosticStore { prompt: PromptDiagnostic, ) -> Result { let prompt = prompt.into_bounded(); - let update = DiagnosticUpdateKind::PromptUpdated { - component_count: u32::try_from(prompt.components.len()).unwrap_or(u32::MAX), - total_estimated_tokens: prompt.total_estimated_tokens, - truncated: prompt.any_content_truncated(), - }; - self.record(scope, update, |run, _| { + self.record(scope, move |run, _| { + let update = DiagnosticUpdateKind::PromptUpdated { + component_count: u32::try_from(prompt.components.len()).unwrap_or(u32::MAX), + total_estimated_tokens: prompt.total_estimated_tokens, + truncated: prompt.any_content_truncated(), + }; run.prompt = Some(prompt); + update }) } @@ -377,15 +395,28 @@ impl InMemoryDiagnosticStore { model_call: ModelCallDiagnostic, ) -> Result { let model_call = model_call.into_bounded(); - let update = DiagnosticUpdateKind::ModelCall(model_call.clone()); let cap = self.limits.max_model_calls_per_run; - self.record(scope, update, move |run, _| { + self.record(scope, move |run, _| { + let contribution = ModelStatsContribution::from(&model_call); + let existing = replace_bounded_state( + &mut run.model_stats_records, + &mut run.model_stats_order, + model_call.call_id, + contribution.clone(), + MAX_STATS_RECORDS_PER_RUN, + ); + if let Some(existing) = existing.as_ref() { + update_model_stats(&mut run.stats, existing, false); + } + update_model_stats(&mut run.stats, &contribution, true); + let update = DiagnosticUpdateKind::ModelCall(model_call.clone()); replace_or_push( &mut run.model_calls, model_call, cap, |existing, incoming| existing.call_id == incoming.call_id, ); + update }) } @@ -395,20 +426,33 @@ impl InMemoryDiagnosticStore { tool: ToolExecutionDiagnostic, ) -> Result { let tool = tool.into_bounded(); - let update = DiagnosticUpdateKind::ToolExecutionUpdated { - activity_id: tool.activity_id, - model_call_id: tool.model_call_id, - capability_name: tool.capability_name.clone(), - status: tool.status, - duration_ms: tool.duration_ms, - output_bytes: tool.output_bytes, - result_truncated: tool.result_truncated(), - }; let cap = self.limits.max_tool_executions_per_run; - self.record(scope, update, move |run, _| { + self.record(scope, move |run, _| { + let contribution = ToolStatsContribution::from(&tool); + let existing = replace_bounded_state( + &mut run.tool_stats_records, + &mut run.tool_stats_order, + tool.activity_id, + contribution, + MAX_STATS_RECORDS_PER_RUN, + ); + if let Some(existing) = existing.as_ref() { + update_tool_stats(&mut run.stats, existing, false); + } + update_tool_stats(&mut run.stats, &contribution, true); + let update = DiagnosticUpdateKind::ToolExecutionUpdated { + activity_id: tool.activity_id, + model_call_id: tool.model_call_id, + capability_name: tool.capability_name.clone(), + status: tool.status, + duration_ms: tool.duration_ms, + output_bytes: tool.output_bytes, + result_truncated: tool.result_truncated(), + }; replace_or_push(&mut run.tool_executions, tool, cap, |existing, incoming| { existing.activity_id == incoming.activity_id }); + update }) } @@ -418,14 +462,15 @@ impl InMemoryDiagnosticStore { event: DiagnosticActivityEvent, ) -> Result { let event = event.into_bounded(); - let update = DiagnosticUpdateKind::Activity(event.clone()); let cap = self.limits.max_activity_entries_per_run; - self.record(scope, update, move |run, sequence| { + self.record(scope, move |run, sequence| { + let update = DiagnosticUpdateKind::Activity(event.clone()); push_bounded( &mut run.activity, DiagnosticActivityEntry { sequence, event }, cap, ); + update }) } @@ -435,11 +480,10 @@ impl InMemoryDiagnosticStore { stats: SessionDiagnosticStats, ) -> Result { let stats = stats.into_bounded(); - self.record( - scope, - DiagnosticUpdateKind::Stats(stats.clone()), - move |run, _| run.stats = stats, - ) + self.record(scope, move |run, _| { + run.stats = stats.clone(); + DiagnosticUpdateKind::Stats(stats) + }) } pub fn snapshot( @@ -557,8 +601,7 @@ impl InMemoryDiagnosticStore { fn record( &self, scope: DiagnosticScope, - update: DiagnosticUpdateKind, - mutate: impl FnOnce(&mut DiagnosticRunState, DiagnosticSequence), + mutate: impl FnOnce(&mut DiagnosticRunState, DiagnosticSequence) -> DiagnosticUpdateKind, ) -> Result { let mut state = self .state @@ -571,7 +614,7 @@ impl InMemoryDiagnosticStore { .checked_add(1) .ok_or(DiagnosticStoreError::SequenceExhausted)?; let sequence = DiagnosticSequence::new(next); - mutate(run, sequence); + let update = mutate(run, sequence); let cursor = DiagnosticCursor::new(run.stream_id, sequence); let envelope = DiagnosticUpdateEnvelope { scope, @@ -591,6 +634,146 @@ impl InMemoryDiagnosticStore { } } +fn update_counter(counter: &mut u64, value: u64, add: bool) { + *counter = if add { + counter.saturating_add(value) + } else { + counter.saturating_sub(value) + }; +} + +fn update_metric(metric: &mut DiagnosticMetricTotal, value: Option, add: bool) { + match value { + Some(value) => update_counter(&mut metric.known_total, value, add), + None => update_counter(&mut metric.unavailable_samples, 1, add), + } +} + +#[derive(Debug, Clone)] +struct ModelStatsContribution { + model: String, + usage: Option, + duration_ms: Option, +} + +impl From<&ModelCallDiagnostic> for ModelStatsContribution { + fn from(call: &ModelCallDiagnostic) -> Self { + Self { + model: call + .effective_model + .as_ref() + .unwrap_or(&call.requested_model) + .content() + .to_string(), + usage: call.usage.clone(), + duration_ms: call.duration_ms, + } + } +} + +#[derive(Debug, Clone, Copy)] +struct ToolStatsContribution { + status: ToolExecutionStatus, +} + +impl From<&ToolExecutionDiagnostic> for ToolStatsContribution { + fn from(tool: &ToolExecutionDiagnostic) -> Self { + Self { + status: tool.status, + } + } +} + +fn update_model_count(stats: &mut SessionDiagnosticStats, model: &str, add: bool) { + if let Some(index) = stats + .calls_per_model + .iter() + .position(|entry| entry.model.content() == model) + { + update_counter(&mut stats.calls_per_model[index].calls, 1, add); + if stats.calls_per_model[index].calls == 0 { + stats.calls_per_model.remove(index); + } + } else if add { + if stats.calls_per_model.len() < MAX_MODELS_IN_STATS { + stats + .calls_per_model + .push(DiagnosticModelCount::new(model, 1)); + } else { + // This is a monotonic completeness marker. Once a model bucket was + // omitted, later removals cannot reconstruct its historical count, + // so clearing the flag would incorrectly claim an exact breakdown. + stats.calls_per_model_truncated = true; + } + } +} + +fn update_model_stats( + stats: &mut SessionDiagnosticStats, + call: &ModelStatsContribution, + add: bool, +) { + update_counter(&mut stats.total_model_calls, 1, add); + update_model_count(stats, &call.model, add); + match call.usage.as_ref() { + Some(usage) => { + update_metric(&mut stats.input_tokens, usage.input_tokens, add); + update_metric(&mut stats.output_tokens, usage.output_tokens, add); + update_metric( + &mut stats.cache_read_input_tokens, + usage.cache_read_input_tokens, + add, + ); + update_metric( + &mut stats.cache_creation_input_tokens, + usage.cache_creation_input_tokens, + add, + ); + } + None => { + update_metric(&mut stats.input_tokens, None, add); + update_metric(&mut stats.output_tokens, None, add); + update_metric(&mut stats.cache_read_input_tokens, None, add); + update_metric(&mut stats.cache_creation_input_tokens, None, add); + } + } + update_metric(&mut stats.total_latency_ms, call.duration_ms, add); +} + +fn update_tool_stats(stats: &mut SessionDiagnosticStats, tool: &ToolStatsContribution, add: bool) { + update_counter(&mut stats.total_tool_calls, 1, add); + match tool.status { + ToolExecutionStatus::Succeeded => { + update_counter(&mut stats.successful_tool_calls, 1, add); + } + ToolExecutionStatus::Failed => { + update_counter(&mut stats.failed_tool_calls, 1, add); + } + ToolExecutionStatus::Started => {} + } +} + +fn replace_bounded_state( + values: &mut HashMap, + order: &mut VecDeque, + key: K, + value: V, + capacity: usize, +) -> Option +where + K: Clone + Eq + Hash, +{ + let previous = values.insert(key.clone(), value); + touch(order, key); + while values.len() > capacity { + let Some(evicted) = order.pop_front() else { + break; + }; + values.remove(&evicted); + } + previous +} + impl DiagnosticStorePort for InMemoryDiagnosticStore { fn record_activity( &self, @@ -774,6 +957,58 @@ impl HostManagedPromptDiagnosticSink for InMemoryDiagnosticStore { tracing::debug!(%error, "prompt diagnostics could not be retained"); } } + + fn record_model_call(&self, capture: HostManagedModelCallDiagnosticCapture) { + let user_id = capture + .context + .actor + .as_ref() + .map(|actor| actor.user_id.clone()) + .or_else(|| capture.context.scope.explicit_owner_user_id().cloned()); + let Some(user_id) = user_id else { + tracing::debug!( + run_id = %capture.context.run_id, + "model-call diagnostics skipped because the run has no user scope" + ); + return; + }; + let scope = DiagnosticScope::new( + capture.context.scope.tenant_id.clone(), + user_id, + capture.context.thread_id.clone(), + capture.context.run_id, + ); + let usage = capture.usage.map(|usage| ModelTokenUsage { + input_tokens: Some(u64::from(usage.input_tokens)), + output_tokens: Some(u64::from(usage.output_tokens)), + cache_read_input_tokens: Some(u64::from(usage.cache_read_input_tokens)), + cache_creation_input_tokens: Some(u64::from(usage.cache_creation_input_tokens)), + }); + let status = match capture.status { + HostManagedModelCallDiagnosticStatus::Succeeded => InspectorModelCallStatus::Succeeded, + HostManagedModelCallDiagnosticStatus::Failed => InspectorModelCallStatus::Failed, + }; + let detector = LeakDetector::new(); + let model_call = ModelCallDiagnostic::new( + ironclaw_product_contracts::inspector::DiagnosticModelCallId::new(), + capture.iteration, + diagnostic_prompt_text(&detector, &capture.requested_model), + capture + .effective_model + .map(|model| diagnostic_prompt_text(&detector, &model)), + capture.started_at, + Some(capture.completed_at), + Some(capture.duration_ms), + status, + usage, + capture + .failure_summary + .map(|summary| diagnostic_prompt_text(&detector, &summary)), + ); + if let Err(error) = InMemoryDiagnosticStore::record_model_call(self, scope, model_call) { + tracing::debug!(%error, "model-call diagnostics could not be retained"); + } + } } impl Default for InMemoryDiagnosticStore { @@ -833,7 +1068,9 @@ mod tests { ids::{CapabilityId, TenantId, ThreadId, UserId}, turn::{RunProfileId, RunProfileVersion, TurnActor, TurnId, TurnRunId, TurnScope}, }; - use ironclaw_loop_contracts::{LoopRunContext, ModelProfileId, ResolvedRunProfile, SkillName}; + use ironclaw_loop_contracts::{ + LoopModelUsage, LoopRunContext, ModelProfileId, ResolvedRunProfile, SkillName, + }; use ironclaw_loop_host::ProviderModelId; use ironclaw_product_contracts::inspector::{ BoundedDiagnosticText, DIAGNOSTIC_LABEL_MAX_BYTES, DIAGNOSTIC_SUMMARY_MAX_BYTES, @@ -1220,6 +1457,7 @@ mod tests { }, ) .expect("activity"); + store .record_stats( scope.clone(), @@ -1298,6 +1536,380 @@ mod tests { ); } + #[test] + fn metrics_cover_provider_failover_replacements_missing_usage_and_retention() { + let store = InMemoryDiagnosticStore::new(tiny_limits()).expect("store"); + let scope = scope("tenant", "user", "thread", TurnRunId::new()); + let first_id = DiagnosticModelCallId::new(); + let first_started = ModelCallDiagnostic::new( + first_id, + 2, + "requested-a", + Some("model-a".to_string()), + Utc::now(), + None, + None, + InspectorModelCallStatus::Started, + None, + None, + ); + store + .record_model_call(scope.clone(), first_started) + .expect("record start"); + let first_completed = ModelCallDiagnostic::new( + first_id, + 2, + "requested-a", + Some("model-a".to_string()), + Utc::now(), + Some(Utc::now()), + Some(100), + InspectorModelCallStatus::Succeeded, + Some(ModelTokenUsage { + input_tokens: Some(15), + output_tokens: Some(5), + cache_read_input_tokens: Some(3), + cache_creation_input_tokens: Some(2), + }), + None, + ); + store + .record_model_call(scope.clone(), first_completed) + .expect("complete first"); + store + .record_model_call( + scope.clone(), + ModelCallDiagnostic::new( + DiagnosticModelCallId::new(), + 3, + "requested-a", + Some("model-b".to_string()), + Utc::now(), + Some(Utc::now()), + Some(300), + InspectorModelCallStatus::Failed, + Some(ModelTokenUsage { + input_tokens: Some(7), + output_tokens: Some(1), + cache_read_input_tokens: Some(0), + cache_creation_input_tokens: Some(0), + }), + Some("provider failed".to_string()), + ), + ) + .expect("record failed call with usage"); + store + .record_model_call( + scope.clone(), + ModelCallDiagnostic::new( + DiagnosticModelCallId::new(), + 4, + "requested-a", + Some("model-a".to_string()), + Utc::now(), + Some(Utc::now()), + None, + InspectorModelCallStatus::Succeeded, + None, + None, + ), + ) + .expect("record call without usage"); + + let tool_id = ironclaw_host_api::turn::CapabilityActivityId::new(); + for status in [ToolExecutionStatus::Started, ToolExecutionStatus::Succeeded] { + store + .record_tool_execution( + scope.clone(), + ToolExecutionDiagnostic::new( + tool_id, + None, + "filesystem.read", + None, + None, + status, + None, + None, + None, + None, + ), + ) + .expect("record tool transition"); + } + store + .record_tool_execution( + scope.clone(), + ToolExecutionDiagnostic::new( + ironclaw_host_api::turn::CapabilityActivityId::new(), + None, + "network.fetch", + None, + None, + ToolExecutionStatus::Failed, + None, + None, + None, + None, + ), + ) + .expect("record failed tool"); + + let snapshot = store.snapshot(&scope).expect("snapshot").expect("present"); + assert_eq!(snapshot.model_calls.len(), 2, "run detail stays bounded"); + assert_eq!(snapshot.model_calls[0].iteration, 3); + assert_eq!(snapshot.model_calls[1].iteration, 4); + assert_eq!(snapshot.stats.total_model_calls, 3); + assert_eq!(snapshot.stats.input_tokens.known_total, 22); + assert_eq!(snapshot.stats.input_tokens.unavailable_samples, 1); + assert_eq!(snapshot.stats.output_tokens.known_total, 6); + assert_eq!(snapshot.stats.total_latency_ms.known_total, 400); + assert_eq!(snapshot.stats.total_latency_ms.unavailable_samples, 1); + assert_eq!(snapshot.stats.total_tool_calls, 2); + assert_eq!(snapshot.stats.successful_tool_calls, 1); + assert_eq!(snapshot.stats.failed_tool_calls, 1); + assert_eq!(snapshot.stats.calls_per_model.len(), 2); + assert_eq!(snapshot.stats.calls_per_model[0].calls, 2); + assert_eq!(snapshot.stats.calls_per_model[1].calls, 1); + } + + #[test] + fn terminal_updates_do_not_recount_records_evicted_from_display_retention() { + let store = InMemoryDiagnosticStore::new(tiny_limits()).expect("store"); + let scope = scope("tenant", "user", "thread", TurnRunId::new()); + let first_model_id = DiagnosticModelCallId::new(); + + for (call_id, iteration, status) in [ + (first_model_id, 1, InspectorModelCallStatus::Started), + ( + DiagnosticModelCallId::new(), + 2, + InspectorModelCallStatus::Succeeded, + ), + ( + DiagnosticModelCallId::new(), + 3, + InspectorModelCallStatus::Failed, + ), + ] { + store + .record_model_call( + scope.clone(), + ModelCallDiagnostic::new( + call_id, + iteration, + "requested", + Some("provider-model".to_string()), + Utc::now(), + None, + None, + status, + None, + None, + ), + ) + .expect("record model call"); + } + + let first_tool_id = ironclaw_host_api::turn::CapabilityActivityId::new(); + for (activity_id, status) in [ + (first_tool_id, ToolExecutionStatus::Started), + ( + ironclaw_host_api::turn::CapabilityActivityId::new(), + ToolExecutionStatus::Succeeded, + ), + ( + ironclaw_host_api::turn::CapabilityActivityId::new(), + ToolExecutionStatus::Failed, + ), + ] { + store + .record_tool_execution( + scope.clone(), + ToolExecutionDiagnostic::new( + activity_id, + None, + "filesystem.read", + None, + None, + status, + None, + None, + None, + None, + ), + ) + .expect("record tool execution"); + } + + let retained = store.snapshot(&scope).expect("snapshot").expect("run"); + assert_eq!(retained.model_calls.len(), 2); + assert!( + retained + .model_calls + .iter() + .all(|call| call.call_id != first_model_id) + ); + assert_eq!(retained.tool_executions.len(), 2); + assert!( + retained + .tool_executions + .iter() + .all(|tool| tool.activity_id != first_tool_id) + ); + + store + .record_model_call( + scope.clone(), + ModelCallDiagnostic::new( + first_model_id, + 1, + "requested", + Some("provider-model".to_string()), + Utc::now(), + Some(Utc::now()), + Some(50), + InspectorModelCallStatus::Succeeded, + Some(ModelTokenUsage { + input_tokens: Some(11), + output_tokens: Some(4), + cache_read_input_tokens: Some(0), + cache_creation_input_tokens: Some(0), + }), + None, + ), + ) + .expect("complete evicted model call"); + store + .record_tool_execution( + scope.clone(), + ToolExecutionDiagnostic::new( + first_tool_id, + None, + "filesystem.read", + None, + None, + ToolExecutionStatus::Succeeded, + None, + None, + None, + None, + ), + ) + .expect("complete evicted tool execution"); + + let snapshot = store.snapshot(&scope).expect("snapshot").expect("run"); + assert_eq!(snapshot.stats.total_model_calls, 3); + assert_eq!(snapshot.stats.input_tokens.known_total, 11); + assert_eq!(snapshot.stats.input_tokens.unavailable_samples, 2); + assert_eq!(snapshot.stats.total_tool_calls, 3); + assert_eq!(snapshot.stats.successful_tool_calls, 2); + assert_eq!(snapshot.stats.failed_tool_calls, 1); + } + + #[test] + fn model_breakdown_truncation_remains_latched_after_a_bucket_is_removed() { + let store = InMemoryDiagnosticStore::default(); + let scope = scope("tenant", "user", "thread", TurnRunId::new()); + let first_call_id = DiagnosticModelCallId::new(); + for index in 0..=MAX_MODELS_IN_STATS { + let call_id = if index == 0 { + first_call_id + } else { + DiagnosticModelCallId::new() + }; + store + .record_model_call( + scope.clone(), + ModelCallDiagnostic::new( + call_id, + u32::try_from(index).expect("model index"), + "requested", + Some(format!("model-{index}")), + Utc::now(), + Some(Utc::now()), + Some(1), + InspectorModelCallStatus::Succeeded, + None, + None, + ), + ) + .expect("record distinct model"); + } + store + .record_model_call( + scope.clone(), + ModelCallDiagnostic::new( + first_call_id, + 0, + "requested", + Some("model-1".to_string()), + Utc::now(), + Some(Utc::now()), + Some(1), + InspectorModelCallStatus::Succeeded, + None, + None, + ), + ) + .expect("replace first model contribution"); + + let snapshot = store.snapshot(&scope).expect("snapshot").expect("run"); + assert!(snapshot.stats.calls_per_model_truncated); + assert_eq!( + snapshot.stats.total_model_calls, + u64::try_from(MAX_MODELS_IN_STATS + 1).expect("model count") + ); + assert_eq!( + snapshot.stats.calls_per_model.len(), + MAX_MODELS_IN_STATS - 1 + ); + } + + #[test] + fn metric_aggregation_saturates_instead_of_wrapping() { + let store = InMemoryDiagnosticStore::new(tiny_limits()).expect("store"); + let scope = scope("tenant", "user", "thread", TurnRunId::new()); + store + .record_stats( + scope.clone(), + SessionDiagnosticStats { + total_model_calls: u64::MAX, + input_tokens: DiagnosticMetricTotal { + known_total: u64::MAX, + unavailable_samples: u64::MAX, + }, + ..SessionDiagnosticStats::default() + }, + ) + .expect("seed saturated stats"); + store + .record_model_call( + scope.clone(), + ModelCallDiagnostic::new( + DiagnosticModelCallId::new(), + 0, + "requested", + Some("model".to_string()), + Utc::now(), + Some(Utc::now()), + Some(1), + InspectorModelCallStatus::Succeeded, + None, + None, + ), + ) + .expect("aggregate at saturation"); + + let stats = store + .snapshot(&scope) + .expect("snapshot") + .expect("present") + .stats; + assert_eq!(stats.total_model_calls, u64::MAX); + assert_eq!(stats.input_tokens.known_total, u64::MAX); + assert_eq!(stats.input_tokens.unavailable_samples, u64::MAX); + } + #[test] fn host_prompt_capture_is_scoped_redacted_validated_and_bounded() { let store = InMemoryDiagnosticStore::default(); @@ -1321,7 +1933,7 @@ mod tests { HostManagedPromptDiagnosticSink::record_prompt( &store, HostManagedPromptDiagnosticCapture { - context, + context: context.clone(), messages: vec![ HostManagedPromptDiagnosticMessage { role: HostManagedModelMessageRole::System, @@ -1375,9 +1987,30 @@ mod tests { context_limit: 128_000, }, ); + HostManagedPromptDiagnosticSink::record_model_call( + &store, + HostManagedModelCallDiagnosticCapture { + context, + iteration: 4, + requested_model: "interactive_model".to_string(), + effective_model: None, + started_at: Utc::now(), + completed_at: Utc::now(), + duration_ms: 42, + status: HostManagedModelCallDiagnosticStatus::Succeeded, + usage: Some(LoopModelUsage { + input_tokens: 12, + output_tokens: 4, + cache_read_input_tokens: 2, + cache_creation_input_tokens: 1, + }), + failure_summary: None, + }, + ); + let diagnostic_scope = DiagnosticScope::new(tenant_id, user_id, thread_id, run_id); let prompt = store - .prompt(&DiagnosticScope::new(tenant_id, user_id, thread_id, run_id)) + .prompt(&diagnostic_scope) .expect("prompt read") .expect("prompt captured"); assert_eq!(prompt.message_count, 4); @@ -1412,6 +2045,15 @@ mod tests { prompt.components.last().map(|component| component.kind), Some(PromptComponentKind::Capability) ); + let snapshot = store + .snapshot(&diagnostic_scope) + .expect("snapshot") + .expect("present"); + assert_eq!(snapshot.model_calls.len(), 1); + assert_eq!(snapshot.model_calls[0].iteration, 4); + assert_eq!(snapshot.model_calls[0].effective_model, None); + assert_eq!(snapshot.stats.total_model_calls, 1); + assert_eq!(snapshot.stats.input_tokens.known_total, 12); } #[test] diff --git a/crates/product/ironclaw_assistant/tests/support/planned_agent_loop.rs b/crates/product/ironclaw_assistant/tests/support/planned_agent_loop.rs index 7da1a1ec074..1e146aa27bf 100644 --- a/crates/product/ironclaw_assistant/tests/support/planned_agent_loop.rs +++ b/crates/product/ironclaw_assistant/tests/support/planned_agent_loop.rs @@ -225,6 +225,7 @@ pub fn capability_call_response( safe_reasoning_deltas: Vec::new(), usage: None, effective_fallback_index: Some(0), + diagnostic_effective_model: None, output: ParentLoopOutput::CapabilityCalls(vec![CapabilityCallCandidate { activity_id: ironclaw_turns::CapabilityActivityId::new(), surface_version: harness_surface_version(), @@ -683,6 +684,7 @@ impl ScriptedHostRuntimeToolCall { safe_reasoning_deltas: Vec::new(), usage: None, effective_fallback_index: Some(request.fallback_index), + diagnostic_effective_model: None, output: ParentLoopOutput::CapabilityCalls(vec![CapabilityCallCandidate { activity_id: ironclaw_turns::CapabilityActivityId::new(), surface_version, diff --git a/crates/product/ironclaw_webui/frontend/src/pages/chat/inspector/inspector-panel.test.tsx b/crates/product/ironclaw_webui/frontend/src/pages/chat/inspector/inspector-panel.test.tsx index d986b0bb405..d94cd75b8d0 100644 --- a/crates/product/ironclaw_webui/frontend/src/pages/chat/inspector/inspector-panel.test.tsx +++ b/crates/product/ironclaw_webui/frontend/src/pages/chat/inspector/inspector-panel.test.tsx @@ -120,6 +120,42 @@ test.each(truncationCases)("prompt tab reports a truncated %s", async (_label, t ); }); +test("stats tab formats aggregates and unavailable samples without zero fabrication", async () => { + inspectorState.snapshot = { + stats: { + total_model_calls: 3, + calls_per_model: [ + { + model: { content: "provider-model", original_bytes: 14, truncated: false }, + calls: 3, + }, + ], + calls_per_model_truncated: false, + input_tokens: { known_total: 1_200, unavailable_samples: 1 }, + output_tokens: { known_total: 80, unavailable_samples: 1 }, + cache_read_input_tokens: { known_total: 0, unavailable_samples: 3 }, + cache_creation_input_tokens: { known_total: 20, unavailable_samples: 1 }, + total_latency_ms: { known_total: 900, unavailable_samples: 1 }, + }, + }; + + await act(async () => + root?.render(), + ); + await act(async () => + document.querySelector("[data-testid='inspector-tab-stats']")?.click(), + ); + + const stats = document.querySelector("[data-testid='inspector-stats-content']"); + assert.ok(stats); + assert.match(stats.textContent || "", /1,200/); + assert.match(stats.textContent || "", /450 ms/); + assert.match(stats.textContent || "", /Unavailable/); + assert.match(stats.textContent || "", /provider-model3/); + assert.doesNotMatch(stats.textContent || "", /Tool calls|Tool outcomes/); + assert.match(stats.textContent || "", /7 metric samples were unavailable/); +}); + afterEach(async () => { await act(async () => root?.unmount()); document.body.replaceChildren(); diff --git a/crates/product/ironclaw_webui/frontend/src/pages/chat/inspector/inspector-panel.tsx b/crates/product/ironclaw_webui/frontend/src/pages/chat/inspector/inspector-panel.tsx index ba2e346b0cd..b72f2e29b35 100644 --- a/crates/product/ironclaw_webui/frontend/src/pages/chat/inspector/inspector-panel.tsx +++ b/crates/product/ironclaw_webui/frontend/src/pages/chat/inspector/inspector-panel.tsx @@ -239,27 +239,96 @@ function ActivityShell({ ); } +interface DiagnosticMetricTotal { + known_total: number; + unavailable_samples: number; +} + +interface SessionDiagnosticStats { + total_model_calls: number; + calls_per_model: Array<{ model: BoundedDiagnosticText; calls: number }>; + calls_per_model_truncated: boolean; + input_tokens: DiagnosticMetricTotal; + output_tokens: DiagnosticMetricTotal; + cache_read_input_tokens: DiagnosticMetricTotal; + cache_creation_input_tokens: DiagnosticMetricTotal; + total_latency_ms: DiagnosticMetricTotal; +} + +function metricValue( + metric: DiagnosticMetricTotal, + sampleCount: number, + suffix = "", +): string { + if (sampleCount > 0 && metric.unavailable_samples >= sampleCount) return "Unavailable"; + return `${metric.known_total.toLocaleString()}${suffix}`; +} + +function MetricCard({ label, value }: { label: string; value: string }) { + return ( +

+

{label}

+

{value}

+
+ ); +} + function StatsShell({ snapshot }: { snapshot: Record | null }) { - const stats = snapshot?.stats as { total_model_calls?: unknown } | undefined; + const stats = snapshot?.stats as SessionDiagnosticStats | undefined; if (!stats) { return ( ); } - const totalCalls = typeof stats.total_model_calls === "number" ? stats.total_model_calls : 0; + const knownLatencySamples = Math.max( + 0, + stats.total_model_calls - stats.total_latency_ms.unavailable_samples, + ); + const averageLatency = knownLatencySamples > 0 + ? `${Math.round(stats.total_latency_ms.known_total / knownLatencySamples).toLocaleString()} ms` + : "Unavailable"; + const partialMetricCount = [ + stats.input_tokens, + stats.output_tokens, + stats.cache_read_input_tokens, + stats.cache_creation_input_tokens, + stats.total_latency_ms, + ].reduce((total, metric) => total + metric.unavailable_samples, 0); return ( -
-
-

Model calls

-

{totalCalls}

+
+
+ + + + + + +
-
-

Live updates

-

Available

+
+

Calls per model

+ {stats.calls_per_model.length === 0 ? ( +

No model breakdown available.

+ ) : ( +
+ {stats.calls_per_model.map((entry, index) => ( +
+
{entry.model.content}
+
{entry.calls.toLocaleString()}
+
+ ))} +
+ )}
+ {(partialMetricCount > 0 || stats.calls_per_model_truncated) && ( +

+ Statistics are partial: {partialMetricCount.toLocaleString()} metric samples were unavailable + {stats.calls_per_model_truncated ? " and the model breakdown was truncated" : ""}. +

+ )}
); } diff --git a/crates/product/ironclaw_webui/frontend/src/pages/chat/inspector/useInspector.test.tsx b/crates/product/ironclaw_webui/frontend/src/pages/chat/inspector/useInspector.test.tsx index 113401cc984..ce1ac8e9576 100644 --- a/crates/product/ironclaw_webui/frontend/src/pages/chat/inspector/useInspector.test.tsx +++ b/crates/product/ironclaw_webui/frontend/src/pages/chat/inspector/useInspector.test.tsx @@ -163,38 +163,85 @@ test("bounds transient snapshot retries", async () => { } }); -test("deduplicates cursors, rebases snapshots, and stops on forbidden", async () => { - await act(async () => root?.render()); - const stream = eventStreams[0]; - await act(async () => stream.respond()); +test("deduplicates cursors, bounds refresh bursts, rebases, and stops on forbidden", async () => { + vi.useFakeTimers(); + try { + await act(async () => root?.render()); + const stream = eventStreams[0]; + await act(async () => stream.respond()); - const streamId = "550e8400-e29b-41d4-a716-446655440000"; - await act(async () => { - stream.message("diagnostic_update", `${streamId}:1`, { sequence: 1 }); - stream.message("diagnostic_update", `${streamId}:1`, { sequence: 1 }); - stream.message("diagnostic_update", `${streamId}:2`, { sequence: 2 }); - }); - assert.equal(latestState?.updates.length, 2); - assert.equal(latestState?.lastCursor, `${streamId}:2`); + const streamId = "550e8400-e29b-41d4-a716-446655440000"; + await act(async () => { + stream.message("diagnostic_update", `${streamId}:1`, { sequence: 1 }); + stream.message("diagnostic_update", `${streamId}:1`, { sequence: 1 }); + stream.message("diagnostic_update", `${streamId}:2`, { sequence: 2 }); + }); + assert.equal(latestState?.updates.length, 2); + assert.equal(latestState?.lastCursor, `${streamId}:2`); - await act(async () => { - stream.message("diagnostic_update", `${streamId}:3`, { - update: { type: "prompt_updated" }, + await act(async () => { + stream.message("diagnostic_update", `${streamId}:3`, { + update: { type: "prompt_updated" }, + }); + stream.message("diagnostic_update", `${streamId}:4`, { + update: { type: "model_call", data: {} }, + }); + stream.message("diagnostic_update", `${streamId}:5`, { + update: { type: "tool_execution_updated" }, + }); + stream.message("diagnostic_update", `${streamId}:6`, { + update: { type: "stats" }, + }); + }); + assert.equal(vi.mocked(fetch).mock.calls.length, 1); + await act(async () => vi.advanceTimersByTimeAsync(49)); + assert.equal(vi.mocked(fetch).mock.calls.length, 1); + await act(async () => vi.advanceTimersByTimeAsync(1)); + assert.equal(vi.mocked(fetch).mock.calls.length, 2); + assert.deepEqual(stream.options.headers(), { + Authorization: "Bearer operator-token", + "Last-Event-ID": `${streamId}:6`, }); - }); - assert.equal(vi.mocked(fetch).mock.calls.length, 2); - await act(async () => { - stream.message("diagnostic_rebase", `${streamId}:4`, { - latest_cursor: { stream_id: streamId, sequence: 4 }, + for (let sequence = 7; sequence <= 12; sequence += 1) { + await act(async () => { + stream.message("diagnostic_update", `${streamId}:${sequence}`, { + update: { type: "stats" }, + }); + await vi.advanceTimersByTimeAsync(40); + }); + } + await act(async () => { + stream.message("diagnostic_update", `${streamId}:13`, { + update: { type: "stats" }, + }); + }); + assert.equal(vi.mocked(fetch).mock.calls.length, 2); + await act(async () => vi.advanceTimersByTimeAsync(9)); + assert.equal(vi.mocked(fetch).mock.calls.length, 2); + await act(async () => vi.advanceTimersByTimeAsync(1)); + assert.equal(vi.mocked(fetch).mock.calls.length, 3); + assert.deepEqual(stream.options.headers(), { + Authorization: "Bearer operator-token", + "Last-Event-ID": `${streamId}:13`, }); - }); - assert.equal(latestState?.updates.length, 0); - assert.equal(vi.mocked(fetch).mock.calls.length, 3); - await act(async () => stream.respond(403, "application/json")); - assert.equal(latestState?.health, INSPECTOR_HEALTH.FORBIDDEN); - assert.equal(stream.controller.abort.mock.calls.length, 1); + await act(async () => { + stream.message("diagnostic_rebase", `${streamId}:14`, { + latest_cursor: { stream_id: streamId, sequence: 14 }, + }); + }); + assert.equal(latestState?.updates.length, 0); + assert.equal(vi.mocked(fetch).mock.calls.length, 3); + await act(async () => vi.advanceTimersByTimeAsync(50)); + assert.equal(vi.mocked(fetch).mock.calls.length, 4); + + await act(async () => stream.respond(403, "application/json")); + assert.equal(latestState?.health, INSPECTOR_HEALTH.FORBIDDEN); + assert.equal(stream.controller.abort.mock.calls.length, 1); + } finally { + vi.useRealTimers(); + } }); test("hidden tabs release the stream and reconnect when visible", async () => { diff --git a/crates/product/ironclaw_webui/frontend/src/pages/chat/inspector/useInspector.ts b/crates/product/ironclaw_webui/frontend/src/pages/chat/inspector/useInspector.ts index 4a10a073047..8510ae723ab 100644 --- a/crates/product/ironclaw_webui/frontend/src/pages/chat/inspector/useInspector.ts +++ b/crates/product/ironclaw_webui/frontend/src/pages/chat/inspector/useInspector.ts @@ -13,6 +13,8 @@ const MAX_RETRY_INTERVAL_MS = 30_000; const MAX_RETAINED_UPDATES = 1_024; const MAX_SNAPSHOT_ATTEMPTS = 3; const SNAPSHOT_RETRY_BASE_DELAY_MS = 500; +const SNAPSHOT_REFRESH_DEBOUNCE_MS = 50; +const SNAPSHOT_REFRESH_MAX_WAIT_MS = 250; interface DiagnosticUpdate { stream_id?: string; @@ -123,22 +125,53 @@ export function useInspector({ let disposed = false; let connectedOnce = false; let terminalState = false; + let snapshotRefreshTimer: number | null = null; + let snapshotRefreshMaxWaitTimer: number | null = null; let controller: ReturnType | null = null; const request = inspectorEventStreamRequest({ threadId, runId }); const stream = new EventSourcePlus(request.url, { credentials: "same-origin", - headers: request.headers, + headers: () => ({ + ...request.headers(), + ...(lastCursorRef.current ? { "Last-Event-ID": lastCursorRef.current } : {}), + }), maxRetryInterval: MAX_RETRY_INTERVAL_MS, retryStrategy: "always", }); function terminal(healthState: InspectorHealth, message: string): void { terminalState = true; + cancelSnapshotRefresh(); setHealth(healthState); setError(message); controller?.abort("terminal inspector response"); } + function cancelSnapshotRefresh(): void { + if (snapshotRefreshTimer !== null) window.clearTimeout(snapshotRefreshTimer); + if (snapshotRefreshMaxWaitTimer !== null) { + window.clearTimeout(snapshotRefreshMaxWaitTimer); + } + snapshotRefreshTimer = null; + snapshotRefreshMaxWaitTimer = null; + } + + function refreshSnapshot(): void { + cancelSnapshotRefresh(); + if (!disposed) setSnapshotGeneration((generation) => generation + 1); + } + + function scheduleSnapshotRefresh(): void { + if (snapshotRefreshMaxWaitTimer === null) { + snapshotRefreshMaxWaitTimer = window.setTimeout( + refreshSnapshot, + SNAPSHOT_REFRESH_MAX_WAIT_MS, + ); + } + if (snapshotRefreshTimer !== null) window.clearTimeout(snapshotRefreshTimer); + snapshotRefreshTimer = window.setTimeout(refreshSnapshot, SNAPSHOT_REFRESH_DEBOUNCE_MS); + } + function connect(): void { if (disposed) return; setHealth(connectedOnce ? INSPECTOR_HEALTH.RECONNECTING : INSPECTOR_HEALTH.CONNECTING); @@ -190,7 +223,7 @@ export function useInspector({ if (message.event === "diagnostic_rebase") { if (cursor) lastCursorRef.current = cursor; setUpdates([]); - setSnapshotGeneration((generation) => generation + 1); + scheduleSnapshotRefresh(); return; } if (message.event !== "diagnostic_update") return; @@ -198,8 +231,13 @@ export function useInspector({ lastCursorRef.current = cursor; setUpdates((current) => [...current, payload as DiagnosticUpdate].slice(-MAX_RETAINED_UPDATES)); const update = payload.update as { type?: unknown } | undefined; - if (update?.type === "prompt_updated") { - setSnapshotGeneration((generation) => generation + 1); + if ( + update?.type === "prompt_updated" + || update?.type === "model_call" + || update?.type === "tool_execution_updated" + || update?.type === "stats" + ) { + scheduleSnapshotRefresh(); } setHealth(INSPECTOR_HEALTH.CONNECTED); }, @@ -225,6 +263,7 @@ export function useInspector({ document.addEventListener("visibilitychange", onVisibilityChange); return () => { disposed = true; + cancelSnapshotRefresh(); document.removeEventListener("visibilitychange", onVisibilityChange); controller?.abort("inspector disposed"); }; diff --git a/tests/CLAUDE.md b/tests/CLAUDE.md index 1b896f6b065..3e1ccbe001f 100644 --- a/tests/CLAUDE.md +++ b/tests/CLAUDE.md @@ -320,7 +320,7 @@ entries. | Delete a thread behind a shared confirmation dialog | `test_reborn_webui_v2_smoke.py::test_reborn_v2_thread_delete_uses_shared_confirmation_dialog` | | Collapse the sidebar, pick a theme, pick a language — and have it persist | `test_reborn_webui_v2_smoke.py` | | Opt into the inspector, preserve its selected tab while closing, resizing, and reloading, and leave the ordinary chat shell unchanged when debug mode is off | `test_reborn_webui_v2_smoke.py::test_inspector_debug_activation_and_responsive_shell` | -| Inspect the bounded host-resolved prompt for a completed run | `test_reborn_webui_v2_smoke.py::test_inspector_prompt_renders_host_resolved_diagnostics` | +| Inspect the bounded host-resolved prompt and model-call statistics for a completed run | `test_reborn_webui_v2_smoke.py::test_inspector_prompt_and_stats_render_host_diagnostics` | | Reconnect SSE without gaps or duplicates; multiple tabs both get the reply; excess connections are rate-limited | `test_reborn_webui_v2_legacy_sse_history.py`, `test_reborn_webui_v2_streaming_run_control_api.py` (10) | | Keep execution-only engine threads out of the chat sidebar while preserving deep-linked history | `test_v2_thread_visibility.py` (2; pending legacy migration #6369) | diff --git a/tests/e2e/helpers.py b/tests/e2e/helpers.py index 73e6a58cd57..0212e779d0a 100644 --- a/tests/e2e/helpers.py +++ b/tests/e2e/helpers.py @@ -324,6 +324,7 @@ async def dismiss_native_dialog(dialog) -> None: "inspector_panel": "[data-testid='inspector-panel']", "inspector_prompt_content": "[data-testid='inspector-prompt-content']", "inspector_tab_stats": "[data-testid='inspector-tab-stats']", + "inspector_stats_content": "[data-testid='inspector-stats-content']", "inspector_close": "[data-testid='inspector-close']", "inspector_open": "[data-testid='inspector-open']", "msg_user": "[data-testid='msg-user']", # user message bubble diff --git a/tests/e2e/scenarios/test_reborn_webui_v2_smoke.py b/tests/e2e/scenarios/test_reborn_webui_v2_smoke.py index 2412cef3be0..41190e64175 100644 --- a/tests/e2e/scenarios/test_reborn_webui_v2_smoke.py +++ b/tests/e2e/scenarios/test_reborn_webui_v2_smoke.py @@ -431,11 +431,11 @@ async def test_inspector_debug_activation_and_responsive_shell( await context.close() -async def test_inspector_prompt_renders_host_resolved_diagnostics( +async def test_inspector_prompt_and_stats_render_host_diagnostics( reborn_v2_server, reborn_v2_browser, ): - """A real model turn reaches the bounded operator-only Prompt tab.""" + """A real model turn reaches the bounded operator-only Prompt and Stats tabs.""" marker = f"prompt-inspector-e2e-{uuid.uuid4()}" async with httpx.AsyncClient(headers=reborn_bearer_headers()) as client: thread_id = await _create_thread(client, reborn_v2_server) @@ -473,6 +473,29 @@ async def test_inspector_prompt_renders_host_resolved_diagnostics( "Reconstructed content reflects the latest host prompt boundary", ) ).to_have_count(1) + + await page.locator(SEL_V2["inspector_tab_stats"]).click() + stats = page.locator(SEL_V2["inspector_stats_content"]) + await expect(stats).to_be_visible() + model_calls = ( + stats.get_by_text("Model calls", exact=True) + .locator("..") + .locator("p") + .nth(1) + ) + await expect(model_calls).to_have_text("1") + input_tokens = ( + stats.get_by_text("Input tokens", exact=True) + .locator("..") + .locator("p") + .nth(1) + ) + await expect(input_tokens).to_have_text("10") + await expect(stats.get_by_text("Output tokens", exact=True).locator("..")).not_to_contain_text( + "Unavailable" + ) + await expect(stats.get_by_text("mock-model", exact=True)).to_be_visible() + await expect(stats.get_by_text("Statistics are partial:")).to_have_count(0) finally: await context.close() diff --git a/tests/support/reborn_parity_qa/binary_e2e.rs b/tests/support/reborn_parity_qa/binary_e2e.rs index 158a3ca624f..56562a4eed1 100644 --- a/tests/support/reborn_parity_qa/binary_e2e.rs +++ b/tests/support/reborn_parity_qa/binary_e2e.rs @@ -1627,6 +1627,7 @@ pub fn trace_tool_call_response() -> ironclaw_loop_host::HostManagedModelRespons safe_reasoning_deltas: Vec::new(), usage: None, effective_fallback_index: Some(0), + diagnostic_effective_model: None, output: ParentLoopOutput::CapabilityCalls(vec![CapabilityCallCandidate { activity_id: ironclaw_turns::CapabilityActivityId::new(), surface_version: CapabilitySurfaceVersion::new(TEST_CAPABILITY_SURFACE_VERSION)