diff --git a/Cargo.toml b/Cargo.toml index a990a756902..629201a7915 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -506,6 +506,10 @@ path = "tests/integration/webui_v2_router_smoke.rs" name = "reborn_integration_wiring_parity" path = "tests/integration/wiring_parity.rs" +[[test]] +name = "reborn_integration_run_artifact_timings" +path = "tests/integration/run_artifact_timings.rs" + [[test]] name = "reborn_integration_db_write_canonical" path = "tests/integration/db_write_canonical.rs" diff --git a/crates/product/ironclaw_assistant/AGENTS.md b/crates/product/ironclaw_assistant/AGENTS.md index d04173d3bcf..5a04537aaf4 100644 --- a/crates/product/ironclaw_assistant/AGENTS.md +++ b/crates/product/ironclaw_assistant/AGENTS.md @@ -338,7 +338,7 @@ by file. | `commands` | The product command grammar surface: listing declared commands, executing one, and the capability handler dispatch table | A command's *effect* — each handler delegates to the concern that owns it | `mod product_capability_handlers`, `mod product_commands`, `PRODUCT_COMMAND_LIST_COMMAND_ID`, `PRODUCT_COMMAND_LIST_COMMAND`, `PRODUCT_COMMAND_EXECUTE_COMMAND_ID`, `PRODUCT_COMMAND_EXECUTE_COMMAND`, `command_result_field`, `model_command_view`, `user_model_preference_command_view`, `idle_status_command_view`, `nothing_to_stop_command_view`, `new_conversation_started_view`, `reborn_services/types.rs::RebornProductCommandEffect`, `product_command_input`, `reborn_services/product_capability_handlers.rs::ProductCommandHandler`, `reborn_services/product_capability_handlers.rs::command_output`, `reborn_services/product_capability_handlers.rs::ProductCapabilityHandler`, `reborn_services/product_commands.rs::lifecycle_command_title`, `reborn_services/product_commands.rs::lifecycle_command_view`, `reborn_services/product_commands.rs::lifecycle_rows_view`, `reborn_services/product_commands.rs::lifecycle_confirmation_view`, `reborn_services/product_commands.rs::package_ref_fields`, `reborn_services/product_commands.rs::package_kind_label`, `reborn_services/product_commands.rs::blocker_line`, `reborn_services/product_commands.rs::capability_lines`, `reborn_services/product_commands.rs::extension_row`, `reborn_services/product_commands.rs::skill_row`, `reborn_services/product_commands.rs::yes_no`, `reborn_services/types.rs::RebornExecuteProductCommandResponse` | | `admin-users` | Admin user CRUD, role/status, per-user secrets, and last-admin protection | Authorization of the caller — that is `dispatch`'s view gate | `mod admin_users`, `ADMIN_USER_UPDATE_CAPABILITY_ID`, `ADMIN_USER_UPDATE_CAPABILITY`, `ADMIN_USER_SET_STATUS_CAPABILITY_ID`, `ADMIN_USER_SET_STATUS_CAPABILITY`, `ADMIN_USER_SET_ROLE_CAPABILITY_ID`, `ADMIN_USER_SET_ROLE_CAPABILITY`, `ADMIN_USER_DELETE_CAPABILITY_ID`, `ADMIN_USER_DELETE_CAPABILITY`, `ADMIN_USER_PUT_SECRET_CAPABILITY_ID`, `ADMIN_USER_PUT_SECRET_CAPABILITY`, `ADMIN_USER_DELETE_SECRET_CAPABILITY_ID`, `ADMIN_USER_DELETE_SECRET_CAPABILITY`, `ADMIN_USER_CREATE_COMMAND_ID`, `ADMIN_USER_CREATE_COMMAND`, `ADMIN_USER_DELETE_SECRET_COMMAND_ID`, `ADMIN_USER_DELETE_SECRET_COMMAND`, `ADMIN_USERS_VIEW`, `ADMIN_USER_VIEW`, `ADMIN_USER_SECRETS_VIEW`, `map_admin_user_error`, `last_admin_error`, `reborn_services/admin_users.rs::RejectingAdminUserService` | | `dispatch` | The `RebornServices` object itself, the caller/scope plumbing every concern shares (`ProductAgentBoundCaller`, resource scope, secret handles), the `ProductSurfaceError` shaping helpers, and the view-authorization gate | A concern's own logic — an item belongs here only when **more than one** sub-owner calls it | `mod views`, `ProductAgentBoundCaller`, `caller_resource_scope`, `product_view_forbidden`, `authorize_product_view`, `ProductCapabilityInvoker`, `UnavailableProductCapabilityInvoker`, `RebornServices`, `product_capability_input_error`, `product_secret_handle`, `segment`, `map_adapter_error`, `code_for_status`, `kind_for_surface_rejection`, `truncate_utf8_to_bytes`, `product_agent_bound_caller_from_webui`, `reborn_services/views.rs::EmptyViewParams`, `reborn_services/views.rs::parse_empty_view_params`, `reborn_services/views.rs::required_string_view_param`, `reborn_services/views.rs::view_page`, `reborn_services/views.rs::view_page_with_cursor`, `reborn_services/views.rs::UnavailableRebornViewProvider` | -| `run-artifact` | Run and thread artifact export views | Live run state — that is `runs` | `mod run_artifact`, `mod thread_artifact`, `ADMIN_THREAD_SCRAPE_THREADS_VIEW`, `ADMIN_THREAD_SCRAPE_ARTIFACT_VIEW`, `ADMIN_THREAD_SCRAPE_RUN_ARTIFACT_VIEW`, `reborn_services/run_artifact.rs::RUN_ARTIFACT_SCHEMA`, `reborn_services/run_artifact.rs::RUN_ARTIFACT_VIEW`, `reborn_services/run_artifact.rs::ARTIFACT_REDACTION_PIPELINE`, `reborn_services/run_artifact.rs::RebornRunArtifact`, `reborn_services/run_artifact.rs::RunArtifactMessage`, `reborn_services/run_artifact.rs::RunArtifactToolCall`, `reborn_services/run_artifact.rs::RunArtifactLogs`, `reborn_services/run_artifact.rs::RunArtifactRedaction`, `reborn_services/run_artifact.rs::context_messages_by_id`, `reborn_services/run_artifact.rs::artifact_messages`, `reborn_services/run_artifact.rs::redact_text`, `reborn_services/run_artifact.rs::redact_json`, `reborn_services/run_artifact.rs::redact_json_strings`, `reborn_services/thread_artifact.rs::THREAD_ARTIFACT_SCHEMA`, `reborn_services/thread_artifact.rs::THREAD_ARTIFACT_MAX_MESSAGES`, `reborn_services/thread_artifact.rs::THREAD_ARTIFACT_MAX_STORED_BYTES`, `reborn_services/thread_artifact.rs::THREAD_ARTIFACT_MAX_SERIALIZED_BYTES`, `reborn_services/thread_artifact.rs::THREAD_ARTIFACT_VIEW`, `reborn_services/thread_artifact.rs::RebornThreadArtifact`, `reborn_services/thread_artifact.rs::thread_artifact_too_large` | +| `run-artifact` | Run and thread artifact export views | Live run state — that is `runs` | `mod run_artifact`, `mod thread_artifact`, `mod timings_source`, `ADMIN_THREAD_SCRAPE_THREADS_VIEW`, `ADMIN_THREAD_SCRAPE_ARTIFACT_VIEW`, `ADMIN_THREAD_SCRAPE_RUN_ARTIFACT_VIEW`, `reborn_services/run_artifact.rs::RUN_ARTIFACT_SCHEMA`, `reborn_services/run_artifact.rs::RUN_ARTIFACT_VIEW`, `reborn_services/run_artifact.rs::ARTIFACT_REDACTION_PIPELINE`, `reborn_services/run_artifact.rs::RebornRunArtifact`, `reborn_services/run_artifact.rs::RunArtifactMessage`, `reborn_services/run_artifact.rs::RunArtifactToolCall`, `reborn_services/run_artifact.rs::RunArtifactLogs`, `reborn_services/run_artifact.rs::RunArtifactRedaction`, `reborn_services/run_artifact.rs::mod timings`, `reborn_services/run_artifact.rs::context_messages_by_id`, `reborn_services/run_artifact.rs::artifact_messages`, `reborn_services/run_artifact.rs::redact_text`, `reborn_services/run_artifact.rs::redact_json`, `reborn_services/run_artifact.rs::redact_json_strings`, `reborn_services/thread_artifact.rs::THREAD_ARTIFACT_SCHEMA`, `reborn_services/thread_artifact.rs::THREAD_ARTIFACT_MAX_MESSAGES`, `reborn_services/thread_artifact.rs::THREAD_ARTIFACT_MAX_STORED_BYTES`, `reborn_services/thread_artifact.rs::THREAD_ARTIFACT_MAX_SERIALIZED_BYTES`, `reborn_services/thread_artifact.rs::THREAD_ARTIFACT_VIEW`, `reborn_services/thread_artifact.rs::RebornThreadArtifact`, `reborn_services/thread_artifact.rs::RunArtifactRunTimings`, `reborn_services/thread_artifact.rs::group_messages_by_run`, `reborn_services/thread_artifact.rs::thread_artifact_too_large`, `reborn_services/timings_source.rs::derive_wall_clock_ms` | | `skills` | Skill install/update/remove, search and content reads, auto-activation settings, and the activation recorder/clearer ports | Skill *selection* — that is `ironclaw_skills` | `SkillActivationRecorder`, `SkillActivationClearer`, `SKILL_INSTALL_CAPABILITY_ID`, `SKILL_INSTALL_CAPABILITY`, `SKILL_UPDATE_CAPABILITY_ID`, `SKILL_UPDATE_CAPABILITY`, `SKILL_REMOVE_CAPABILITY_ID`, `SKILL_REMOVE_CAPABILITY`, `SKILL_AUTO_ACTIVATE_SET_CAPABILITY_ID`, `SKILL_AUTO_ACTIVATE_SET_CAPABILITY`, `SKILL_AUTO_ACTIVATE_LEARNED_SET_CAPABILITY_ID`, `SKILL_AUTO_ACTIVATE_LEARNED_SET_CAPABILITY`, `SKILLS_VIEW`, `SKILL_SEARCH_VIEW`, `SKILL_CONTENT_VIEW`, `SkillsProductService`, `UnsupportedSkillsProductService` | | `gates` | Approval and auth gate routing: which resolver a `gate_ref` reaches, stale/attacker-supplied ref rejection, and the blocked-gate notices | The approval or auth *policy* — that lives in `approval_interaction` / `auth_interaction`, not in this module | `mod approval_settings`, `RESOLVE_GATE_COMMAND_ID`, `RESOLVE_GATE_COMMAND`, `NOTICE_BLOCKED_APPROVAL`, `NOTICE_BLOCKED_AUTH`, `GateResolutionRoute`, `validate_current_gate_ref`, `participant_denied`, `assert_generic_run_parked_on_gate`, `reject_generic_auth_gate_resolution`, `persistent_approval_unavailable`, `blocked_approval_unavailable`, `blocked_authentication_unavailable`, `map_auth_interaction_error`, `reborn_services/types.rs::reborn_resume_gate_response` | | `traces` | Trace credits, trace holds, and the trace account login link | Trace *content* — that is `ironclaw_trace_commons` | `mod trace_credits`, `TRACE_ACCOUNT_LOGIN_LINK_COMMAND_ID`, `TRACE_ACCOUNT_LOGIN_LINK_COMMAND`, `TRACE_HOLD_AUTHORIZE_COMMAND_ID`, `TRACE_HOLD_AUTHORIZE_COMMAND`, `reborn_services/trace_credits.rs::TRACE_CREDITS_VIEW`, `reborn_services/trace_credits.rs::TRACE_ACCOUNT_TRACES_VIEW`, `reborn_services/trace_credits.rs::TRACE_CREDITS_NOTE`, `reborn_services/trace_credits.rs::AccountLoginLinkMintError`, `reborn_services/trace_credits.rs::account_login_link_for_user`, `reborn_services/trace_credits.rs::AccountTracesError`, `reborn_services/trace_credits.rs::account_traces_for_user`, `reborn_services/trace_credits.rs::authorize_trace_hold_for_user`, `reborn_services/trace_credits.rs::local_trace_credits_for_user` | diff --git a/crates/product/ironclaw_assistant/src/inspector_store.rs b/crates/product/ironclaw_assistant/src/inspector_store.rs index 682b9e11c04..e02417c55b5 100644 --- a/crates/product/ironclaw_assistant/src/inspector_store.rs +++ b/crates/product/ironclaw_assistant/src/inspector_store.rs @@ -349,6 +349,77 @@ pub struct InMemoryDiagnosticStore { state: RwLock, } +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct DiagnosticTimingModelCall { + pub call_id: DiagnosticModelCallId, + pub iteration: u32, + pub requested_model: BoundedDiagnosticText, + pub effective_model: Option, + pub started_at: chrono::DateTime, + pub completed_at: Option>, + pub duration_ms: Option, + pub status: InspectorModelCallStatus, +} + +impl From<&ModelCallDiagnostic> for DiagnosticTimingModelCall { + fn from(call: &ModelCallDiagnostic) -> Self { + Self { + call_id: call.call_id, + iteration: call.iteration, + requested_model: call.requested_model.clone(), + effective_model: call.effective_model.clone(), + started_at: call.started_at, + completed_at: call.completed_at, + duration_ms: call.duration_ms, + status: call.status, + } + } +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct DiagnosticTimingToolExecution { + pub model_call_id: Option, + pub capability_name: BoundedDiagnosticText, + pub status: ToolExecutionStatus, + pub duration_ms: Option, +} + +impl From<&ToolExecutionDiagnostic> for DiagnosticTimingToolExecution { + fn from(tool: &ToolExecutionDiagnostic) -> Self { + Self { + model_call_id: tool.model_call_id, + capability_name: tool.capability_name.clone(), + status: tool.status, + duration_ms: tool.duration_ms, + } + } +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct DiagnosticTimingSnapshot { + pub model_calls: Vec, + pub tool_executions: Vec, + pub stats: SessionDiagnosticStats, +} + +impl From for DiagnosticTimingSnapshot { + fn from(snapshot: DiagnosticSnapshot) -> Self { + Self { + model_calls: snapshot + .model_calls + .iter() + .map(DiagnosticTimingModelCall::from) + .collect(), + tool_executions: snapshot + .tool_executions + .iter() + .map(DiagnosticTimingToolExecution::from) + .collect(), + stats: snapshot.stats, + } + } +} + /// Operator inspection store surface exposed by product composition. /// /// Capture remains behind [`HostManagedPromptDiagnosticSink`]; consumers of @@ -365,6 +436,17 @@ pub trait DiagnosticStorePort: Send + Sync { scope: &DiagnosticScope, ) -> Result, DiagnosticStoreError>; + /// Read only the data needed by timing projections. Implementations with + /// larger snapshots should override this to avoid cloning prompt and + /// activity payloads under the store lock. + fn timing_snapshot( + &self, + scope: &DiagnosticScope, + ) -> Result, DiagnosticStoreError> { + self.snapshot(scope) + .map(|snapshot| snapshot.map(DiagnosticTimingSnapshot::from)) + } + fn prompt( &self, scope: &DiagnosticScope, @@ -592,6 +674,32 @@ impl InMemoryDiagnosticStore { })) } + pub fn timing_snapshot( + &self, + scope: &DiagnosticScope, + ) -> Result, DiagnosticStoreError> { + let state = self + .state + .read() + .map_err(|_| DiagnosticStoreError::StateUnavailable)?; + let Some(run) = state.run(scope) else { + return Ok(None); + }; + Ok(Some(DiagnosticTimingSnapshot { + model_calls: run + .model_calls + .iter() + .map(DiagnosticTimingModelCall::from) + .collect(), + tool_executions: run + .tool_executions + .iter() + .map(DiagnosticTimingToolExecution::from) + .collect(), + stats: run.stats.clone(), + })) + } + pub fn prompt( &self, scope: &DiagnosticScope, @@ -884,6 +992,13 @@ impl DiagnosticStorePort for InMemoryDiagnosticStore { InMemoryDiagnosticStore::snapshot(self, scope) } + fn timing_snapshot( + &self, + scope: &DiagnosticScope, + ) -> Result, DiagnosticStoreError> { + InMemoryDiagnosticStore::timing_snapshot(self, scope) + } + fn prompt( &self, scope: &DiagnosticScope, @@ -2531,6 +2646,44 @@ mod tests { ); } + #[test] + fn timing_snapshot_omits_prompt_and_activity_payloads() { + let store = InMemoryDiagnosticStore::default(); + let scope = scope("tenant", "user", "thread", TurnRunId::new()); + let call_id = DiagnosticModelCallId::new(); + store + .record_model_call( + scope.clone(), + ModelCallDiagnostic::new( + call_id, + 1, + "requested-model", + None, + Utc::now(), + None, + None, + InspectorModelCallStatus::Started, + None, + None, + ), + ) + .expect("model call"); + store + .record_activity(scope.clone(), activity("not part of timings")) + .expect("activity"); + + let full = store.snapshot(&scope).expect("full snapshot").expect("run"); + assert_eq!(full.activity.len(), 1); + + let timings = store + .timing_snapshot(&scope) + .expect("timing snapshot") + .expect("run"); + assert_eq!(timings.model_calls.len(), 1); + assert!(timings.tool_executions.is_empty()); + assert_eq!(timings.stats.total_model_calls, 1); + } + #[test] fn session_eviction_is_deterministic_and_write_lru() { let store = InMemoryDiagnosticStore::new(tiny_limits()).expect("store"); diff --git a/crates/product/ironclaw_assistant/src/lib.rs b/crates/product/ironclaw_assistant/src/lib.rs index cd84fc2947a..0ef5baef788 100644 --- a/crates/product/ironclaw_assistant/src/lib.rs +++ b/crates/product/ironclaw_assistant/src/lib.rs @@ -252,6 +252,9 @@ pub use run_delivery::{ // Adapter, projection, and event DTOs are re-exported from // `ironclaw_host_api::product_adapter` above so product terminals consume a // single product service. +pub use reborn_services::run_artifact::timings::{ + RunArtifactIterationTiming, RunArtifactTimingTotals, RunArtifactTimings, RunArtifactToolTiming, +}; pub use reborn_services::{ ADMIN_CONFIGURATION_REPLACE_CAPABILITY, ADMIN_CONFIGURATION_REPLACE_CAPABILITY_ID, ADMIN_CONFIGURATION_VIEW, ADMIN_THREAD_SCRAPE_ARTIFACT_VIEW, @@ -370,7 +373,7 @@ pub use reborn_services::{ RebornTimelineResponse, RebornTraceCreditsResponse, RebornTraceHoldAuthorizeProductRequest, RebornTraceHoldAuthorizeResponse, RebornUpdateMemberRoleRequest, RebornUpdateProjectRequest, RebornVendorAuthAccounts, RegistrationChannelNotificationSetupService, RunArtifactLogs, - RunArtifactMessage, RunArtifactRedaction, RunArtifactToolCall, + RunArtifactMessage, RunArtifactRedaction, RunArtifactRunTimings, RunArtifactToolCall, SKILL_AUTO_ACTIVATE_LEARNED_SET_CAPABILITY, SKILL_AUTO_ACTIVATE_LEARNED_SET_CAPABILITY_ID, SKILL_AUTO_ACTIVATE_SET_CAPABILITY, SKILL_AUTO_ACTIVATE_SET_CAPABILITY_ID, SKILL_CONTENT_VIEW, SKILL_INSTALL_CAPABILITY, SKILL_INSTALL_CAPABILITY_ID, SKILL_REMOVE_CAPABILITY, diff --git a/crates/product/ironclaw_assistant/src/reborn_services.rs b/crates/product/ironclaw_assistant/src/reborn_services.rs index 7a03811d38a..0708ea142ac 100644 --- a/crates/product/ironclaw_assistant/src/reborn_services.rs +++ b/crates/product/ironclaw_assistant/src/reborn_services.rs @@ -148,8 +148,12 @@ mod product_capability_handlers; mod product_commands; mod project_fs; mod projects; -mod run_artifact; +// pub(crate): lib.rs re-exports `reborn_services::run_artifact::timings` +// directly, which needs this segment of the path visible crate-wide; the +// module's own contents stay unexported except through that re-export. +pub(crate) mod run_artifact; mod thread_artifact; +mod timings_source; mod trace_credits; mod types; mod views; @@ -303,8 +307,8 @@ pub use run_artifact::{ RunArtifactLogs, RunArtifactMessage, RunArtifactRedaction, RunArtifactToolCall, }; pub use thread_artifact::{ - RebornThreadArtifact, RebornThreadArtifactRequest, THREAD_ARTIFACT_MAX_MESSAGES, - THREAD_ARTIFACT_SCHEMA, THREAD_ARTIFACT_VIEW, + RebornThreadArtifact, RebornThreadArtifactRequest, RunArtifactRunTimings, + THREAD_ARTIFACT_MAX_MESSAGES, THREAD_ARTIFACT_SCHEMA, THREAD_ARTIFACT_VIEW, }; pub use types::{ RebornAuthAccount, RebornCreateThreadResponse, RebornExecuteProductCommandResponse, diff --git a/crates/product/ironclaw_assistant/src/reborn_services/run_artifact.rs b/crates/product/ironclaw_assistant/src/reborn_services/run_artifact.rs index 4f66cf006ed..0e437038cfd 100644 --- a/crates/product/ironclaw_assistant/src/reborn_services/run_artifact.rs +++ b/crates/product/ironclaw_assistant/src/reborn_services/run_artifact.rs @@ -31,6 +31,10 @@ use ironclaw_product_contracts::surface::{ pub use ironclaw_product_contracts::product_wire::RebornRunArtifactRequest; +pub mod timings; + +use timings::RunArtifactTimings; + pub const RUN_ARTIFACT_SCHEMA: &str = "ironclaw.run_artifact.v1"; pub const RUN_ARTIFACT_VIEW: RebornViewDescriptor = RebornViewDescriptor { id: "run_artifact", @@ -46,6 +50,10 @@ pub struct RebornRunArtifact { pub run: RebornGetRunStateResponse, pub messages: Vec, pub logs: RunArtifactLogs, + /// Best-effort per-iteration and per-tool timing. `#[serde(default)]` so a + /// `v1` artifact downloaded before timings existed still deserializes. + #[serde(default)] + pub timings: RunArtifactTimings, pub redaction: RunArtifactRedaction, } @@ -55,6 +63,15 @@ pub struct RunArtifactMessage { pub sequence: u64, #[serde(default, skip_serializing_if = "Option::is_none")] pub run_id: Option, + /// When the message row was first persisted. `None` for records written + /// before per-message timestamps existed. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub created_at: Option>, + /// Last time the message row materially changed. For an assistant reply + /// this is the finalization time, which is what makes step-to-step gaps + /// meaningful. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub updated_at: Option>, pub kind: MessageKind, pub status: MessageStatus, pub content: String, @@ -140,6 +157,14 @@ where let (messages, message_redaction_applied) = self .artifact_messages_for_records(thread_scope, &thread_id, run_records, &redactor) .await?; + // Borrow `caller` here, before `artifact_logs` takes it by value below. + let timings = self.artifact_timings( + &caller, + &thread_id, + &run_id, + run.received_at, + messages.iter(), + ); let (logs, log_redaction_applied) = self .artifact_logs(caller, &thread_id, Some(&run_id), &redactor) .await; @@ -151,6 +176,7 @@ where run, messages, logs, + timings, redaction: RunArtifactRedaction { pipeline: ARTIFACT_REDACTION_PIPELINE.to_string(), applied: message_redaction_applied || log_redaction_applied, @@ -301,6 +327,8 @@ pub(super) fn artifact_messages( message_id: record.message_id.to_string(), sequence: record.sequence, run_id: record.turn_run_id.clone(), + created_at: record.created_at, + updated_at: record.updated_at, kind: record.kind, status: record.status, content, @@ -425,6 +453,32 @@ mod tests { assert!(serialized.contains("[REDACTED]")); } + #[test] + fn assembler_exports_durable_message_timestamps() { + let thread_id = ThreadId::new("thread-a").expect("thread id"); + let created = Utc::now(); + let updated = created + chrono::Duration::seconds(41); + let mut source = record( + &thread_id, + ThreadMessageId::new(), + 1, + MessageKind::Assistant, + "hello", + ); + source.created_at = Some(created); + source.updated_at = Some(updated); + + let (messages, _) = artifact_messages( + vec![source], + &HashMap::new(), + &DeterministicTraceRedactor::new(Vec::new()), + ); + + assert_eq!(messages.len(), 1); + assert_eq!(messages[0].created_at, Some(created)); + assert_eq!(messages[0].updated_at, Some(updated)); + } + fn record( thread_id: &ThreadId, message_id: ThreadMessageId, diff --git a/crates/product/ironclaw_assistant/src/reborn_services/run_artifact/timings.rs b/crates/product/ironclaw_assistant/src/reborn_services/run_artifact/timings.rs new file mode 100644 index 00000000000..7f8b40d52c7 --- /dev/null +++ b/crates/product/ironclaw_assistant/src/reborn_services/run_artifact/timings.rs @@ -0,0 +1,495 @@ +//! Timing-only projection of one run's process-local diagnostic snapshot. +//! +//! This module deliberately carries NO `BoundedDiagnosticText` payloads +//! across into the artifact: no prompt text, no tool arguments, no tool +//! results. Capability names, statuses, counts, and durations only. The +//! artifact is a user-downloadable file, and the redaction pipeline that +//! guards `messages` does not run over this block. + +use std::collections::HashMap; + +use crate::inspector_store::DiagnosticTimingSnapshot; +use chrono::{DateTime, Utc}; +use ironclaw_product_contracts::inspector::{ + DiagnosticMetricTotal, DiagnosticModelCallId, InspectorModelCallStatus, ToolExecutionStatus, +}; +use serde::{Deserialize, Serialize}; + +/// Names the capture source in the exported file so a reader knows which +/// buffer produced (or failed to produce) these numbers. +pub const TIMINGS_SOURCE: &str = "diagnostic_store"; + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct RunArtifactTimings { + pub source: String, + pub available: bool, + /// ponytail: always false. Timings come from the process-local, bounded + /// `InMemoryDiagnosticStore` (`crates/product/ironclaw_assistant/src/inspector_store.rs:3` + /// — "deliberately has no persistence backend"), capped by + /// `DiagnosticStoreLimits` (same file, :49). A restart or an eviction + /// removes a run's timings with no durable marker, exactly like the + /// sibling `RunArtifactLogs.complete`. Ceiling: timings are unavailable + /// for any run that left the buffer. Upgrade path: a durable diagnostics + /// store, deliberately deferred — see this plan's Global Constraints. + pub complete: bool, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub unavailable_reason: Option, + pub iterations: Vec, + /// Tool executions the store never correlated to a model call, or whose + /// model call left the buffer first. Counted, never dropped. + #[serde(default)] + pub unattributed_tools: Vec, + pub totals: RunArtifactTimingTotals, +} + +impl Default for RunArtifactTimings { + fn default() -> Self { + unavailable("timings_absent") + } +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct RunArtifactIterationTiming { + /// Agent-loop iteration number this model call served. + pub iteration: u32, + /// Effective model when the provider resolved one, else the requested model. + pub model: String, + pub started_at: DateTime, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub completed_at: Option>, + /// Wall-clock time inside the provider call. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub inference_ms: Option, + pub status: InspectorModelCallStatus, + /// Tool calls this iteration requested — the "how many tools ran before + /// the assistant replied" number, per iteration. + pub tool_calls: u32, + /// Sum of the tool durations below. `None` when no tool reported one. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub tool_ms_total: Option, + pub tools: Vec, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct RunArtifactToolTiming { + pub capability_name: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub duration_ms: Option, + pub status: ToolExecutionStatus, +} + +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] +pub struct RunArtifactTimingTotals { + pub iterations: u64, + pub tool_calls: u64, + pub failed_tool_calls: u64, + /// Summed provider latency. `unavailable_samples` counts calls that + /// reported no duration, so a reader can tell "fast" from "unmeasured". + pub inference_ms: DiagnosticMetricTotal, + /// Sum of durations for tool executions still retained by the bounded + /// diagnostic store. This is never presented as a complete run total. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub retained_tool_ms: Option, + /// False when the store's cumulative tool count exceeds its retained + /// execution records, so `retained_tool_ms` is known to be partial. + #[serde(default, skip_serializing_if = "is_false")] + pub retained_tool_ms_complete: bool, + /// run.received_at → the newest message `updated_at`. Approximate: it + /// includes queue and persistence latency around the run, and it is the + /// only end-to-end number available without widening into turn state. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub wall_clock_ms: Option, +} + +fn is_false(value: &bool) -> bool { + !*value +} + +/// The "no numbers, and here is why" value. Never an error — a missing +/// timing block must not fail an artifact export. +pub fn unavailable(reason: &str) -> RunArtifactTimings { + RunArtifactTimings { + source: TIMINGS_SOURCE.to_string(), + available: false, + complete: false, + unavailable_reason: Some(reason.to_string()), + iterations: Vec::new(), + unattributed_tools: Vec::new(), + totals: RunArtifactTimingTotals::default(), + } +} + +pub fn project_timings( + snapshot: DiagnosticTimingSnapshot, + wall_clock_ms: Option, +) -> RunArtifactTimings { + let mut tools_by_call: HashMap> = + HashMap::new(); + let mut unattributed_tools = Vec::new(); + let known_calls: std::collections::HashSet = snapshot + .model_calls + .iter() + .map(|call| call.call_id) + .collect(); + + let mut retained_tool_ms = 0_u64; + let mut retained_tool_ms_seen = false; + for execution in &snapshot.tool_executions { + let timing = RunArtifactToolTiming { + capability_name: execution.capability_name.content().to_string(), + duration_ms: execution.duration_ms, + status: execution.status, + }; + if let Some(duration) = execution.duration_ms { + retained_tool_ms = retained_tool_ms.saturating_add(duration); + retained_tool_ms_seen = true; + } + match execution + .model_call_id + .filter(|call_id| known_calls.contains(call_id)) + { + Some(call_id) => tools_by_call.entry(call_id).or_default().push(timing), + None => unattributed_tools.push(timing), + } + } + + let mut iterations: Vec = snapshot + .model_calls + .iter() + .map(|call| { + let tools = tools_by_call.remove(&call.call_id).unwrap_or_default(); + let tool_ms_total = sum_durations(&tools); + RunArtifactIterationTiming { + iteration: call.iteration, + model: call + .effective_model + .as_ref() + .unwrap_or(&call.requested_model) + .content() + .to_string(), + started_at: call.started_at, + completed_at: call.completed_at, + inference_ms: call.duration_ms, + status: call.status, + tool_calls: u32::try_from(tools.len()).unwrap_or(u32::MAX), + tool_ms_total, + tools, + } + }) + .collect(); + // Capture order follows completion, not iteration order; sort so the file + // reads as the loop ran. + iterations.sort_by_key(|iteration| iteration.iteration); + + RunArtifactTimings { + source: TIMINGS_SOURCE.to_string(), + available: true, + complete: false, + unavailable_reason: None, + iterations, + unattributed_tools, + totals: RunArtifactTimingTotals { + iterations: snapshot.stats.total_model_calls, + tool_calls: snapshot.stats.total_tool_calls, + failed_tool_calls: snapshot.stats.failed_tool_calls, + inference_ms: snapshot.stats.total_latency_ms, + retained_tool_ms: retained_tool_ms_seen.then_some(retained_tool_ms), + retained_tool_ms_complete: u64::try_from(snapshot.tool_executions.len()) + .unwrap_or(u64::MAX) + >= snapshot.stats.total_tool_calls, + wall_clock_ms, + }, + } +} + +fn sum_durations(tools: &[RunArtifactToolTiming]) -> Option { + let mut total = 0_u64; + let mut seen = false; + for tool in tools { + if let Some(duration) = tool.duration_ms { + total = total.saturating_add(duration); + seen = true; + } + } + seen.then_some(total) +} + +#[cfg(test)] +mod tests { + use super::*; + use ironclaw_host_api::turn::CapabilityActivityId; + use ironclaw_product_contracts::inspector::{ + DiagnosticMetricTotal, InspectorModelCallStatus, ModelCallDiagnostic, + SessionDiagnosticStats, ToolExecutionDiagnostic, ToolExecutionStatus, + }; + + fn model_call( + call_id: DiagnosticModelCallId, + iteration: u32, + duration_ms: u64, + ) -> ModelCallDiagnostic { + let started = Utc::now(); + ModelCallDiagnostic::new( + call_id, + iteration, + "claude-opus-5", + Some("claude-opus-5-20260101".to_string()), + started, + Some(started + chrono::Duration::milliseconds(duration_ms as i64)), + Some(duration_ms), + InspectorModelCallStatus::Succeeded, + None, + None, + ) + } + + fn tool( + model_call_id: Option, + name: &str, + duration_ms: u64, + status: ToolExecutionStatus, + ) -> ToolExecutionDiagnostic { + ToolExecutionDiagnostic::new( + CapabilityActivityId::new(), + model_call_id, + name, + None, + None, + status, + Some(duration_ms), + None, + None, + None, + ) + } + + fn snapshot( + model_calls: Vec, + tool_executions: Vec, + stats: SessionDiagnosticStats, + ) -> DiagnosticTimingSnapshot { + DiagnosticTimingSnapshot { + model_calls: model_calls + .iter() + .map(crate::inspector_store::DiagnosticTimingModelCall::from) + .collect(), + tool_executions: tool_executions + .iter() + .map(crate::inspector_store::DiagnosticTimingToolExecution::from) + .collect(), + stats, + } + } + + #[test] + fn tools_are_grouped_under_the_iteration_that_requested_them() { + let first = DiagnosticModelCallId::new(); + let second = DiagnosticModelCallId::new(); + let projected = project_timings( + snapshot( + vec![model_call(first, 1, 4_100), model_call(second, 2, 900)], + vec![ + tool(Some(first), "shell", 22_000, ToolExecutionStatus::Succeeded), + tool(Some(first), "read_file", 30, ToolExecutionStatus::Succeeded), + tool(Some(second), "shell", 11, ToolExecutionStatus::Failed), + ], + SessionDiagnosticStats::default(), + ), + Some(91_000), + ); + + assert!(projected.available); + assert!( + !projected.complete, + "process-local capture is never complete" + ); + assert_eq!(projected.iterations.len(), 2); + assert_eq!(projected.iterations[0].iteration, 1); + assert_eq!(projected.iterations[0].inference_ms, Some(4_100)); + assert_eq!(projected.iterations[0].tool_calls, 2); + assert_eq!(projected.iterations[0].tool_ms_total, Some(22_030)); + assert_eq!(projected.iterations[1].tool_calls, 1); + assert!(projected.unattributed_tools.is_empty()); + assert_eq!(projected.totals.wall_clock_ms, Some(91_000)); + } + + #[test] + fn iterations_are_ordered_by_iteration_number_not_capture_order() { + let first = DiagnosticModelCallId::new(); + let second = DiagnosticModelCallId::new(); + let projected = project_timings( + snapshot( + vec![model_call(second, 7, 10), model_call(first, 2, 10)], + Vec::new(), + SessionDiagnosticStats::default(), + ), + None, + ); + + let seen: Vec = projected + .iterations + .iter() + .map(|iteration| iteration.iteration) + .collect(); + assert_eq!(seen, vec![2, 7]); + } + + #[test] + fn uncorrelated_tools_land_in_the_unattributed_bucket_and_still_count() { + let call = DiagnosticModelCallId::new(); + let projected = project_timings( + snapshot( + vec![model_call(call, 1, 50)], + vec![ + tool(None, "shell", 700, ToolExecutionStatus::Succeeded), + tool(Some(call), "read_file", 5, ToolExecutionStatus::Succeeded), + ], + SessionDiagnosticStats::default(), + ), + None, + ); + + assert_eq!(projected.unattributed_tools.len(), 1); + assert_eq!(projected.unattributed_tools[0].capability_name, "shell"); + assert_eq!(projected.iterations[0].tool_calls, 1); + assert_eq!(projected.totals.retained_tool_ms, Some(705)); + } + + #[test] + fn aggregate_totals_preserve_counts_and_unavailable_samples() { + let call = DiagnosticModelCallId::new(); + let projected = project_timings( + snapshot( + vec![model_call(call, 1, 50)], + vec![tool(Some(call), "shell", 7, ToolExecutionStatus::Failed)], + SessionDiagnosticStats { + total_model_calls: 3, + total_tool_calls: 4, + failed_tool_calls: 2, + total_latency_ms: DiagnosticMetricTotal { + known_total: 50, + unavailable_samples: 1, + }, + ..SessionDiagnosticStats::default() + }, + ), + None, + ); + + assert_eq!(projected.totals.iterations, 3); + assert_eq!(projected.totals.tool_calls, 4); + assert_eq!(projected.totals.failed_tool_calls, 2); + assert_eq!(projected.totals.inference_ms.known_total, 50); + assert_eq!(projected.totals.inference_ms.unavailable_samples, 1); + assert_eq!(projected.totals.retained_tool_ms, Some(7)); + assert!(!projected.totals.retained_tool_ms_complete); + } + + #[test] + fn timing_totals_decode_legacy_payload_without_completion_flag() { + let legacy = r#"{ + "iterations": 1, + "tool_calls": 0, + "failed_tool_calls": 0, + "inference_ms": { + "known_total": 5, + "unavailable_samples": 0 + } + }"#; + + let totals: RunArtifactTimingTotals = + serde_json::from_str(legacy).expect("legacy timing totals should decode"); + assert!(!totals.retained_tool_ms_complete); + + let serialized = serde_json::to_string(&totals).expect("timing totals should serialize"); + assert!(!serialized.contains("retained_tool_ms_complete")); + } + + #[test] + fn timing_sums_saturate_at_u64_max() { + let call = DiagnosticModelCallId::new(); + let projected = project_timings( + snapshot( + vec![model_call(call, 1, 1)], + vec![ + tool( + Some(call), + "first", + u64::MAX, + ToolExecutionStatus::Succeeded, + ), + tool(Some(call), "second", 1, ToolExecutionStatus::Succeeded), + ], + SessionDiagnosticStats::default(), + ), + None, + ); + + assert_eq!(projected.iterations[0].tool_ms_total, Some(u64::MAX)); + assert_eq!(projected.totals.retained_tool_ms, Some(u64::MAX)); + } + + #[test] + fn a_tool_pointing_at_an_unknown_call_is_unattributed_not_dropped() { + let known = DiagnosticModelCallId::new(); + let dangling = DiagnosticModelCallId::new(); + let projected = project_timings( + snapshot( + vec![model_call(known, 1, 50)], + vec![tool( + Some(dangling), + "shell", + 3, + ToolExecutionStatus::Succeeded, + )], + SessionDiagnosticStats::default(), + ), + None, + ); + + assert_eq!(projected.unattributed_tools.len(), 1); + assert_eq!(projected.iterations[0].tool_calls, 0); + } + + #[test] + fn no_bounded_payload_text_reaches_the_projection() { + let call = DiagnosticModelCallId::new(); + // Built through `new` rather than by mutating fields: `new` derives + // `output_bytes` from `result`, and the `TryFrom` wire guard rejects a + // record whose byte metadata disagrees with its result. + let leaky = ToolExecutionDiagnostic::new( + CapabilityActivityId::new(), + Some(call), + "shell", + Some("{\"path\":\"/home/alice/secret.txt\"}".to_string()), + Some("sk-ant-super-secret".to_string()), + ToolExecutionStatus::Succeeded, + Some(1), + None, + None, + None, + ); + + let projected = project_timings( + snapshot( + vec![model_call(call, 1, 50)], + vec![leaky], + SessionDiagnosticStats::default(), + ), + None, + ); + + let serialized = serde_json::to_string(&projected).expect("serialize timings"); + assert!(!serialized.contains("secret.txt")); + assert!(!serialized.contains("sk-ant-super-secret")); + assert!(!serialized.contains("/home/alice")); + } + + #[test] + fn the_default_value_reports_itself_as_unavailable() { + let absent = RunArtifactTimings::default(); + assert!(!absent.available); + assert!(!absent.complete); + assert!(absent.iterations.is_empty()); + } +} diff --git a/crates/product/ironclaw_assistant/src/reborn_services/thread_artifact.rs b/crates/product/ironclaw_assistant/src/reborn_services/thread_artifact.rs index ca85ff4a18f..8bfb7f08a6b 100644 --- a/crates/product/ironclaw_assistant/src/reborn_services/thread_artifact.rs +++ b/crates/product/ironclaw_assistant/src/reborn_services/thread_artifact.rs @@ -7,6 +7,7 @@ use ironclaw_product_contracts::surface::{ use ironclaw_trace_commons::contribution::DeterministicTraceRedactor; use ironclaw_threads::{BoundedThreadMessages, BoundedThreadMessagesRequest}; +use ironclaw_turns::TurnRunId; use serde::{Deserialize, Serialize}; use ironclaw_product_contracts::views::{RebornViewDescriptor, RebornViewProvider}; @@ -14,7 +15,10 @@ use ironclaw_product_contracts::views::{RebornViewDescriptor, RebornViewProvider use super::{ ProductCapabilityInvoker, RebornServices, RunArtifactLogs, RunArtifactMessage, RunArtifactRedaction, map_timeline_probe_error, parse_thread_id_field, - run_artifact::{ARTIFACT_REDACTION_PIPELINE, artifact_messages, context_messages_by_id}, + run_artifact::{ + ARTIFACT_REDACTION_PIPELINE, artifact_messages, context_messages_by_id, + timings::RunArtifactTimings, + }, thread_scope_from_turn_scope, }; @@ -36,9 +40,19 @@ pub struct RebornThreadArtifact { pub thread_id: String, pub messages: Vec, pub logs: RunArtifactLogs, + /// One timing block per run with a durable message timestamp. Runs absent + /// from the process-local store retain an explicit unavailable block. + #[serde(default)] + pub timings_by_run: Vec, pub redaction: RunArtifactRedaction, } +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct RunArtifactRunTimings { + pub run_id: String, + pub timings: RunArtifactTimings, +} + impl RebornServices where I: ProductCapabilityInvoker + Clone + 'static, @@ -76,6 +90,41 @@ where let context_by_id = context_messages_by_id(snapshot.context.messages); let (messages, message_redaction_applied) = artifact_messages(snapshot.history.messages, &context_by_id, &redactor); + + // Compute per-run timings before `artifact_logs` takes `caller` by + // value. Single pass: bucket messages by run first, so each run's + // timing computation (`derive_wall_clock_ms`'s `updated_at` scan in + // particular) sees only its own messages instead of the whole + // thread's — passing the full `messages` list per run would let one + // run's wall-clock span reach into a later run's activity, and would + // also re-scan the whole list once per distinct run. + let (runs_in_order, messages_by_run) = group_messages_by_run(&messages); + + let mut timings_by_run = Vec::new(); + for run_id in &runs_in_order { + let run_messages = &messages_by_run[run_id]; + // No per-run `received_at` here; the earliest message creation in + // the run is the closest durable origin available. + // silent-ok: no message in the run carries `created_at` (only + // possible for pre-timestamp records); wall-clock timing has no + // origin to measure from, so this run is left out rather than + // reported with a fabricated span. + let Some(origin) = run_messages.iter().filter_map(|m| m.created_at).min() else { + continue; + }; + let timings = self.artifact_timings( + &caller, + &thread_id, + run_id, + origin, + run_messages.iter().copied(), + ); + timings_by_run.push(RunArtifactRunTimings { + run_id: run_id.to_string(), + timings, + }); + } + let (logs, log_redaction_applied) = self .artifact_logs(caller, &thread_id, None, &redactor) .await; @@ -86,6 +135,7 @@ where thread_id: thread_id.to_string(), messages, logs, + timings_by_run, redaction: RunArtifactRedaction { pipeline: ARTIFACT_REDACTION_PIPELINE.to_string(), applied: message_redaction_applied || log_redaction_applied, @@ -104,3 +154,94 @@ where fn thread_artifact_too_large() -> ProductSurfaceError { ProductSurfaceError::from_status(ProductSurfaceErrorCode::InvalidRequest, 413, false) } + +fn group_messages_by_run( + messages: &[RunArtifactMessage], +) -> ( + Vec, + std::collections::HashMap>, +) { + let mut runs_in_order = Vec::new(); + let mut messages_by_run = std::collections::HashMap::new(); + for message in messages { + // silent-ok: a message with no run_id (pre-turn history, or a + // non-run message kind) simply contributes no timings entry. + let Some(run_id_text) = message.run_id.as_deref() else { + continue; + }; + let Ok(run_id) = TurnRunId::parse(run_id_text) else { + continue; + }; + messages_by_run + .entry(run_id) + .or_insert_with(|| { + runs_in_order.push(run_id); + Vec::new() + }) + .push(message); + } + (runs_in_order, messages_by_run) +} + +#[cfg(test)] +mod tests { + use super::*; + use chrono::{DateTime, Duration, Utc}; + use ironclaw_threads::{MessageKind, MessageStatus}; + + fn message( + run_id: &str, + created_at: DateTime, + updated_at: DateTime, + ) -> RunArtifactMessage { + RunArtifactMessage { + message_id: format!("message-{run_id}"), + sequence: 1, + run_id: Some(run_id.to_string()), + created_at: Some(created_at), + updated_at: Some(updated_at), + kind: MessageKind::Assistant, + status: MessageStatus::Finalized, + content: "hello".to_string(), + tool_call: None, + } + } + + #[test] + fn per_run_message_groups_keep_exact_wall_clock_inputs_without_cloning() { + let origin = Utc::now(); + let run_a = TurnRunId::new(); + let run_b = TurnRunId::new(); + let messages = vec![ + message( + &run_a.to_string(), + origin, + origin + Duration::milliseconds(12), + ), + message( + &run_b.to_string(), + origin + Duration::seconds(2), + origin + Duration::seconds(2) + Duration::milliseconds(34), + ), + ]; + + let (runs, messages_by_run) = group_messages_by_run(&messages); + assert_eq!(runs, vec![run_a, run_b]); + assert_eq!(messages_by_run[&run_a].len(), 1); + assert_eq!(messages_by_run[&run_b].len(), 1); + assert_eq!( + crate::reborn_services::timings_source::derive_wall_clock_ms( + origin, + messages_by_run[&run_a].iter().copied(), + ), + Some(12), + ); + assert_eq!( + crate::reborn_services::timings_source::derive_wall_clock_ms( + origin + Duration::seconds(2), + messages_by_run[&run_b].iter().copied(), + ), + Some(34), + ); + } +} diff --git a/crates/product/ironclaw_assistant/src/reborn_services/timings_source.rs b/crates/product/ironclaw_assistant/src/reborn_services/timings_source.rs new file mode 100644 index 00000000000..7b291f5b6a6 --- /dev/null +++ b/crates/product/ironclaw_assistant/src/reborn_services/timings_source.rs @@ -0,0 +1,137 @@ +//! The one impure edge of the timing lane: read the process-local diagnostic +//! store for a run and hand the snapshot to the pure projection. +//! +//! Mirrors `artifact_logs`: best-effort, returns a value on every path, and +//! never fails an artifact export. A user downloading evidence for a bug +//! report must always get a file. + +use chrono::{DateTime, Utc}; +use ironclaw_host_api::ids::ThreadId; +use ironclaw_product_contracts::inspector::DiagnosticScope; +use ironclaw_product_contracts::surface::ProductSurfaceCaller; +use ironclaw_product_contracts::views::RebornViewProvider; +use ironclaw_turns::TurnRunId; + +use crate::reborn_services::run_artifact::RunArtifactMessage; +use crate::reborn_services::run_artifact::timings::{ + RunArtifactTimings, project_timings, unavailable, +}; +use crate::reborn_services::{ProductCapabilityInvoker, RebornServices}; + +impl RebornServices +where + I: ProductCapabilityInvoker + Clone + 'static, + V: RebornViewProvider + Clone + 'static, +{ + pub(super) fn artifact_timings<'a, M>( + &self, + caller: &ProductSurfaceCaller, + thread_id: &ThreadId, + run_id: &TurnRunId, + run_received_at: DateTime, + messages: M, + ) -> RunArtifactTimings + where + M: IntoIterator, + { + // Same keying as the operator inspector (`inspector.rs::diagnostic_scope`). + // On the admin thread-scrape route the caller was already rebound to the + // scraped user by `thread_scrape_subject`, so this needs no branch. + let scope = DiagnosticScope::new( + caller.tenant_id.clone(), + caller.user_id.clone(), + thread_id.clone(), + *run_id, + ); + let wall_clock_ms = derive_wall_clock_ms(run_received_at, messages); + match self.diagnostic_store.timing_snapshot(&scope) { + Ok(Some(snapshot)) => project_timings(snapshot, wall_clock_ms), + Ok(None) => { + let mut timings = unavailable("run_not_resident"); + timings.totals.wall_clock_ms = wall_clock_ms; + timings + } + Err(error) => { + // debug!, not info!/warn!: the operator log buffer captures + // INFO+ and those entries are embedded into this same + // artifact's `logs` block (see `build_run_artifact` in + // `run_artifact.rs`). + tracing::debug!( + ?error, + "run artifact exported without optional process-local timings" + ); + let mut timings = unavailable("diagnostic_store_unavailable"); + timings.totals.wall_clock_ms = wall_clock_ms; + timings + } + } + } +} + +/// run.received_at → newest message `updated_at`, in milliseconds. +/// +/// Approximate by design: it folds in queue and persistence latency around +/// the run. `None` when no message carries a timestamp (pre-timestamp +/// records) or when the span is negative, which only happens if clocks +/// disagree — report nothing rather than a nonsense number. +pub(super) fn derive_wall_clock_ms<'a>( + run_received_at: DateTime, + messages: impl IntoIterator, +) -> Option { + let newest = messages + .into_iter() + .filter_map(|message| message.updated_at) + .max()?; + u64::try_from((newest - run_received_at).num_milliseconds()).ok() +} + +#[cfg(test)] +mod tests { + use super::*; + use ironclaw_threads::{MessageKind, MessageStatus}; + + fn message(updated_at: Option>) -> RunArtifactMessage { + RunArtifactMessage { + message_id: "message-a".to_string(), + sequence: 1, + run_id: Some("run-a".to_string()), + kind: MessageKind::Assistant, + status: MessageStatus::Finalized, + content: "hello".to_string(), + tool_call: None, + created_at: None, + updated_at, + } + } + + #[test] + fn wall_clock_spans_run_receipt_to_the_newest_message_update() { + let received = Utc::now(); + let messages = [ + message(Some(received + chrono::Duration::seconds(12))), + message(Some(received + chrono::Duration::seconds(91))), + message(Some(received + chrono::Duration::seconds(40))), + ]; + + assert_eq!( + derive_wall_clock_ms(received, messages.iter()), + Some(91_000) + ); + } + + #[test] + fn wall_clock_is_absent_when_no_message_carries_a_timestamp() { + let received = Utc::now(); + assert_eq!( + derive_wall_clock_ms(received, [message(None), message(None)].iter()), + None + ); + } + + #[test] + fn wall_clock_is_absent_rather_than_negative_when_clocks_disagree() { + let received = Utc::now(); + let messages = [message(Some(received - chrono::Duration::seconds(5)))]; + assert_eq!(derive_wall_clock_ms(received, messages.iter()), None); + } +} diff --git a/crates/product/ironclaw_assistant/tests/reborn_services_contract.rs b/crates/product/ironclaw_assistant/tests/reborn_services_contract.rs index 7aecdbc347c..4ad664bad45 100644 --- a/crates/product/ironclaw_assistant/tests/reborn_services_contract.rs +++ b/crates/product/ironclaw_assistant/tests/reborn_services_contract.rs @@ -23,6 +23,9 @@ use ironclaw_approvals::{ ToolPermissionOverrideKey, }; use ironclaw_assistant::EXTENSION_REGISTER_HOSTED_MCP_CAPABILITY_ID; +use ironclaw_assistant::inspector_store::{ + DiagnosticStoreError, DiagnosticStoreLimits, DiagnosticStorePort, InMemoryDiagnosticStore, +}; use ironclaw_assistant::{ ADMIN_THREAD_SCRAPE_ARTIFACT_VIEW, ADMIN_THREAD_SCRAPE_RUN_ARTIFACT_VIEW, ADMIN_THREAD_SCRAPE_THREADS_VIEW, ADMIN_USER_DELETE_CAPABILITY_ID, @@ -62,7 +65,7 @@ use ironclaw_assistant::{ ProductCapabilityInvoker, ProductNewCommandInput, ProductNewCommandOutput, ProductStatusCommandInput, ProductSurfaceFailure, ProjectCaller, ProjectFilesystemReader, ProjectFsEntry, ProjectFsEntryKind, ProjectFsError, ProjectFsFile, ProjectFsStat, - RUN_ARTIFACT_VIEW, RebornAccountTracesResponse, RebornAddMemberRequest, + RUN_ARTIFACT_SCHEMA, RUN_ARTIFACT_VIEW, RebornAccountTracesResponse, RebornAddMemberRequest, RebornAttachmentRequest, RebornAutomationInfo, RebornAutomationMutationResponse, RebornAutomationRecentRunInfo, RebornAutomationRecentRunStatus, RebornAutomationRequest, RebornAutomationRunStatus, RebornAutomationSource, RebornAutomationState, @@ -126,8 +129,9 @@ use ironclaw_host_api::product_adapter::{ ProtocolAuthFailure, RedactedString, }; use ironclaw_host_api::turn::{ - AcceptedMessageRef, EventCursor, ReplyTargetBindingRef, RunProfileId, RunProfileVersion, - SanitizedFailure, TurnActor, TurnGateRef, TurnId, TurnRunId, TurnScope, TurnStatus, + AcceptedMessageRef, CapabilityActivityId, EventCursor, ReplyTargetBindingRef, RunProfileId, + RunProfileVersion, SanitizedFailure, TurnActor, TurnGateRef, TurnId, TurnRunId, TurnScope, + TurnStatus, }; use ironclaw_host_api::{ capability::{EffectKind, PermissionMode}, @@ -153,6 +157,11 @@ use ironclaw_product_contracts::inbound_requests::{ ProductListThreadsRequest, ProductRenameAutomationRequest, ProductResolveGateRequest, ProductRetryRunRequest, ProductSetupExtensionRequest, ProductSubmitTurnRequest, }; +use ironclaw_product_contracts::inspector::{ + BoundedDiagnosticText, DiagnosticActivityEvent, DiagnosticCursor, DiagnosticModelCallId, + DiagnosticScope, DiagnosticSnapshot, DiagnosticUpdateBatch, InspectorModelCallStatus, + ModelCallDiagnostic, PromptDiagnostic, ToolExecutionDiagnostic, +}; use ironclaw_product_contracts::ironhub::{ IRONHUB_DELIVER_INSTALL_COMMAND_ID, IronhubInstallDeliveryRequest, IronhubInstallDeliveryResult, IronhubLinkError, IronhubLinkService, IronhubRegisterRequest, @@ -8634,6 +8643,123 @@ async fn run_artifact_selects_one_owned_run_and_queries_only_its_scoped_logs() { Some(run_id.to_string().as_str()) ); assert_eq!(requests[0].limit, Some(500)); + + // The diagnostic store was never populated for this run: the export must + // still succeed, say so honestly, and keep the durable timestamp floor. + assert!(!artifact.timings.available); + assert_eq!( + artifact.timings.unavailable_reason.as_deref(), + Some("run_not_resident") + ); + assert!(artifact.timings.iterations.is_empty()); + assert!( + artifact + .messages + .iter() + .any(|message| message.created_at.is_some()), + "durable message timestamps must survive an absent diagnostic store" + ); + assert_eq!(artifact.schema, RUN_ARTIFACT_SCHEMA); +} + +/// A diagnostic store whose every method fails, standing in for a backend +/// outage. The single guarantee under test: a user filing a bug report must +/// always get a file, so a diagnostic-store failure must never fail the +/// artifact export. +#[derive(Default)] +struct FailingDiagnosticStore; + +impl DiagnosticStorePort for FailingDiagnosticStore { + fn record_activity( + &self, + _scope: DiagnosticScope, + _event: DiagnosticActivityEvent, + ) -> Result { + Err(DiagnosticStoreError::StateUnavailable) + } + + fn snapshot( + &self, + _scope: &DiagnosticScope, + ) -> Result, DiagnosticStoreError> { + Err(DiagnosticStoreError::StateUnavailable) + } + + fn prompt( + &self, + _scope: &DiagnosticScope, + ) -> Result, DiagnosticStoreError> { + Err(DiagnosticStoreError::StateUnavailable) + } + + fn tool_execution( + &self, + _scope: &DiagnosticScope, + _activity_id: CapabilityActivityId, + ) -> Result, DiagnosticStoreError> { + Err(DiagnosticStoreError::StateUnavailable) + } + + fn updates_after( + &self, + _scope: &DiagnosticScope, + _after: Option, + ) -> Result { + Err(DiagnosticStoreError::StateUnavailable) + } +} + +#[tokio::test] +async fn run_artifact_reports_diagnostic_store_failure_without_failing_the_export() { + let owner = caller(); + let thread_scope = thread_scope_for(&owner); + let thread_id = ThreadId::new("thread-artifact-store-failure").expect("thread id"); + let run_id = TurnRunId::parse(&run_id_string()).expect("run id"); + let thread_service = Arc::new(InMemorySessionThreadService::default()); + thread_service + .ensure_thread(EnsureThreadRequest { + scope: thread_scope.clone(), + thread_id: Some(thread_id.clone()), + created_by_actor_id: owner.user_id.as_str().to_string(), + title: None, + metadata_json: None, + }) + .await + .expect("thread"); + seed_submitted_message( + &thread_service, + &thread_scope, + &thread_id, + &run_id, + "diagnostic store is down", + ) + .await; + let services = session_services(thread_service, Arc::new(FakeTurnCoordinator::default())) + .with_diagnostic_store(Arc::new(FailingDiagnosticStore)); + + let page = services + .query( + owner, + RebornViewQuery { + view_id: RUN_ARTIFACT_VIEW.id.to_string(), + params: serde_json::to_value(RebornRunArtifactRequest { + thread_id: thread_id.to_string(), + run_id: run_id.to_string(), + }) + .expect("artifact params"), + cursor: None, + }, + ) + .await + .expect("artifact export must succeed even when the diagnostic store errors"); + let artifact: RebornRunArtifact = + serde_json::from_value(page.payload).expect("artifact payload"); + + assert!(!artifact.timings.available); + assert_eq!( + artifact.timings.unavailable_reason.as_deref(), + Some("diagnostic_store_unavailable") + ); } #[tokio::test] @@ -8762,6 +8888,120 @@ async fn thread_artifact_includes_all_owned_runs_and_queries_thread_scoped_logs( assert_eq!(requests[0].limit, Some(500)); } +fn diagnostic_model_call(status: InspectorModelCallStatus) -> ModelCallDiagnostic { + ModelCallDiagnostic { + call_id: DiagnosticModelCallId::new(), + iteration: 1, + requested_model: BoundedDiagnosticText::label("test-model"), + effective_model: None, + started_at: Utc::now(), + completed_at: Some(Utc::now()), + duration_ms: Some(1), + status, + usage: None, + failure_summary: None, + } +} + +/// A thread with two runs must expose one timing entry per run. The exact +/// per-run wall-clock projection is pinned with deterministic timestamps in +/// the `thread_artifact` unit test; this route test proves both entries survive +/// the caller-owned export path. +#[tokio::test] +async fn thread_artifact_per_run_timings_do_not_reach_into_another_runs_activity() { + let owner = caller(); + let thread_scope = thread_scope_for(&owner); + let thread_id = ThreadId::new("thread-multi-run-timings").expect("thread id"); + let run_a = TurnRunId::parse(&run_id_string()).expect("run id"); + let run_b = TurnRunId::new(); + let run_c = TurnRunId::new(); + let thread_service = Arc::new(InMemorySessionThreadService::default()); + thread_service + .ensure_thread(EnsureThreadRequest { + scope: thread_scope.clone(), + thread_id: Some(thread_id.clone()), + created_by_actor_id: owner.user_id.as_str().to_string(), + title: None, + metadata_json: None, + }) + .await + .expect("thread"); + + seed_submitted_message(&thread_service, &thread_scope, &thread_id, &run_a, "run a").await; + seed_submitted_message(&thread_service, &thread_scope, &thread_id, &run_b, "run b").await; + seed_submitted_message(&thread_service, &thread_scope, &thread_id, &run_c, "run c").await; + + let diagnostic_store = + InMemoryDiagnosticStore::new(DiagnosticStoreLimits::default()).expect("diagnostic store"); + diagnostic_store + .record_model_call( + DiagnosticScope::new( + owner.tenant_id.clone(), + owner.user_id.clone(), + thread_id.clone(), + run_a, + ), + diagnostic_model_call(InspectorModelCallStatus::Succeeded), + ) + .expect("run a model call recorded"); + diagnostic_store + .record_model_call( + DiagnosticScope::new( + owner.tenant_id.clone(), + owner.user_id.clone(), + thread_id.clone(), + run_b, + ), + diagnostic_model_call(InspectorModelCallStatus::Succeeded), + ) + .expect("run b model call recorded"); + + let services = session_services(thread_service, Arc::new(FakeTurnCoordinator::default())) + .with_diagnostic_store(Arc::new(diagnostic_store)); + + let page = services + .query( + owner, + RebornViewQuery { + view_id: THREAD_ARTIFACT_VIEW.id.to_string(), + params: serde_json::to_value(RebornThreadArtifactRequest { + thread_id: thread_id.to_string(), + }) + .expect("artifact params"), + cursor: None, + }, + ) + .await + .expect("thread artifact"); + let artifact: RebornThreadArtifact = + serde_json::from_value(page.payload).expect("artifact payload"); + + assert_eq!(artifact.timings_by_run.len(), 3); + let run_a_timing = artifact + .timings_by_run + .iter() + .find(|entry| entry.run_id == run_a.to_string()) + .expect("run a timing entry"); + assert!(run_a_timing.timings.available); + assert!(run_a_timing.timings.totals.wall_clock_ms.is_some()); + let run_b_timing = artifact + .timings_by_run + .iter() + .find(|entry| entry.run_id == run_b.to_string()) + .expect("run b timing entry"); + assert!(run_b_timing.timings.available); + let run_c_timing = artifact + .timings_by_run + .iter() + .find(|entry| entry.run_id == run_c.to_string()) + .expect("run c timing entry"); + assert!(!run_c_timing.timings.available); + assert_eq!( + run_c_timing.timings.unavailable_reason.as_deref(), + Some("run_not_resident") + ); +} + #[tokio::test] async fn thread_artifact_projects_messages_from_the_bounded_snapshot() { let owner = caller(); diff --git a/crates/product/ironclaw_webui/tests/webui_v2_handlers_contract.rs b/crates/product/ironclaw_webui/tests/webui_v2_handlers_contract.rs index b75daaf5710..960e142fc40 100644 --- a/crates/product/ironclaw_webui/tests/webui_v2_handlers_contract.rs +++ b/crates/product/ironclaw_webui/tests/webui_v2_handlers_contract.rs @@ -85,7 +85,7 @@ use ironclaw_assistant::{ RebornStreamEventsRequest, RebornStreamEventsResponse, RebornSubmitTurnResponse, RebornThreadArtifact, RebornThreadArtifactRequest, RebornTimelineRequest, RebornTimelineResponse, RebornTraceCreditsResponse, RebornTraceHoldAuthorizeProductRequest, - RebornTraceHoldAuthorizeResponse, RunArtifactLogs, RunArtifactRedaction, + RebornTraceHoldAuthorizeResponse, RunArtifactLogs, RunArtifactRedaction, RunArtifactTimings, SKILL_AUTO_ACTIVATE_LEARNED_SET_CAPABILITY_ID, SKILL_AUTO_ACTIVATE_SET_CAPABILITY_ID, SKILL_CONTENT_VIEW, SKILL_INSTALL_CAPABILITY_ID, SKILL_REMOVE_CAPABILITY_ID, SKILL_SEARCH_VIEW, SKILL_UPDATE_CAPABILITY_ID, SKILLS_VIEW, THREAD_ARTIFACT_SCHEMA, THREAD_ARTIFACT_VIEW, @@ -424,6 +424,7 @@ fn stub_thread_artifact(thread_id: String) -> RebornThreadArtifact { unavailable_reason: None, entries: Vec::new(), }, + timings_by_run: Vec::new(), redaction: RunArtifactRedaction { pipeline: "deterministic-trace-redactor-v1".to_string(), applied: false, @@ -460,6 +461,7 @@ fn stub_run_artifact(thread_id: String, run_id: TurnRunId) -> RebornRunArtifact unavailable_reason: None, entries: Vec::new(), }, + timings: RunArtifactTimings::default(), redaction: RunArtifactRedaction { pipeline: "deterministic-trace-redactor-v1".to_string(), applied: false, diff --git a/tests/CLAUDE.md b/tests/CLAUDE.md index 43e7e541a6e..3b5a5a0ef6d 100644 --- a/tests/CLAUDE.md +++ b/tests/CLAUDE.md @@ -62,7 +62,7 @@ Tier-selection rule: `.claude/rules/testing.md`. | Providers (Google/Slack/GitHub contracts) | — | — | ✓ | ✓ | | Coverage/meta gates | — | 2 | ✓ | ✓ | -Totals: **59** group scenarios · **60** flat integration bins (53 in +Totals: **59** group scenarios · **61** flat integration bins (54 in `tests/integration/`, 7 in `tests/integration/auth/`) · **39** top-level Rust bins · **102** Python scenario files (**869** test functions) registered in the active Reborn coverage map below. Section 6 separately inventories retained and legacy @@ -181,7 +181,7 @@ ones speak MTProto over a raw socket with no injectable seam. --- -## 4. Flat integration bins — `tests/integration/*.rs` and `tests/integration/auth/*.rs` (60) +## 4. Flat integration bins — `tests/integration/*.rs` and `tests/integration/auth/*.rs` (61) One thread, whole real turn. Grouped by what the user experiences. @@ -281,8 +281,9 @@ One thread, whole real turn. Grouped by what the user experiences. | Enroll/refresh/remove a browser for web push over the real routes — advertised VAPID key, endpoint redacted to its push-service host, undeclared push hosts rejected, and the `web-app` catalog row selectable through the same notification-channels wire as every vendor channel | `webui_v2_product_api.rs::browser_channel_notification_setup_round_trip_through_production_facade` | | Identity resolution runs on the coverage lane | `identity_resolution_smoke.rs` | | A canonical 10-tool-call agent turn's database write volume is measured and reported (for tracking, not gated) on both libSQL and Postgres, and custom-actor group threads are rejected from canonical durable milestones | `db_write_canonical.rs` | +| A downloaded run artifact carries per-iteration model-call timing evidence for a completed run, and still carries durable per-message timestamps (with an explicit `run_not_resident` reason) when the process-local timing buffer was evicted or the process restarted | `run_artifact_timings.rs` | -One of the 60 registered bins, `delivery_user_journeys.rs`, holds the explicit +One of the 61 registered bins, `delivery_user_journeys.rs`, holds the explicit channel-delivery journeys (two-lane model): | A user can… | Scenario | diff --git a/tests/integration/run_artifact_timings.rs b/tests/integration/run_artifact_timings.rs new file mode 100644 index 00000000000..1bd80f4add7 --- /dev/null +++ b/tests/integration/run_artifact_timings.rs @@ -0,0 +1,155 @@ +//! Timing evidence in the user-downloadable run artifact, driven through the +//! real product surface (the real `webui_v2` router over a real +//! `RebornServices`, mirroring `webui_v2_product_api.rs`'s pattern). +//! +//! SCOPE NOTE: the integration harness wires the prompt diagnostic sink (so +//! model-call/inference timings are observable) but has NO tool diagnostic +//! sink (`staged_capability_io_for_test` / +//! `staged_capability_io_with_observer_for_test` in +//! `crates/app/ironclaw_composition/src/runtime/capability_host.rs` hardcode +//! `tool_diagnostic_sink: None`, unlike production's `capability_wiring` in +//! `runtime.rs`). Tool-execution timings (`iterations[].tool_calls`, +//! `totals.tool_calls`, per-tool durations) are therefore NOT observable here +//! and are deliberately not asserted below — see the plan's Task 6 report for +//! the follow-up. The brief's third test +//! (`exported_timings_carry_no_tool_arguments_or_results`) is deliberately +//! omitted: with no tool sink there are no tool entries at all, so it would +//! pass vacuously. The no-payload-leak guarantee is already pinned +//! non-vacuously at crate tier by +//! `no_bounded_payload_text_reaches_the_projection` in +//! `crates/product/ironclaw_assistant/src/reborn_services/run_artifact/timings.rs`. + +#[allow(dead_code)] +#[path = "support/mod.rs"] +mod reborn_support; +#[allow(dead_code)] +#[path = "../support/mod.rs"] +mod support; + +use std::sync::Arc; +use std::time::Duration; + +use axum::Router; +use ironclaw_assistant::{RebornRunArtifact, RebornServices}; +use ironclaw_product_contracts::surface::ProductSurface; +use ironclaw_webui::webui_v2::{ + DEFAULT_SSE_MAX_CONCURRENT_PER_CALLER, WebUiV2Capabilities, WebUiV2State, webui_v2_router, +}; +use reborn_support::builder::RebornIntegrationHarness; +use reborn_support::reply::RebornScriptedReply; +use reborn_support::webui_mount::{get_json, webui_caller_for}; + +/// The run/thread artifact export routes are mounted only when the +/// deployment opts in (`WebUiV2State::with_regression_artifact_export_enabled`, +/// off by default in production). `reborn_support::webui_mount::mount_webui_v2_router` +/// deliberately leaves that flag off for its other callers, so this test +/// builds its own router with the flag on — mirroring +/// `crates/product/ironclaw_webui/tests/webui_v2_handlers_contract.rs::artifact_router_with`. +fn artifact_export_router( + services: Arc, + caller: ironclaw_product_contracts::surface::ProductSurfaceCaller, +) -> Router { + webui_v2_router( + WebUiV2State::new(services, DEFAULT_SSE_MAX_CONCURRENT_PER_CALLER) + .with_regression_artifact_export_enabled(true), + ) + .layer(axum::Extension(caller)) + .layer(axum::Extension(WebUiV2Capabilities::default())) +} + +#[tokio::test] +async fn exported_run_artifact_carries_per_iteration_timings() { + // Three model iterations (two tool-call turns, one final reply). Tool + // execution timing is unobservable in this harness (see module doc), so + // this drives the run purely to exercise the model-call/inference timing + // path — the one lane the harness's diagnostic sink actually captures. + let h = RebornIntegrationHarness::builder("thread-timings") + .with_builtin_http_tools() + .with_model_call_delay_for_test(Duration::from_millis(5)) + .script([ + RebornScriptedReply::tool_call( + "builtin.http", + serde_json::json!({"url": "https://example.com/a"}), + ), + RebornScriptedReply::tool_call( + "builtin.http", + serde_json::json!({"url": "https://example.com/b"}), + ), + RebornScriptedReply::text("done"), + ]) + .build() + .await + .expect("harness"); + let run_id = h.submit_turn("check the build").await.expect("run"); + + let services = RebornServices::new(h.thread_harness.service.clone(), h.coordinator.clone()) + .with_diagnostic_store(h.diagnostic_store()); + let router = artifact_export_router(Arc::new(services), webui_caller_for(&h.binding)); + + let (status, body) = get_json( + router, + &format!( + "/api/webchat/v2/threads/{}/runs/{}/artifact", + h.binding.thread_id.as_str(), + run_id + ), + ) + .await; + assert_eq!(status, axum::http::StatusCode::OK, "artifact body: {body}"); + let artifact: RebornRunArtifact = serde_json::from_value(body).expect("artifact deserializes"); + + assert!(artifact.timings.available); + assert!(!artifact.timings.complete); + assert_eq!(artifact.timings.iterations.len(), 3); + assert_eq!(artifact.timings.totals.iterations, 3); + assert!(artifact.timings.totals.inference_ms.known_total > 0); + assert_eq!(artifact.timings.totals.inference_ms.unavailable_samples, 0); + assert!(artifact.timings.totals.wall_clock_ms.is_some()); + assert!( + artifact.timings.iterations[0].inference_ms.is_some(), + "model-call latency arrives through the prompt diagnostic sink" + ); +} + +#[tokio::test] +async fn exported_run_artifact_keeps_timestamps_when_timings_were_evicted() { + let h = RebornIntegrationHarness::builder("thread-evicted") + .script([RebornScriptedReply::text("hello")]) + .build() + .await + .expect("harness"); + let run_id = h.submit_turn("hello").await.expect("run"); + + // No `with_diagnostic_store`: services fall back to their own empty store, + // which is exactly the post-restart / post-eviction case a user hits when + // filing a bug report the next day. + let services = RebornServices::new(h.thread_harness.service.clone(), h.coordinator.clone()); + let router = artifact_export_router(Arc::new(services), webui_caller_for(&h.binding)); + + let (status, body) = get_json( + router, + &format!( + "/api/webchat/v2/threads/{}/runs/{}/artifact", + h.binding.thread_id.as_str(), + run_id + ), + ) + .await; + assert_eq!(status, axum::http::StatusCode::OK, "artifact body: {body}"); + let artifact: RebornRunArtifact = + serde_json::from_value(body).expect("artifact export must still succeed"); + + assert!(!artifact.timings.available); + assert_eq!( + artifact.timings.unavailable_reason.as_deref(), + Some("run_not_resident") + ); + assert!( + artifact + .messages + .iter() + .any(|message| message.created_at.is_some()), + "the durable timestamp floor must survive an absent diagnostic store" + ); + assert!(artifact.timings.totals.wall_clock_ms.is_some()); +} diff --git a/tests/integration/support/builder.rs b/tests/integration/support/builder.rs index 43934b4cca4..80bf764f560 100644 --- a/tests/integration/support/builder.rs +++ b/tests/integration/support/builder.rs @@ -1131,6 +1131,16 @@ impl RebornIntegrationHarness { std::sync::Arc::clone(&self.workflow) } + /// The SAME `InMemoryDiagnosticStore` the loop's diagnostic sinks write + /// into. Hand this to `RebornServices::with_diagnostic_store` so a test + /// reads what the run actually recorded, exactly as production connects + /// the two in `product_surface.rs:98`. + pub(crate) fn diagnostic_store( + &self, + ) -> std::sync::Arc { + std::sync::Arc::clone(&self._shared.diagnostic_store) + } + /// A fresh binding-service instance over the GROUP-shared product /// harness storage — the SAME durable binding ledger the workflow above /// resolves through at admission. Delivery proofs hand this to the diff --git a/tests/integration/support/group.rs b/tests/integration/support/group.rs index 33390499020..071a409c267 100644 --- a/tests/integration/support/group.rs +++ b/tests/integration/support/group.rs @@ -284,6 +284,13 @@ pub(crate) struct GroupSharedStorage { /// Read by `RebornThreadBuilder::build()` to decide whether to wire the /// real approval/auth interaction services into the thread's workflow. pub(crate) real_gate_dispatch_services: bool, + /// Task 5 (run-artifact-timings): the ONE `InMemoryDiagnosticStore` wired + /// into the group's ONE planned runtime as `prompt_diagnostic_sink`, + /// mirroring production's single-store shape (`runtime.rs:3455`, + /// `product_surface.rs:98`). Retained so `RebornIntegrationHarness:: + /// diagnostic_store()` reads exactly what the loop actually recorded, + /// instead of a second throwaway instance nothing observes. + pub(crate) diagnostic_store: Arc, } impl GroupSharedStorage { @@ -1252,6 +1259,12 @@ impl RebornIntegrationGroupBuilder { None => Arc::clone(&user_profile_source), }; let reply_attachment_intent_port = capability.reply_attachment_intent_port(); + // Production parity: ONE store, connected to the loop's diagnostic + // sinks and readable by product services — see runtime.rs:3455 and + // product_surface.rs:98. The previous inline `default()` was write-only, + // so nothing could observe what the loop recorded. + let diagnostic_store = + Arc::new(ironclaw_assistant::inspector_store::InMemoryDiagnosticStore::default()); let parts = DefaultPlannedRuntimeParts { process_system: process_system.clone(), thread_service: runtime_thread_service, @@ -1362,9 +1375,8 @@ impl RebornIntegrationGroupBuilder { attachment_read_port: capability_recorder .attachment_test_support() .map(|support| support.read_port), - prompt_diagnostic_sink: Some(Arc::new( - ironclaw_assistant::inspector_store::InMemoryDiagnosticStore::default(), - )), + prompt_diagnostic_sink: Some(Arc::clone(&diagnostic_store) + as Arc), reply_attachment_intent_port: Some(reply_attachment_intent_port), // §5.2.9 render-from-record: the SAME durable gate-record store this // group's capability port persists `GateRecord::Auth` into, so the @@ -1413,6 +1425,7 @@ impl RebornIntegrationGroupBuilder { planned_runtime_parts_shape, real_gate_dispatch_services: self.real_gate_dispatch_services, channel_connection: self.channel_connection, + diagnostic_store, }), }) } diff --git a/tests/integration/wiring_parity.rs b/tests/integration/wiring_parity.rs index 31face8243f..90359d7fedd 100644 --- a/tests/integration/wiring_parity.rs +++ b/tests/integration/wiring_parity.rs @@ -34,11 +34,13 @@ mod support; use std::collections::HashSet; use ironclaw_host_api::ids::CapabilityId; +use ironclaw_product_contracts::inspector::DiagnosticScope; use reborn_support::builder::RebornIntegrationHarness; use reborn_support::group::RebornIntegrationGroup; use reborn_support::harness::HarnessResult; use reborn_support::harness::options::ToolsProfile; use reborn_support::planned_runtime_parts_shape::DefaultPlannedRuntimePartsShape; +use reborn_support::reply::RebornScriptedReply; // --------------------------------------------------------------------------- // Part 1: DefaultPlannedRuntimeParts Some/None shape parity @@ -218,6 +220,37 @@ async fn test_default_planned_runtime_parts_shape_matches_production() { ); } +/// Task 5 (run-artifact-timings): the harness must hand out the SAME +/// `InMemoryDiagnosticStore` instance the loop's diagnostic sinks write into, +/// mirroring production's one-store shape (`runtime.rs:3455`, +/// `product_surface.rs:98`). Before this test, `group.rs` built a throwaway +/// store inline and kept no handle, so nothing could read what a run actually +/// recorded. +#[tokio::test] +async fn harness_shares_one_diagnostic_store_with_the_loop() { + let harness = RebornIntegrationHarness::builder("wiring-parity") + .script([RebornScriptedReply::text("diagnostic")]) + .build() + .await + .expect("harness"); + let run_id = harness + .submit_turn("write diagnostics") + .await + .expect("scripted turn"); + + let snapshot = harness + .diagnostic_store() + .snapshot(&DiagnosticScope::new( + harness.binding.tenant_id.clone(), + harness.binding.actor_user_id.clone(), + harness.binding.thread_id.clone(), + run_id, + )) + .expect("diagnostic snapshot") + .expect("the loop must write to the shared diagnostic store"); + assert!(!snapshot.model_calls.is_empty()); +} + #[tokio::test] async fn builtin_tools_planned_runtime_parts_shape_matches_production() { let group = RebornIntegrationGroup::builtin_tools() diff --git a/tests/reborn_qa_recorded_behavior.rs b/tests/reborn_qa_recorded_behavior.rs index c531444fccc..4f4177d095c 100644 --- a/tests/reborn_qa_recorded_behavior.rs +++ b/tests/reborn_qa_recorded_behavior.rs @@ -1275,12 +1275,14 @@ async fn replay_routine_phrase_fires(case: &QaPhrase, cron_fragment: &str) { ); let now = Utc::now(); - // At an exact cron boundary the live poller may claim the trigger before - // the test forces it due. Do not schedule a second fire in that case. + // Make the persisted slot due without moving it behind the current cron + // boundary. If this used `now - 120s`, a `*/30` trigger started just after + // `:00` or `:30` would reschedule to that already-passed boundary after + // the first fire, letting the poller submit the same test routine twice. let already_due_or_fired = trigger.is_due_at(now) || trigger.has_active_fire() || trigger.last_fired_slot.is_some(); if !already_due_or_fired { - trigger.next_run_at = now - chrono::Duration::try_seconds(120).expect("duration"); + trigger.next_run_at = now; repo.upsert_trigger(trigger) .await .expect("make replayed routine due");