diff --git a/crates/app/ironclaw_architecture_tests/tests/reborn_dependency_boundaries.rs b/crates/app/ironclaw_architecture_tests/tests/reborn_dependency_boundaries.rs index b947c74ade8..2ef2ab33c39 100644 --- a/crates/app/ironclaw_architecture_tests/tests/reborn_dependency_boundaries.rs +++ b/crates/app/ironclaw_architecture_tests/tests/reborn_dependency_boundaries.rs @@ -626,7 +626,12 @@ fn reborn_contracts_crates_carry_a_checked_size_ceiling() { // reviewed in the PR body's architecture-audit section. ("ironclaw_host_api", 18_570), ("ironclaw_loop_contracts", 14_479), - ("ironclaw_product_contracts", 14_471), + // Raised 14_471 -> 15_685 by #7230 (operator inspector diagnostics): + // the growth is neutral snapshot, update, cursor, prompt, model-call, + // tool-execution, activity, and stats wire vocabulary plus serde and + // bounded-field invariants. Process-local retention and sequencing + // remain implemented by ironclaw_assistant. + ("ironclaw_product_contracts", 15_685), ("ironclaw_prompt_envelope", 832), ]; diff --git a/crates/contracts/ironclaw_product_contracts/src/inspector.rs b/crates/contracts/ironclaw_product_contracts/src/inspector.rs new file mode 100644 index 00000000000..602d7239e27 --- /dev/null +++ b/crates/contracts/ironclaw_product_contracts/src/inspector.rs @@ -0,0 +1,1447 @@ +//! Operator-only diagnostic vocabulary for the Web Debug Inspector. +//! +//! These DTOs cross the product boundary through the dedicated inspection +//! surface; snapshots and updates flow out, while cursors round-trip to resume +//! reads. They are deliberately separate from product projection events: raw +//! prompt components and tool details must never enter the normal product +//! stream. + +use chrono::{DateTime, Utc}; +use ironclaw_host_api::{ + ids::{TenantId, ThreadId, UserId}, + turn::{CapabilityActivityId, TurnRunId}, +}; +use serde::{Deserialize, Serialize}; +use uuid::Uuid; + +pub const PROMPT_COMPONENT_CONTENT_MAX_BYTES: usize = 64 * 1024; +pub const PROMPT_COMPONENT_TOTAL_MAX_BYTES: usize = 256 * 1024; +pub const RECONSTRUCTED_PROMPT_MAX_BYTES: usize = 256 * 1024; +pub const TOOL_ARGUMENTS_MAX_BYTES: usize = 64 * 1024; +pub const TOOL_RESULT_MAX_BYTES: usize = 50 * 1024; +pub const DIAGNOSTIC_LABEL_MAX_BYTES: usize = 256; +pub const DIAGNOSTIC_SUMMARY_MAX_BYTES: usize = 2 * 1024; +pub const MAX_PROMPT_COMPONENTS: usize = 128; +pub const MAX_ACTIVE_SKILLS: usize = 64; +pub const MAX_MODELS_IN_STATS: usize = 64; +// Keep the process-wide defaults conservative because retained tool payloads +// may each contain both bounded arguments and a bounded result. +pub const DEFAULT_MAX_ACTIVITY_ENTRIES: usize = 1_000; +pub const DEFAULT_MAX_TRACKED_SESSIONS: usize = 8; +pub const DEFAULT_MAX_RETAINED_RUNS_PER_SESSION: usize = 2; +pub const DEFAULT_MAX_LIVE_UPDATE_SCOPES: usize = + DEFAULT_MAX_TRACKED_SESSIONS * DEFAULT_MAX_RETAINED_RUNS_PER_SESSION; +pub const DEFAULT_MAX_MODEL_CALLS_PER_RUN: usize = 128; +pub const DEFAULT_MAX_TOOL_EXECUTIONS_PER_RUN: usize = 16; +pub const DEFAULT_MAX_RETAINED_UPDATES_PER_RUN: usize = 1_024; + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)] +#[serde(transparent)] +pub struct DiagnosticModelCallId(Uuid); + +impl DiagnosticModelCallId { + pub fn new() -> Self { + Self(Uuid::new_v4()) + } + + pub fn from_uuid(value: Uuid) -> Self { + Self(value) + } + + pub fn as_uuid(self) -> Uuid { + self.0 + } +} + +impl Default for DiagnosticModelCallId { + fn default() -> Self { + Self::new() + } +} + +impl std::fmt::Display for DiagnosticModelCallId { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!(formatter, "{}", self.0) + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)] +#[serde(transparent)] +pub struct DiagnosticStreamId(Uuid); + +impl DiagnosticStreamId { + pub fn new() -> Self { + Self(Uuid::new_v4()) + } + + pub fn as_uuid(self) -> Uuid { + self.0 + } +} + +impl Default for DiagnosticStreamId { + fn default() -> Self { + Self::new() + } +} + +impl std::fmt::Display for DiagnosticStreamId { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!(formatter, "{}", self.0) + } +} + +#[derive( + Debug, Clone, Copy, Default, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize, +)] +#[serde(transparent)] +pub struct DiagnosticSequence(u64); + +impl DiagnosticSequence { + pub const ZERO: Self = Self(0); + + pub const fn new(value: u64) -> Self { + Self(value) + } + + pub const fn as_u64(self) -> u64 { + self.0 + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)] +pub struct DiagnosticCursor { + pub stream_id: DiagnosticStreamId, + pub sequence: DiagnosticSequence, +} + +impl DiagnosticCursor { + pub const fn new(stream_id: DiagnosticStreamId, sequence: DiagnosticSequence) -> Self { + Self { + stream_id, + sequence, + } + } +} + +#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize)] +pub struct DiagnosticScope { + pub tenant_id: TenantId, + pub user_id: UserId, + pub thread_id: ThreadId, + pub run_id: TurnRunId, +} + +impl DiagnosticScope { + pub fn new( + tenant_id: TenantId, + user_id: UserId, + thread_id: ThreadId, + run_id: TurnRunId, + ) -> Self { + Self { + tenant_id, + user_id, + thread_id, + run_id, + } + } +} + +/// UTF-8 text with explicit original-size and truncation metadata. +/// +/// Construction is limited to the purpose-specific constructors so callers +/// cannot silently select an unbounded maximum. +#[derive(Clone, PartialEq, Eq, Serialize)] +pub struct BoundedDiagnosticText { + content: String, + original_bytes: u64, + truncated: bool, +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct BoundedDiagnosticTextWire { + content: String, + original_bytes: u64, + truncated: bool, +} + +impl<'de> Deserialize<'de> for BoundedDiagnosticText { + fn deserialize(deserializer: D) -> Result + where + D: serde::Deserializer<'de>, + { + let wire = BoundedDiagnosticTextWire::deserialize(deserializer)?; + if wire.content.len() > RECONSTRUCTED_PROMPT_MAX_BYTES { + return Err(serde::de::Error::custom( + "diagnostic text exceeds the maximum retained byte length", + )); + } + let retained_bytes = u64::try_from(wire.content.len()).map_err(|_| { + serde::de::Error::custom("diagnostic text retained byte length is not representable") + })?; + if wire.original_bytes < retained_bytes { + return Err(serde::de::Error::custom( + "diagnostic text original byte length is smaller than retained text", + )); + } + if wire.truncated != (wire.original_bytes > retained_bytes) { + return Err(serde::de::Error::custom( + "diagnostic text truncation metadata is inconsistent", + )); + } + Ok(Self { + content: wire.content, + original_bytes: wire.original_bytes, + truncated: wire.truncated, + }) + } +} + +impl std::fmt::Debug for BoundedDiagnosticText { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter + .debug_struct("BoundedDiagnosticText") + .field("content", &"[diagnostic content redacted]") + .field("retained_bytes", &self.content.len()) + .field("original_bytes", &self.original_bytes) + .field("truncated", &self.truncated) + .finish() + } +} + +impl BoundedDiagnosticText { + pub fn label(value: impl Into) -> Self { + Self::bounded(value.into(), DIAGNOSTIC_LABEL_MAX_BYTES) + } + + pub fn summary(value: impl Into) -> Self { + Self::bounded(value.into(), DIAGNOSTIC_SUMMARY_MAX_BYTES) + } + + pub fn prompt_component(value: impl Into) -> Self { + Self::bounded(value.into(), PROMPT_COMPONENT_CONTENT_MAX_BYTES) + } + + pub fn reconstructed_prompt(value: impl Into) -> Self { + Self::bounded(value.into(), RECONSTRUCTED_PROMPT_MAX_BYTES) + } + + pub fn tool_arguments(value: impl Into) -> Self { + Self::bounded(value.into(), TOOL_ARGUMENTS_MAX_BYTES) + } + + pub fn tool_result(value: impl Into) -> Self { + Self::bounded(value.into(), TOOL_RESULT_MAX_BYTES) + } + + pub fn content(&self) -> &str { + &self.content + } + + pub const fn original_bytes(&self) -> u64 { + self.original_bytes + } + + pub const fn truncated(&self) -> bool { + self.truncated + } + + fn validate_retained_max(self, max_bytes: usize) -> Result { + if self.content.len() > max_bytes { + return Err("diagnostic text exceeds its field byte limit"); + } + Ok(self) + } + + fn rebound(self, max_bytes: usize) -> Self { + if self.content.len() <= max_bytes { + return self; + } + let original_bytes = self.original_bytes; + let mut bounded = Self::bounded(self.content, max_bytes); + bounded.original_bytes = original_bytes; + bounded.truncated = true; + bounded + } + + fn bounded(value: String, max_bytes: usize) -> Self { + let original_bytes = u64::try_from(value.len()).unwrap_or(u64::MAX); + if value.len() <= max_bytes { + return Self { + content: value, + original_bytes, + truncated: false, + }; + } + let mut end = max_bytes.min(value.len()); + while end > 0 && !value.is_char_boundary(end) { + end -= 1; + } + Self { + content: value[..end].to_string(), // safety: `end` is a verified UTF-8 boundary. + original_bytes, + truncated: true, + } + } +} + +fn deserialize_bounded_label<'de, D>(deserializer: D) -> Result +where + D: serde::Deserializer<'de>, +{ + BoundedDiagnosticText::deserialize(deserializer)? + .validate_retained_max(DIAGNOSTIC_LABEL_MAX_BYTES) + .map_err(serde::de::Error::custom) +} + +fn deserialize_bounded_prompt_component<'de, D>( + deserializer: D, +) -> Result +where + D: serde::Deserializer<'de>, +{ + BoundedDiagnosticText::deserialize(deserializer)? + .validate_retained_max(PROMPT_COMPONENT_CONTENT_MAX_BYTES) + .map_err(serde::de::Error::custom) +} + +fn deserialize_optional_bounded_tool_arguments<'de, D>( + deserializer: D, +) -> Result, D::Error> +where + D: serde::Deserializer<'de>, +{ + Option::::deserialize(deserializer)? + .map(|value| value.validate_retained_max(TOOL_ARGUMENTS_MAX_BYTES)) + .transpose() + .map_err(serde::de::Error::custom) +} + +fn deserialize_optional_bounded_tool_result<'de, D>( + deserializer: D, +) -> Result, D::Error> +where + D: serde::Deserializer<'de>, +{ + Option::::deserialize(deserializer)? + .map(|value| value.validate_retained_max(TOOL_RESULT_MAX_BYTES)) + .transpose() + .map_err(serde::de::Error::custom) +} + +fn deserialize_optional_bounded_label<'de, D>( + deserializer: D, +) -> Result, D::Error> +where + D: serde::Deserializer<'de>, +{ + Option::::deserialize(deserializer)? + .map(|value| value.validate_retained_max(DIAGNOSTIC_LABEL_MAX_BYTES)) + .transpose() + .map_err(serde::de::Error::custom) +} + +fn deserialize_optional_bounded_summary<'de, D>( + deserializer: D, +) -> Result, D::Error> +where + D: serde::Deserializer<'de>, +{ + Option::::deserialize(deserializer)? + .map(|value| value.validate_retained_max(DIAGNOSTIC_SUMMARY_MAX_BYTES)) + .transpose() + .map_err(serde::de::Error::custom) +} + +fn deserialize_bounded_model_counts<'de, D>( + deserializer: D, +) -> Result, D::Error> +where + D: serde::Deserializer<'de>, +{ + let values = Vec::::deserialize(deserializer)?; + if values.len() > MAX_MODELS_IN_STATS { + return Err(serde::de::Error::custom( + "diagnostic model counts exceed the item limit", + )); + } + Ok(values) +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum PromptComponentKind { + System, + Identity, + Instruction, + Skill, + Capability, + Conversation, + Other, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct PromptComponentDiagnostic { + pub kind: PromptComponentKind, + #[serde(deserialize_with = "deserialize_bounded_label")] + pub label: BoundedDiagnosticText, + #[serde(deserialize_with = "deserialize_bounded_prompt_component")] + pub content: BoundedDiagnosticText, + pub estimated_tokens: Option, +} + +impl PromptComponentDiagnostic { + pub fn new( + kind: PromptComponentKind, + label: impl Into, + content: impl Into, + estimated_tokens: Option, + ) -> Self { + Self { + kind, + label: BoundedDiagnosticText::label(label), + content: BoundedDiagnosticText::prompt_component(content), + estimated_tokens, + } + } +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +pub struct PromptDiagnostic { + pub captured_at: DateTime, + pub components: Vec, + pub components_truncated: bool, + pub reconstructed_prompt: BoundedDiagnosticText, + pub total_estimated_tokens: Option, + pub message_count: u32, + pub identity_message_count: u32, + pub instruction_snippet_count: u32, + pub active_skills: Vec, + pub active_skills_truncated: bool, + pub capability_count: u32, + pub requested_model: Option, + pub effective_model: Option, + pub context_limit: Option, +} + +impl PromptDiagnostic { + // arch-exempt: too_many_args, one validated path bounds the prompt DTO, plan #7219 + #[allow(clippy::too_many_arguments)] + pub fn new( + captured_at: DateTime, + components: Vec, + reconstructed_prompt: impl Into, + total_estimated_tokens: Option, + message_count: u32, + identity_message_count: u32, + instruction_snippet_count: u32, + active_skills: Vec, + capability_count: u32, + requested_model: Option, + effective_model: Option, + context_limit: Option, + ) -> Self { + let active_skills_truncated = active_skills.len() > MAX_ACTIVE_SKILLS; + let active_skills = active_skills + .into_iter() + .take(MAX_ACTIVE_SKILLS) + .map(BoundedDiagnosticText::label) + .collect(); + + Self { + captured_at, + components, + components_truncated: false, + reconstructed_prompt: BoundedDiagnosticText::reconstructed_prompt(reconstructed_prompt), + total_estimated_tokens, + message_count, + identity_message_count, + instruction_snippet_count, + active_skills, + active_skills_truncated, + capability_count, + requested_model: requested_model.map(BoundedDiagnosticText::label), + effective_model: effective_model.map(BoundedDiagnosticText::label), + context_limit, + } + .into_bounded() + } + + /// Reapplies every owning prompt limit at a trust boundary. + /// + /// The fields remain public wire DTO fields, so callers may construct this + /// type without using [`Self::new`]. Store and transport boundaries should + /// call this method before retaining an externally constructed value. + pub fn into_bounded(mut self) -> Self { + let original_component_count = self.components.len(); + let mut remaining = PROMPT_COMPONENT_TOTAL_MAX_BYTES; + let mut bounded_components = + Vec::with_capacity(original_component_count.min(MAX_PROMPT_COMPONENTS)); + for mut component in self.components.into_iter().take(MAX_PROMPT_COMPONENTS) { + if remaining == 0 { + break; + } + component.label = component.label.rebound(DIAGNOSTIC_LABEL_MAX_BYTES); + component.content = component + .content + .rebound(PROMPT_COMPONENT_CONTENT_MAX_BYTES); + let retained = component.content.content().len().min(remaining); + component.content = component.content.rebound(retained); + remaining = remaining.saturating_sub(component.content.content().len()); + bounded_components.push(component); + } + self.components_truncated |= bounded_components.len() < original_component_count; + self.components = bounded_components; + + let original_skill_count = self.active_skills.len(); + self.active_skills = self + .active_skills + .into_iter() + .take(MAX_ACTIVE_SKILLS) + .map(|skill| skill.rebound(DIAGNOSTIC_LABEL_MAX_BYTES)) + .collect(); + self.active_skills_truncated |= self.active_skills.len() < original_skill_count; + self.reconstructed_prompt = self + .reconstructed_prompt + .rebound(RECONSTRUCTED_PROMPT_MAX_BYTES); + self.requested_model = self + .requested_model + .map(|model| model.rebound(DIAGNOSTIC_LABEL_MAX_BYTES)); + self.effective_model = self + .effective_model + .map(|model| model.rebound(DIAGNOSTIC_LABEL_MAX_BYTES)); + self + } + + pub fn any_content_truncated(&self) -> bool { + self.components_truncated + || self.reconstructed_prompt.truncated() + || self.active_skills_truncated + || self + .components + .iter() + .any(|component| component.label.truncated() || component.content.truncated()) + || self + .active_skills + .iter() + .any(BoundedDiagnosticText::truncated) + || self + .requested_model + .as_ref() + .is_some_and(BoundedDiagnosticText::truncated) + || self + .effective_model + .as_ref() + .is_some_and(BoundedDiagnosticText::truncated) + } +} + +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] +pub struct ModelTokenUsage { + pub input_tokens: Option, + pub output_tokens: Option, + pub cache_read_input_tokens: Option, + pub cache_creation_input_tokens: Option, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum InspectorModelCallStatus { + Started, + Succeeded, + Failed, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct ModelCallDiagnostic { + pub call_id: DiagnosticModelCallId, + pub iteration: u32, + #[serde(deserialize_with = "deserialize_bounded_label")] + pub requested_model: BoundedDiagnosticText, + #[serde(default, deserialize_with = "deserialize_optional_bounded_label")] + pub effective_model: Option, + pub started_at: DateTime, + pub completed_at: Option>, + pub duration_ms: Option, + pub status: InspectorModelCallStatus, + pub usage: Option, + #[serde(default, deserialize_with = "deserialize_optional_bounded_summary")] + pub failure_summary: Option, +} + +impl ModelCallDiagnostic { + // arch-exempt: too_many_args, atomically construct one bounded model call, plan #7219 + #[allow(clippy::too_many_arguments)] + pub fn new( + call_id: DiagnosticModelCallId, + iteration: u32, + requested_model: impl Into, + effective_model: Option, + started_at: DateTime, + completed_at: Option>, + duration_ms: Option, + status: InspectorModelCallStatus, + usage: Option, + failure_summary: Option, + ) -> Self { + Self { + call_id, + iteration, + requested_model: BoundedDiagnosticText::label(requested_model), + effective_model: effective_model.map(BoundedDiagnosticText::label), + started_at, + completed_at, + duration_ms, + status, + usage, + failure_summary: failure_summary.map(BoundedDiagnosticText::summary), + } + } + + pub fn into_bounded(mut self) -> Self { + self.requested_model = self.requested_model.rebound(DIAGNOSTIC_LABEL_MAX_BYTES); + self.effective_model = self + .effective_model + .map(|model| model.rebound(DIAGNOSTIC_LABEL_MAX_BYTES)); + self.failure_summary = self + .failure_summary + .map(|summary| summary.rebound(DIAGNOSTIC_SUMMARY_MAX_BYTES)); + self + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum ToolExecutionStatus { + Started, + Succeeded, + Failed, +} + +#[derive(Deserialize)] +struct ToolExecutionDiagnosticWire { + activity_id: CapabilityActivityId, + model_call_id: Option, + #[serde(deserialize_with = "deserialize_bounded_label")] + capability_name: BoundedDiagnosticText, + #[serde( + default, + deserialize_with = "deserialize_optional_bounded_tool_arguments" + )] + arguments: Option, + #[serde(default, deserialize_with = "deserialize_optional_bounded_tool_result")] + result: Option, + status: ToolExecutionStatus, + duration_ms: Option, + output_bytes: Option, + #[serde(default, deserialize_with = "deserialize_optional_bounded_label")] + failure_category: Option, + #[serde(default, deserialize_with = "deserialize_optional_bounded_summary")] + failure_summary: Option, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(try_from = "ToolExecutionDiagnosticWire")] +pub struct ToolExecutionDiagnostic { + pub activity_id: CapabilityActivityId, + pub model_call_id: Option, + pub capability_name: BoundedDiagnosticText, + pub arguments: Option, + pub result: Option, + pub status: ToolExecutionStatus, + pub duration_ms: Option, + pub output_bytes: Option, + pub failure_category: Option, + pub failure_summary: Option, +} + +impl TryFrom for ToolExecutionDiagnostic { + type Error = &'static str; + + fn try_from(wire: ToolExecutionDiagnosticWire) -> Result { + if let Some(result) = wire.result.as_ref() + && wire.output_bytes != Some(result.original_bytes()) + { + return Err("tool result byte metadata is inconsistent"); + } + Ok(Self { + activity_id: wire.activity_id, + model_call_id: wire.model_call_id, + capability_name: wire.capability_name, + arguments: wire.arguments, + result: wire.result, + status: wire.status, + duration_ms: wire.duration_ms, + output_bytes: wire.output_bytes, + failure_category: wire.failure_category, + failure_summary: wire.failure_summary, + }) + } +} + +impl ToolExecutionDiagnostic { + // arch-exempt: too_many_args, atomically construct one bounded tool record, plan #7219 + #[allow(clippy::too_many_arguments)] + pub fn new( + activity_id: CapabilityActivityId, + model_call_id: Option, + capability_name: impl Into, + arguments: Option, + result: Option, + status: ToolExecutionStatus, + duration_ms: Option, + output_bytes: Option, + failure_category: Option, + failure_summary: Option, + ) -> Self { + let result = result.map(BoundedDiagnosticText::tool_result); + let output_bytes = result + .as_ref() + .map(BoundedDiagnosticText::original_bytes) + .or(output_bytes); + Self { + activity_id, + model_call_id, + capability_name: BoundedDiagnosticText::label(capability_name), + arguments: arguments.map(BoundedDiagnosticText::tool_arguments), + result, + status, + duration_ms, + output_bytes, + failure_category: failure_category.map(BoundedDiagnosticText::label), + failure_summary: failure_summary.map(BoundedDiagnosticText::summary), + } + } + + pub fn into_bounded(mut self) -> Self { + self.capability_name = self.capability_name.rebound(DIAGNOSTIC_LABEL_MAX_BYTES); + self.arguments = self + .arguments + .map(|arguments| arguments.rebound(TOOL_ARGUMENTS_MAX_BYTES)); + self.result = self + .result + .map(|result| result.rebound(TOOL_RESULT_MAX_BYTES)); + if let Some(result) = self.result.as_ref() { + self.output_bytes = Some(result.original_bytes()); + } + self.failure_category = self + .failure_category + .map(|category| category.rebound(DIAGNOSTIC_LABEL_MAX_BYTES)); + self.failure_summary = self + .failure_summary + .map(|summary| summary.rebound(DIAGNOSTIC_SUMMARY_MAX_BYTES)); + self + } + + pub fn result_truncated(&self) -> bool { + self.result + .as_ref() + .is_some_and(BoundedDiagnosticText::truncated) + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum DiagnosticActivityKind { + TurnStarted, + PromptPrepared, + ModelCallStarted, + ModelCallCompleted, + ModelCallFailed, + Progress, + ToolStarted, + ToolCompleted, + ToolFailed, + GateBlocked, + FinalResponseCompleted, + StreamDisconnected, + StreamResumed, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct DiagnosticActivityEvent { + pub occurred_at: DateTime, + pub kind: DiagnosticActivityKind, + pub iteration: Option, + pub activity_id: Option, + pub model_call_id: Option, + #[serde(default, deserialize_with = "deserialize_optional_bounded_summary")] + pub summary: Option, +} + +impl DiagnosticActivityEvent { + pub fn new( + occurred_at: DateTime, + kind: DiagnosticActivityKind, + iteration: Option, + activity_id: Option, + model_call_id: Option, + summary: Option, + ) -> Self { + Self { + occurred_at, + kind, + iteration, + activity_id, + model_call_id, + summary: summary.map(BoundedDiagnosticText::summary), + } + } + + pub fn into_bounded(mut self) -> Self { + self.summary = self + .summary + .map(|summary| summary.rebound(DIAGNOSTIC_SUMMARY_MAX_BYTES)); + self + } +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +pub struct DiagnosticActivityEntry { + pub sequence: DiagnosticSequence, + pub event: DiagnosticActivityEvent, +} + +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] +pub struct DiagnosticMetricTotal { + pub known_total: u64, + pub unavailable_samples: u64, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct DiagnosticModelCount { + #[serde(deserialize_with = "deserialize_bounded_label")] + pub model: BoundedDiagnosticText, + pub calls: u64, +} + +impl DiagnosticModelCount { + pub fn new(model: impl Into, calls: u64) -> Self { + Self { + model: BoundedDiagnosticText::label(model), + calls, + } + } + + pub fn into_bounded(mut self) -> Self { + self.model = self.model.rebound(DIAGNOSTIC_LABEL_MAX_BYTES); + self + } +} + +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] +pub struct SessionDiagnosticStats { + pub total_model_calls: u64, + #[serde(deserialize_with = "deserialize_bounded_model_counts")] + pub calls_per_model: Vec, + pub calls_per_model_truncated: bool, + pub input_tokens: DiagnosticMetricTotal, + pub output_tokens: DiagnosticMetricTotal, + pub cache_read_input_tokens: DiagnosticMetricTotal, + pub cache_creation_input_tokens: DiagnosticMetricTotal, + pub total_latency_ms: DiagnosticMetricTotal, + pub total_tool_calls: u64, + pub successful_tool_calls: u64, + pub failed_tool_calls: u64, +} + +impl SessionDiagnosticStats { + pub fn into_bounded(mut self) -> Self { + if self.calls_per_model.len() > MAX_MODELS_IN_STATS { + self.calls_per_model.truncate(MAX_MODELS_IN_STATS); + self.calls_per_model_truncated = true; + } + self.calls_per_model = self + .calls_per_model + .into_iter() + .map(DiagnosticModelCount::into_bounded) + .collect(); + self + } +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(tag = "type", content = "data", rename_all = "snake_case")] +pub enum DiagnosticUpdateKind { + PromptUpdated { + component_count: u32, + total_estimated_tokens: Option, + truncated: bool, + }, + ModelCall(ModelCallDiagnostic), + ToolExecutionUpdated { + activity_id: CapabilityActivityId, + model_call_id: Option, + #[serde(deserialize_with = "deserialize_bounded_label")] + capability_name: BoundedDiagnosticText, + status: ToolExecutionStatus, + duration_ms: Option, + output_bytes: Option, + result_truncated: bool, + }, + Activity(DiagnosticActivityEvent), + Stats(SessionDiagnosticStats), +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +pub struct DiagnosticUpdateEnvelope { + pub scope: DiagnosticScope, + pub stream_id: DiagnosticStreamId, + pub sequence: DiagnosticSequence, + pub emitted_at: DateTime, + pub update: DiagnosticUpdateKind, +} + +impl DiagnosticUpdateEnvelope { + pub const fn cursor(&self) -> DiagnosticCursor { + DiagnosticCursor::new(self.stream_id, self.sequence) + } +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +pub struct DiagnosticUpdateBatch { + pub updates: Vec, + pub retention_floor: Option, + pub latest_cursor: Option, + pub rebase_required: bool, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +pub struct DiagnosticSnapshot { + pub scope: DiagnosticScope, + pub stream_id: DiagnosticStreamId, + pub prompt: Option, + pub model_calls: Vec, + pub tool_executions: Vec, + pub activity: Vec, + pub stats: SessionDiagnosticStats, + pub latest_sequence: DiagnosticSequence, +} + +#[cfg(test)] +mod tests { + use serde::de::DeserializeOwned; + + use super::*; + + fn assert_unit_enum_wire(cases: &[(T, &str)]) + where + T: std::fmt::Debug + PartialEq + Serialize + DeserializeOwned, + { + for (value, wire_name) in cases { + let encoded = serde_json::to_value(value).expect("serialize wire enum"); + assert_eq!(encoded, serde_json::Value::String((*wire_name).to_string())); + let decoded: T = serde_json::from_value(encoded).expect("deserialize wire enum"); + assert_eq!(&decoded, value); + } + } + + fn oversized_bounded_text(max_bytes: usize) -> serde_json::Value { + serde_json::json!({ + "content": "x".repeat(max_bytes + 1), + "original_bytes": max_bytes + 1, + "truncated": false, + }) + } + + #[test] + fn bounded_text_preserves_utf8_and_reports_original_size() { + let value = "€".repeat(TOOL_RESULT_MAX_BYTES); + let bounded = BoundedDiagnosticText::tool_result(value.clone()); + assert!(bounded.truncated()); + assert!(bounded.content().len() <= TOOL_RESULT_MAX_BYTES); + assert_eq!( + bounded.content().chars().count(), + TOOL_RESULT_MAX_BYTES / '€'.len_utf8() + ); + assert_eq!(bounded.original_bytes(), value.len() as u64); + } + + #[test] + fn diagnostic_cursor_round_trips_through_json() { + let cursor = DiagnosticCursor::new(DiagnosticStreamId::new(), DiagnosticSequence::new(42)); + + let encoded = serde_json::to_value(cursor).expect("serialize cursor"); + let decoded: DiagnosticCursor = + serde_json::from_value(encoded).expect("deserialize cursor"); + + assert_eq!(decoded, cursor); + } + + #[test] + fn unit_wire_enums_round_trip_with_their_snake_case_names() { + assert_unit_enum_wire(&[ + (PromptComponentKind::System, "system"), + (PromptComponentKind::Identity, "identity"), + (PromptComponentKind::Instruction, "instruction"), + (PromptComponentKind::Skill, "skill"), + (PromptComponentKind::Capability, "capability"), + (PromptComponentKind::Conversation, "conversation"), + (PromptComponentKind::Other, "other"), + ]); + assert_unit_enum_wire(&[ + (InspectorModelCallStatus::Started, "started"), + (InspectorModelCallStatus::Succeeded, "succeeded"), + (InspectorModelCallStatus::Failed, "failed"), + ]); + assert_unit_enum_wire(&[ + (ToolExecutionStatus::Started, "started"), + (ToolExecutionStatus::Succeeded, "succeeded"), + (ToolExecutionStatus::Failed, "failed"), + ]); + assert_unit_enum_wire(&[ + (DiagnosticActivityKind::TurnStarted, "turn_started"), + (DiagnosticActivityKind::PromptPrepared, "prompt_prepared"), + ( + DiagnosticActivityKind::ModelCallStarted, + "model_call_started", + ), + ( + DiagnosticActivityKind::ModelCallCompleted, + "model_call_completed", + ), + (DiagnosticActivityKind::ModelCallFailed, "model_call_failed"), + (DiagnosticActivityKind::Progress, "progress"), + (DiagnosticActivityKind::ToolStarted, "tool_started"), + (DiagnosticActivityKind::ToolCompleted, "tool_completed"), + (DiagnosticActivityKind::ToolFailed, "tool_failed"), + (DiagnosticActivityKind::GateBlocked, "gate_blocked"), + ( + DiagnosticActivityKind::FinalResponseCompleted, + "final_response_completed", + ), + ( + DiagnosticActivityKind::StreamDisconnected, + "stream_disconnected", + ), + (DiagnosticActivityKind::StreamResumed, "stream_resumed"), + ]); + } + + #[test] + fn diagnostic_update_variants_round_trip_with_their_tagged_wire_names() { + let model_call_id = DiagnosticModelCallId::new(); + let updates = [ + ( + DiagnosticUpdateKind::PromptUpdated { + component_count: 2, + total_estimated_tokens: Some(3), + truncated: false, + }, + "prompt_updated", + ), + ( + DiagnosticUpdateKind::ModelCall(ModelCallDiagnostic::new( + model_call_id, + 1, + "requested-model", + Some("effective-model".to_string()), + Utc::now(), + Some(Utc::now()), + Some(4), + InspectorModelCallStatus::Succeeded, + Some(ModelTokenUsage { + input_tokens: Some(5), + output_tokens: Some(6), + ..ModelTokenUsage::default() + }), + None, + )), + "model_call", + ), + ( + DiagnosticUpdateKind::ToolExecutionUpdated { + activity_id: CapabilityActivityId::new(), + model_call_id: Some(model_call_id), + capability_name: BoundedDiagnosticText::label("filesystem.read"), + status: ToolExecutionStatus::Succeeded, + duration_ms: Some(7), + output_bytes: Some(8), + result_truncated: false, + }, + "tool_execution_updated", + ), + ( + DiagnosticUpdateKind::Activity(DiagnosticActivityEvent::new( + Utc::now(), + DiagnosticActivityKind::Progress, + Some(1), + None, + Some(model_call_id), + Some("working".to_string()), + )), + "activity", + ), + ( + DiagnosticUpdateKind::Stats(SessionDiagnosticStats { + total_model_calls: 1, + calls_per_model: vec![DiagnosticModelCount::new("model", 1)], + ..SessionDiagnosticStats::default() + }), + "stats", + ), + ]; + + for (update, wire_name) in updates { + let encoded = serde_json::to_value(&update).expect("serialize diagnostic update"); + assert_eq!(encoded["type"], wire_name); + let decoded: DiagnosticUpdateKind = + serde_json::from_value(encoded).expect("deserialize diagnostic update"); + assert_eq!(decoded, update); + } + } + + #[test] + fn owning_diagnostic_records_round_trip_through_json() { + let component = PromptComponentDiagnostic::new( + PromptComponentKind::Instruction, + "policy", + "keep responses concise", + Some(4), + ); + let encoded = serde_json::to_value(&component).expect("serialize prompt component"); + let decoded = serde_json::from_value::(encoded) + .expect("deserialize prompt component"); + assert_eq!(decoded, component); + + let tool = ToolExecutionDiagnostic::new( + CapabilityActivityId::new(), + Some(DiagnosticModelCallId::new()), + "filesystem.read", + Some("{\"path\":\"notes.txt\"}".to_string()), + Some("contents".to_string()), + ToolExecutionStatus::Succeeded, + Some(2), + Some(1), + Some("none".to_string()), + Some("completed".to_string()), + ); + let encoded = serde_json::to_value(&tool).expect("serialize tool execution"); + let decoded = serde_json::from_value::(encoded) + .expect("deserialize tool execution"); + assert_eq!(decoded, tool); + } + + #[test] + fn prompt_component_deserialization_enforces_its_field_limits() { + let component = PromptComponentDiagnostic::new( + PromptComponentKind::Instruction, + "policy", + "content", + None, + ); + let payload = serde_json::to_value(component).expect("serialize prompt component"); + + for (field, max_bytes) in [ + ("label", DIAGNOSTIC_LABEL_MAX_BYTES), + ("content", PROMPT_COMPONENT_CONTENT_MAX_BYTES), + ] { + let mut oversized = payload.clone(); + oversized[field] = oversized_bounded_text(max_bytes); + + let error = serde_json::from_value::(oversized) + .expect_err("oversized prompt component field must be rejected"); + assert!(error.to_string().contains("field byte limit")); + } + } + + #[test] + fn tool_execution_deserialization_enforces_its_field_limits() { + let tool = ToolExecutionDiagnostic::new( + CapabilityActivityId::new(), + None, + "filesystem.read", + Some("arguments".to_string()), + Some("result".to_string()), + ToolExecutionStatus::Failed, + None, + None, + Some("provider".to_string()), + Some("failed".to_string()), + ); + let payload = serde_json::to_value(tool).expect("serialize tool execution"); + + for (field, max_bytes) in [ + ("capability_name", DIAGNOSTIC_LABEL_MAX_BYTES), + ("arguments", TOOL_ARGUMENTS_MAX_BYTES), + ("result", TOOL_RESULT_MAX_BYTES), + ("failure_category", DIAGNOSTIC_LABEL_MAX_BYTES), + ("failure_summary", DIAGNOSTIC_SUMMARY_MAX_BYTES), + ] { + let mut oversized = payload.clone(); + oversized[field] = oversized_bounded_text(max_bytes); + + let error = serde_json::from_value::(oversized) + .expect_err("oversized tool execution field must be rejected"); + assert!(error.to_string().contains("field byte limit")); + } + } + + #[test] + fn tool_execution_deserialization_rejects_inconsistent_result_size() { + let tool = ToolExecutionDiagnostic::new( + CapabilityActivityId::new(), + None, + "filesystem.read", + None, + Some("result".to_string()), + ToolExecutionStatus::Succeeded, + None, + None, + None, + None, + ); + let mut payload = serde_json::to_value(tool).expect("serialize tool execution"); + payload["output_bytes"] = serde_json::json!(1); + + let error = serde_json::from_value::(payload) + .expect_err("inconsistent tool result byte metadata must be rejected"); + assert!(error.to_string().contains("byte metadata is inconsistent")); + } + + #[test] + fn bounded_text_deserialization_rejects_the_global_retained_byte_limit() { + let content = "x".repeat(RECONSTRUCTED_PROMPT_MAX_BYTES + 1); + let payload = serde_json::json!({ + "content": content, + "original_bytes": RECONSTRUCTED_PROMPT_MAX_BYTES + 1, + "truncated": false, + }); + + let error = serde_json::from_value::(payload) + .expect_err("oversized diagnostic text must be rejected"); + + assert!(error.to_string().contains("maximum retained byte length")); + } + + #[test] + fn bounded_text_deserialization_rejects_inconsistent_metadata() { + for payload in [ + serde_json::json!({ + "content": "ok", + "original_bytes": 1, + "truncated": false, + }), + serde_json::json!({ + "content": "ok", + "original_bytes": 2, + "truncated": true, + }), + serde_json::json!({ + "content": "ok", + "original_bytes": 3, + "truncated": false, + }), + ] { + serde_json::from_value::(payload) + .expect_err("inconsistent diagnostic text metadata must be rejected"); + } + } + + #[test] + fn diagnostic_update_deserialization_enforces_the_owning_field_limit() { + let update = DiagnosticUpdateKind::ToolExecutionUpdated { + activity_id: CapabilityActivityId::new(), + model_call_id: None, + capability_name: BoundedDiagnosticText::label("filesystem.read"), + status: ToolExecutionStatus::Succeeded, + duration_ms: None, + output_bytes: None, + result_truncated: false, + }; + let mut payload = serde_json::to_value(update).expect("serialize diagnostic update"); + let oversized_label = "x".repeat(DIAGNOSTIC_LABEL_MAX_BYTES + 1); + payload["data"]["capability_name"] = serde_json::json!({ + "content": oversized_label, + "original_bytes": DIAGNOSTIC_LABEL_MAX_BYTES + 1, + "truncated": false, + }); + + let error = serde_json::from_value::(payload) + .expect_err("oversized capability name must be rejected"); + + assert!(error.to_string().contains("field byte limit")); + } + + #[test] + fn diagnostic_stats_deserialization_rejects_too_many_model_counts() { + let update = DiagnosticUpdateKind::Stats(SessionDiagnosticStats::default()); + let mut payload = serde_json::to_value(update).expect("serialize diagnostic stats"); + payload["data"]["calls_per_model"] = serde_json::to_value( + (0..=MAX_MODELS_IN_STATS) + .map(|index| DiagnosticModelCount::new(format!("model-{index}"), 1)) + .collect::>(), + ) + .expect("serialize model counts"); + + let error = serde_json::from_value::(payload) + .expect_err("too many model counts must be rejected"); + + assert!(error.to_string().contains("item limit")); + } + + #[test] + fn bounded_text_debug_never_exposes_content() { + let bounded = BoundedDiagnosticText::tool_result("super-secret-value"); + let debug = format!("{bounded:?}"); + assert!(!debug.contains("super-secret-value")); + assert!(debug.contains("diagnostic content redacted")); + } + + #[test] + fn prompt_constructor_applies_component_and_skill_caps() { + let components = (0..=MAX_PROMPT_COMPONENTS) + .map(|index| { + PromptComponentDiagnostic::new( + PromptComponentKind::Instruction, + format!("component-{index}"), + "x", + Some(1), + ) + }) + .collect(); + let skills = (0..=MAX_ACTIVE_SKILLS) + .map(|index| format!("skill-{index}")) + .collect(); + let prompt = PromptDiagnostic::new( + Utc::now(), + components, + "prompt", + Some(1), + 1, + 0, + 1, + skills, + 0, + None, + None, + None, + ); + assert_eq!(prompt.components.len(), MAX_PROMPT_COMPONENTS); + assert!(prompt.components_truncated); + assert_eq!(prompt.active_skills.len(), MAX_ACTIVE_SKILLS); + assert!(prompt.active_skills_truncated); + } + + #[test] + fn prompt_component_content_truncation_does_not_claim_the_list_was_shortened() { + let prompt = PromptDiagnostic::new( + Utc::now(), + vec![PromptComponentDiagnostic::new( + PromptComponentKind::Instruction, + "large-component", + "x".repeat(PROMPT_COMPONENT_CONTENT_MAX_BYTES + 1), + None, + )], + "prompt", + None, + 1, + 0, + 1, + Vec::new(), + 0, + None, + None, + None, + ); + + assert_eq!(prompt.components.len(), 1); + assert!(!prompt.components_truncated); + assert!(prompt.components[0].content.truncated()); + assert!(prompt.any_content_truncated()); + } + + #[test] + fn prompt_constructor_enforces_the_total_component_byte_budget() { + let mut components = (0..3) + .map(|index| { + PromptComponentDiagnostic::new( + PromptComponentKind::Instruction, + format!("full-{index}"), + "x".repeat(PROMPT_COMPONENT_CONTENT_MAX_BYTES), + None, + ) + }) + .collect::>(); + components.push(PromptComponentDiagnostic::new( + PromptComponentKind::Instruction, + "partial-prefix", + "x".repeat(40_000), + None, + )); + components.push(PromptComponentDiagnostic::new( + PromptComponentKind::Instruction, + "partially-retained", + "x".repeat(PROMPT_COMPONENT_CONTENT_MAX_BYTES), + None, + )); + components.push(PromptComponentDiagnostic::new( + PromptComponentKind::Instruction, + "dropped", + "x", + None, + )); + + let prompt = PromptDiagnostic::new( + Utc::now(), + components, + "prompt", + None, + 1, + 0, + 1, + Vec::new(), + 0, + None, + None, + None, + ); + + assert_eq!(prompt.components.len(), 5); + assert_eq!( + prompt + .components + .iter() + .map(|component| component.content.content().len()) + .sum::(), + PROMPT_COMPONENT_TOTAL_MAX_BYTES + ); + assert_eq!(prompt.components[4].content.content().len(), 25_536); + assert!(prompt.components[4].content.truncated()); + assert!(prompt.components_truncated); + } + + #[test] + fn tool_result_uses_its_original_length_over_caller_supplied_output_bytes() { + assert_eq!(TOOL_RESULT_MAX_BYTES, 50 * 1024); + let tool = ToolExecutionDiagnostic::new( + CapabilityActivityId::new(), + None, + "filesystem.read", + None, + Some("x".repeat(TOOL_RESULT_MAX_BYTES + 1)), + ToolExecutionStatus::Succeeded, + None, + Some(7), + None, + None, + ); + assert!(tool.result_truncated()); + assert_eq!(tool.output_bytes, Some((TOOL_RESULT_MAX_BYTES + 1) as u64)); + } + + #[test] + fn stats_bound_the_per_model_breakdown_and_mark_truncation() { + let stats = SessionDiagnosticStats { + calls_per_model: (0..=MAX_MODELS_IN_STATS) + .map(|index| DiagnosticModelCount::new(format!("model-{index}"), 1)) + .collect(), + ..SessionDiagnosticStats::default() + } + .into_bounded(); + assert_eq!(stats.calls_per_model.len(), MAX_MODELS_IN_STATS); + assert!(stats.calls_per_model_truncated); + } +} diff --git a/crates/contracts/ironclaw_product_contracts/src/lib.rs b/crates/contracts/ironclaw_product_contracts/src/lib.rs index 2db65e413fd..eb624321132 100644 --- a/crates/contracts/ironclaw_product_contracts/src/lib.rs +++ b/crates/contracts/ironclaw_product_contracts/src/lib.rs @@ -46,6 +46,7 @@ pub mod descriptors; pub mod error; pub mod inbound; pub mod inbound_requests; +pub mod inspector; pub mod interaction_commands; pub mod ironhub; pub mod lifecycle_service; diff --git a/crates/product/ironclaw_assistant/src/inspector_store.rs b/crates/product/ironclaw_assistant/src/inspector_store.rs new file mode 100644 index 00000000000..e1af24cb883 --- /dev/null +++ b/crates/product/ironclaw_assistant/src/inspector_store.rs @@ -0,0 +1,1429 @@ +//! Bounded, process-local storage for operator inspection diagnostics. +//! +//! The store deliberately has no persistence backend. It keeps raw diagnostic +//! content out of durable events and drops all state at process restart. + +use std::collections::{HashMap, VecDeque}; +use std::sync::{Arc, Mutex}; + +use chrono::Utc; +use ironclaw_host_api::{ + ids::{TenantId, ThreadId, UserId}, + turn::TurnRunId, +}; +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, PromptDiagnostic, SessionDiagnosticStats, + ToolExecutionDiagnostic, +}; +use thiserror::Error; +use tokio::sync::broadcast; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct DiagnosticStoreLimits { + pub max_sessions: usize, + pub max_runs_per_session: usize, + pub max_model_calls_per_run: usize, + pub max_tool_executions_per_run: usize, + pub max_activity_entries_per_run: usize, + pub max_updates_per_run: usize, + pub max_live_update_scopes: usize, + pub live_update_capacity: usize, +} + +impl Default for DiagnosticStoreLimits { + fn default() -> Self { + Self { + max_sessions: DEFAULT_MAX_TRACKED_SESSIONS, + max_runs_per_session: DEFAULT_MAX_RETAINED_RUNS_PER_SESSION, + max_model_calls_per_run: DEFAULT_MAX_MODEL_CALLS_PER_RUN, + max_tool_executions_per_run: DEFAULT_MAX_TOOL_EXECUTIONS_PER_RUN, + max_activity_entries_per_run: DEFAULT_MAX_ACTIVITY_ENTRIES, + max_updates_per_run: DEFAULT_MAX_RETAINED_UPDATES_PER_RUN, + max_live_update_scopes: DEFAULT_MAX_LIVE_UPDATE_SCOPES, + live_update_capacity: DEFAULT_MAX_RETAINED_UPDATES_PER_RUN, + } + } +} + +impl DiagnosticStoreLimits { + fn validate(self) -> Result { + // Keep this destructuring exhaustive: adding a limit field must fail + // compilation until its validation ceiling is defined below. + let Self { + max_sessions, + max_runs_per_session, + max_model_calls_per_run, + max_tool_executions_per_run, + max_activity_entries_per_run, + max_updates_per_run, + max_live_update_scopes, + live_update_capacity, + } = self; + let values = [ + ("max_sessions", max_sessions, DEFAULT_MAX_TRACKED_SESSIONS), + ( + "max_runs_per_session", + max_runs_per_session, + DEFAULT_MAX_RETAINED_RUNS_PER_SESSION, + ), + ( + "max_model_calls_per_run", + max_model_calls_per_run, + DEFAULT_MAX_MODEL_CALLS_PER_RUN, + ), + ( + "max_tool_executions_per_run", + max_tool_executions_per_run, + DEFAULT_MAX_TOOL_EXECUTIONS_PER_RUN, + ), + ( + "max_activity_entries_per_run", + max_activity_entries_per_run, + DEFAULT_MAX_ACTIVITY_ENTRIES, + ), + ( + "max_updates_per_run", + max_updates_per_run, + DEFAULT_MAX_RETAINED_UPDATES_PER_RUN, + ), + ( + "max_live_update_scopes", + max_live_update_scopes, + DEFAULT_MAX_LIVE_UPDATE_SCOPES, + ), + ( + "live_update_capacity", + live_update_capacity, + DEFAULT_MAX_RETAINED_UPDATES_PER_RUN, + ), + ]; + if let Some((name, _, _)) = values.iter().copied().find(|(_, value, _)| *value == 0) { + return Err(DiagnosticStoreError::InvalidLimit(name)); + } + if let Some((name, _, maximum)) = values + .iter() + .copied() + .find(|(_, value, maximum)| value > maximum) + { + return Err(DiagnosticStoreError::LimitExceedsMaximum { name, maximum }); + } + Ok(self) + } +} + +#[derive(Debug, Error, Clone, PartialEq, Eq)] +pub enum DiagnosticStoreError { + #[error("diagnostic store limit `{0}` must be non-zero")] + InvalidLimit(&'static str), + #[error("diagnostic store limit `{name}` exceeds its maximum of {maximum}")] + LimitExceedsMaximum { name: &'static str, maximum: usize }, + #[error("diagnostic store state is unavailable")] + StateUnavailable, + #[error("diagnostic sequence space is exhausted")] + SequenceExhausted, + #[error("diagnostic store invariant failed")] + Invariant, + #[error("diagnostic subscriber lagged by {0} updates")] + SubscriberLagged(u64), + #[error("diagnostic subscription is closed")] + SubscriptionClosed, +} + +#[derive(Debug, Clone, PartialEq, Eq, Hash)] +struct DiagnosticSessionKey { + tenant_id: TenantId, + user_id: UserId, + thread_id: ThreadId, +} + +impl From<&DiagnosticScope> for DiagnosticSessionKey { + fn from(scope: &DiagnosticScope) -> Self { + Self { + tenant_id: scope.tenant_id.clone(), + user_id: scope.user_id.clone(), + thread_id: scope.thread_id.clone(), + } + } +} + +#[derive(Debug)] +struct DiagnosticRunState { + stream_id: DiagnosticStreamId, + prompt: Option, + model_calls: VecDeque, + tool_executions: VecDeque, + activity: VecDeque, + stats: SessionDiagnosticStats, + updates: VecDeque, + latest_sequence: DiagnosticSequence, +} + +impl Default for DiagnosticRunState { + fn default() -> Self { + Self { + stream_id: DiagnosticStreamId::new(), + prompt: None, + model_calls: VecDeque::new(), + tool_executions: VecDeque::new(), + activity: VecDeque::new(), + stats: SessionDiagnosticStats::default(), + updates: VecDeque::new(), + latest_sequence: DiagnosticSequence::ZERO, + } + } +} + +#[derive(Debug, Default)] +struct DiagnosticSessionState { + runs: HashMap, + run_order: VecDeque, +} + +#[derive(Debug, Default)] +struct DiagnosticStoreState { + sessions: HashMap, + session_order: VecDeque, + live_updates: HashMap>>, + live_update_order: VecDeque, +} + +impl DiagnosticStoreState { + fn run_mut( + &mut self, + scope: &DiagnosticScope, + limits: DiagnosticStoreLimits, + ) -> Result<&mut DiagnosticRunState, DiagnosticStoreError> { + let session_key = DiagnosticSessionKey::from(scope); + if !self.sessions.contains_key(&session_key) { + while self.sessions.len() >= limits.max_sessions { + let Some(evicted) = self.session_order.pop_front() else { + return Err(DiagnosticStoreError::Invariant); + }; + self.sessions.remove(&evicted); + } + self.sessions + .insert(session_key.clone(), DiagnosticSessionState::default()); + } + touch(&mut self.session_order, session_key.clone()); + + let session = self + .sessions + .get_mut(&session_key) + .ok_or(DiagnosticStoreError::Invariant)?; + if !session.runs.contains_key(&scope.run_id) { + while session.runs.len() >= limits.max_runs_per_session { + let Some(evicted) = session.run_order.pop_front() else { + return Err(DiagnosticStoreError::Invariant); + }; + session.runs.remove(&evicted); + } + session + .runs + .insert(scope.run_id, DiagnosticRunState::default()); + } + touch(&mut session.run_order, scope.run_id); + session + .runs + .get_mut(&scope.run_id) + .ok_or(DiagnosticStoreError::Invariant) + } + + fn run(&self, scope: &DiagnosticScope) -> Option<&DiagnosticRunState> { + self.sessions + .get(&DiagnosticSessionKey::from(scope))? + .runs + .get(&scope.run_id) + } + + fn subscribe( + &mut self, + scope: DiagnosticScope, + capacity: usize, + max_scopes: usize, + ) -> Result>, DiagnosticStoreError> { + self.prune_inactive_live_updates(); + if !self.live_updates.contains_key(&scope) { + while self.live_updates.len() >= max_scopes { + let evicted = self + .live_update_order + .pop_front() + .ok_or(DiagnosticStoreError::Invariant)?; + self.live_updates + .remove(&evicted) + .ok_or(DiagnosticStoreError::Invariant)?; + } + let (sender, _) = broadcast::channel(capacity); + self.live_updates.insert(scope.clone(), sender); + } + touch(&mut self.live_update_order, scope.clone()); + self.live_updates + .get(&scope) + .map(broadcast::Sender::subscribe) + .ok_or(DiagnosticStoreError::Invariant) + } + + fn send_live_update(&mut self, envelope: DiagnosticUpdateEnvelope) { + let scope = envelope.scope.clone(); + let has_no_receivers = self + .live_updates + .get(&scope) + .is_some_and(|sender| sender.send(Arc::new(envelope)).is_err()); + if has_no_receivers { + self.live_updates.remove(&scope); + if let Some(index) = self + .live_update_order + .iter() + .position(|entry| entry == &scope) + { + self.live_update_order.remove(index); + } + } + } + + fn prune_inactive_live_updates(&mut self) { + self.live_updates + .retain(|_, sender| sender.receiver_count() > 0); + self.live_update_order + .retain(|scope| self.live_updates.contains_key(scope)); + } +} + +fn touch(order: &mut VecDeque, value: T) { + if let Some(index) = order.iter().position(|entry| entry == &value) { + order.remove(index); + } + order.push_back(value); +} + +#[derive(Debug)] +pub struct InMemoryDiagnosticStore { + limits: DiagnosticStoreLimits, + state: Mutex, +} + +impl InMemoryDiagnosticStore { + pub fn new(limits: DiagnosticStoreLimits) -> Result { + let limits = limits.validate()?; + Ok(Self { + limits, + state: Mutex::new(DiagnosticStoreState::default()), + }) + } + + pub fn record_prompt( + &self, + scope: DiagnosticScope, + 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, _| { + run.prompt = Some(prompt); + }) + } + + pub fn record_model_call( + &self, + scope: DiagnosticScope, + 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, _| { + replace_or_push( + &mut run.model_calls, + model_call, + cap, + |existing, incoming| existing.call_id == incoming.call_id, + ); + }) + } + + pub fn record_tool_execution( + &self, + scope: DiagnosticScope, + 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, _| { + replace_or_push(&mut run.tool_executions, tool, cap, |existing, incoming| { + existing.activity_id == incoming.activity_id + }); + }) + } + + pub fn record_activity( + &self, + scope: DiagnosticScope, + 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| { + push_bounded( + &mut run.activity, + DiagnosticActivityEntry { sequence, event }, + cap, + ); + }) + } + + pub fn record_stats( + &self, + scope: DiagnosticScope, + stats: SessionDiagnosticStats, + ) -> Result { + let stats = stats.into_bounded(); + self.record( + scope, + DiagnosticUpdateKind::Stats(stats.clone()), + move |run, _| run.stats = stats, + ) + } + + pub fn snapshot( + &self, + scope: &DiagnosticScope, + ) -> Result, DiagnosticStoreError> { + let state = self + .state + .lock() + .map_err(|_| DiagnosticStoreError::StateUnavailable)?; + let Some(run) = state.run(scope) else { + return Ok(None); + }; + Ok(Some(DiagnosticSnapshot { + scope: scope.clone(), + stream_id: run.stream_id, + prompt: run.prompt.clone(), + model_calls: run.model_calls.iter().cloned().collect(), + tool_executions: run.tool_executions.iter().cloned().collect(), + activity: run.activity.iter().cloned().collect(), + stats: run.stats.clone(), + latest_sequence: run.latest_sequence, + })) + } + + pub fn updates_after( + &self, + scope: &DiagnosticScope, + after: Option, + ) -> Result { + let state = self + .state + .lock() + .map_err(|_| DiagnosticStoreError::StateUnavailable)?; + let Some(run) = state.run(scope) else { + return Ok(DiagnosticUpdateBatch { + updates: Vec::new(), + retention_floor: None, + latest_cursor: None, + rebase_required: after.is_some(), + }); + }; + let retention_floor = run.updates.front().map(DiagnosticUpdateEnvelope::cursor); + let rebase_required = match (after, retention_floor) { + (Some(after), _) if after.stream_id != run.stream_id => true, + (Some(after), _) if after.sequence > run.latest_sequence => true, + (Some(after), Some(floor)) => { + after.sequence.as_u64().saturating_add(1) < floor.sequence.as_u64() + } + (None, Some(floor)) => floor.sequence.as_u64() > 1, + _ => false, + }; + let updates = run + .updates + .iter() + .filter(|update| { + after.is_none_or(|cursor| { + cursor.stream_id != run.stream_id || update.sequence > cursor.sequence + }) + }) + .cloned() + .collect(); + Ok(DiagnosticUpdateBatch { + updates, + retention_floor, + latest_cursor: Some(DiagnosticCursor::new(run.stream_id, run.latest_sequence)), + rebase_required, + }) + } + + pub fn subscribe( + &self, + scope: DiagnosticScope, + ) -> Result { + let mut state = self + .state + .lock() + .map_err(|_| DiagnosticStoreError::StateUnavailable)?; + let receiver = state.subscribe( + scope, + self.limits.live_update_capacity, + self.limits.max_live_update_scopes, + )?; + Ok(DiagnosticSubscription { receiver }) + } + + fn record( + &self, + scope: DiagnosticScope, + update: DiagnosticUpdateKind, + mutate: impl FnOnce(&mut DiagnosticRunState, DiagnosticSequence), + ) -> Result { + let mut state = self + .state + .lock() + .map_err(|_| DiagnosticStoreError::StateUnavailable)?; + let run = state.run_mut(&scope, self.limits)?; + let next = run + .latest_sequence + .as_u64() + .checked_add(1) + .ok_or(DiagnosticStoreError::SequenceExhausted)?; + let sequence = DiagnosticSequence::new(next); + mutate(run, sequence); + let cursor = DiagnosticCursor::new(run.stream_id, sequence); + let envelope = DiagnosticUpdateEnvelope { + scope, + stream_id: run.stream_id, + sequence, + emitted_at: Utc::now(), + update, + }; + run.latest_sequence = sequence; + push_bounded( + &mut run.updates, + envelope.clone(), + self.limits.max_updates_per_run, + ); + state.send_live_update(envelope); + Ok(cursor) + } +} + +impl Default for InMemoryDiagnosticStore { + fn default() -> Self { + let limits = DiagnosticStoreLimits::default(); + Self { + limits, + state: Mutex::new(DiagnosticStoreState::default()), + } + } +} + +fn push_bounded(values: &mut VecDeque, value: T, capacity: usize) { + while values.len() >= capacity { + values.pop_front(); + } + values.push_back(value); +} + +fn replace_or_push( + values: &mut VecDeque, + value: T, + capacity: usize, + matches: impl Fn(&T, &T) -> bool, +) { + if let Some(index) = values.iter().position(|existing| matches(existing, &value)) { + values.remove(index); + } + push_bounded(values, value, capacity); +} + +pub struct DiagnosticSubscription { + receiver: broadcast::Receiver>, +} + +impl DiagnosticSubscription { + pub async fn recv(&mut self) -> Result, DiagnosticStoreError> { + match self.receiver.recv().await { + Ok(update) => Ok(update), + Err(broadcast::error::RecvError::Lagged(count)) => { + Err(DiagnosticStoreError::SubscriberLagged(count)) + } + Err(broadcast::error::RecvError::Closed) => { + Err(DiagnosticStoreError::SubscriptionClosed) + } + } + } +} + +#[cfg(test)] +mod tests { + use std::panic::{AssertUnwindSafe, catch_unwind}; + use std::sync::Barrier; + + use chrono::Utc; + use ironclaw_host_api::{ + ids::{TenantId, ThreadId, UserId}, + turn::TurnRunId, + }; + use ironclaw_product_contracts::inspector::{ + BoundedDiagnosticText, DIAGNOSTIC_LABEL_MAX_BYTES, DIAGNOSTIC_SUMMARY_MAX_BYTES, + DiagnosticActivityEvent, DiagnosticActivityKind, DiagnosticModelCallId, + DiagnosticModelCount, DiagnosticScope, InspectorModelCallStatus, MAX_ACTIVE_SKILLS, + MAX_MODELS_IN_STATS, MAX_PROMPT_COMPONENTS, ModelCallDiagnostic, + PROMPT_COMPONENT_CONTENT_MAX_BYTES, PROMPT_COMPONENT_TOTAL_MAX_BYTES, + PromptComponentDiagnostic, PromptComponentKind, PromptDiagnostic, TOOL_ARGUMENTS_MAX_BYTES, + TOOL_RESULT_MAX_BYTES, ToolExecutionDiagnostic, ToolExecutionStatus, + }; + + use super::*; + + fn scope(tenant: &str, user: &str, thread: &str, run_id: TurnRunId) -> DiagnosticScope { + DiagnosticScope::new( + TenantId::new(tenant).expect("tenant"), + UserId::new(user).expect("user"), + ThreadId::new(thread).expect("thread"), + run_id, + ) + } + + fn activity(summary: &str) -> DiagnosticActivityEvent { + DiagnosticActivityEvent::new( + Utc::now(), + DiagnosticActivityKind::Progress, + None, + None, + None, + Some(summary.to_string()), + ) + } + + fn tiny_limits() -> DiagnosticStoreLimits { + DiagnosticStoreLimits { + max_sessions: 2, + max_runs_per_session: 2, + max_model_calls_per_run: 2, + max_tool_executions_per_run: 2, + max_activity_entries_per_run: 2, + max_updates_per_run: 2, + max_live_update_scopes: 2, + live_update_capacity: 8, + } + } + + #[test] + fn rejects_zero_limits() { + let mut limits = tiny_limits(); + limits.max_sessions = 0; + assert_eq!( + InMemoryDiagnosticStore::new(limits).expect_err("zero must fail"), + DiagnosticStoreError::InvalidLimit("max_sessions") + ); + } + + #[test] + fn zero_validation_covers_every_limit_field_with_its_stable_name() { + let valid = tiny_limits(); + let cases = [ + ( + "max_sessions", + DiagnosticStoreLimits { + max_sessions: 0, + ..valid + }, + ), + ( + "max_runs_per_session", + DiagnosticStoreLimits { + max_runs_per_session: 0, + ..valid + }, + ), + ( + "max_model_calls_per_run", + DiagnosticStoreLimits { + max_model_calls_per_run: 0, + ..valid + }, + ), + ( + "max_tool_executions_per_run", + DiagnosticStoreLimits { + max_tool_executions_per_run: 0, + ..valid + }, + ), + ( + "max_activity_entries_per_run", + DiagnosticStoreLimits { + max_activity_entries_per_run: 0, + ..valid + }, + ), + ( + "max_updates_per_run", + DiagnosticStoreLimits { + max_updates_per_run: 0, + ..valid + }, + ), + ( + "max_live_update_scopes", + DiagnosticStoreLimits { + max_live_update_scopes: 0, + ..valid + }, + ), + ( + "live_update_capacity", + DiagnosticStoreLimits { + live_update_capacity: 0, + ..valid + }, + ), + ]; + + for (name, limits) in cases { + assert_eq!( + limits.validate(), + Err(DiagnosticStoreError::InvalidLimit(name)) + ); + } + } + + #[test] + fn rejects_limits_above_the_hard_ceiling() { + let mut limits = tiny_limits(); + limits.max_sessions = DEFAULT_MAX_TRACKED_SESSIONS + 1; + assert_eq!( + InMemoryDiagnosticStore::new(limits).expect_err("oversized limit must fail"), + DiagnosticStoreError::LimitExceedsMaximum { + name: "max_sessions", + maximum: DEFAULT_MAX_TRACKED_SESSIONS, + } + ); + } + + #[test] + fn default_store_limits_pass_the_validated_constructor_contract() { + let limits = DiagnosticStoreLimits::default(); + + assert_eq!(limits.validate(), Ok(limits)); + assert_eq!(InMemoryDiagnosticStore::default().limits, limits); + } + + #[test] + fn exact_scope_keys_prevent_cross_scope_reads() { + let store = InMemoryDiagnosticStore::new(tiny_limits()).expect("store"); + let run_id = TurnRunId::new(); + let allowed = scope("tenant-a", "user-a", "thread-a", run_id); + store + .record_activity(allowed.clone(), activity("allowed")) + .expect("record"); + + for denied in [ + scope("tenant-b", "user-a", "thread-a", run_id), + scope("tenant-a", "user-b", "thread-a", run_id), + scope("tenant-a", "user-a", "thread-b", run_id), + scope("tenant-a", "user-a", "thread-a", TurnRunId::new()), + ] { + assert!(store.snapshot(&denied).expect("snapshot").is_none()); + } + assert!(store.snapshot(&allowed).expect("snapshot").is_some()); + } + + #[test] + fn prompt_model_and_tool_records_share_one_scoped_snapshot() { + let store = InMemoryDiagnosticStore::new(tiny_limits()).expect("store"); + let scope = scope("tenant", "user", "thread", TurnRunId::new()); + let prompt = PromptDiagnostic::new( + Utc::now(), + Vec::new(), + "system prompt", + Some(10), + 1, + 0, + 1, + Vec::new(), + 1, + Some("requested-model".to_string()), + Some("effective-model".to_string()), + Some(100_000), + ); + let model_call = ModelCallDiagnostic::new( + DiagnosticModelCallId::new(), + 1, + "requested-model", + Some("effective-model".to_string()), + Utc::now(), + Some(Utc::now()), + Some(25), + InspectorModelCallStatus::Succeeded, + None, + None, + ); + let tool = ToolExecutionDiagnostic::new( + ironclaw_host_api::turn::CapabilityActivityId::new(), + Some(model_call.call_id), + "filesystem.read", + Some("{}".to_string()), + Some("result".to_string()), + ToolExecutionStatus::Succeeded, + Some(5), + None, + None, + None, + ); + store.record_prompt(scope.clone(), prompt).expect("prompt"); + store + .record_model_call(scope.clone(), model_call) + .expect("model call"); + store + .record_tool_execution(scope.clone(), tool) + .expect("tool"); + + let snapshot = store.snapshot(&scope).expect("snapshot").expect("present"); + assert!(snapshot.prompt.is_some()); + assert_eq!(snapshot.model_calls.len(), 1); + assert_eq!(snapshot.tool_executions.len(), 1); + assert_eq!(snapshot.latest_sequence, DiagnosticSequence::new(3)); + } + + #[test] + fn record_prompt_reapplies_limits_to_a_literal_dto() { + let store = InMemoryDiagnosticStore::new(tiny_limits()).expect("store"); + let scope = scope("tenant", "user", "thread", TurnRunId::new()); + let oversized_label = "l".repeat(DIAGNOSTIC_LABEL_MAX_BYTES + 1); + let oversized_content = "x".repeat(PROMPT_COMPONENT_CONTENT_MAX_BYTES + 1); + let prompt = PromptDiagnostic { + captured_at: Utc::now(), + components: (0..=MAX_PROMPT_COMPONENTS) + .map(|index| PromptComponentDiagnostic { + kind: PromptComponentKind::Instruction, + label: BoundedDiagnosticText::reconstructed_prompt(oversized_label.clone()), + content: BoundedDiagnosticText::reconstructed_prompt(if index < 5 { + oversized_content.clone() + } else { + "x".to_string() + }), + estimated_tokens: None, + }) + .collect(), + components_truncated: false, + reconstructed_prompt: BoundedDiagnosticText::reconstructed_prompt("prompt"), + total_estimated_tokens: None, + message_count: 1, + identity_message_count: 0, + instruction_snippet_count: 1, + active_skills: (0..=MAX_ACTIVE_SKILLS) + .map(|_| BoundedDiagnosticText::reconstructed_prompt(oversized_label.clone())) + .collect(), + active_skills_truncated: false, + capability_count: 0, + requested_model: Some(BoundedDiagnosticText::reconstructed_prompt( + oversized_label.clone(), + )), + effective_model: Some(BoundedDiagnosticText::reconstructed_prompt(oversized_label)), + context_limit: None, + }; + + store.record_prompt(scope.clone(), prompt).expect("prompt"); + + let prompt = store + .snapshot(&scope) + .expect("snapshot") + .expect("run") + .prompt + .expect("prompt"); + assert!(prompt.components_truncated); + assert!(prompt.components.len() <= MAX_PROMPT_COMPONENTS); + assert!( + prompt + .components + .iter() + .all(|component| component.label.content().len() <= DIAGNOSTIC_LABEL_MAX_BYTES) + ); + assert!(prompt.components.iter().all( + |component| component.content.content().len() <= PROMPT_COMPONENT_CONTENT_MAX_BYTES + )); + assert!( + prompt + .components + .iter() + .map(|component| component.content.content().len()) + .sum::() + <= PROMPT_COMPONENT_TOTAL_MAX_BYTES + ); + assert!(prompt.active_skills_truncated); + assert_eq!(prompt.active_skills.len(), MAX_ACTIVE_SKILLS); + assert!( + prompt + .active_skills + .iter() + .all(|skill| skill.content().len() <= DIAGNOSTIC_LABEL_MAX_BYTES) + ); + assert!( + prompt + .requested_model + .as_ref() + .is_some_and(|model| model.content().len() <= DIAGNOSTIC_LABEL_MAX_BYTES) + ); + assert!( + prompt + .effective_model + .as_ref() + .is_some_and(|model| model.content().len() <= DIAGNOSTIC_LABEL_MAX_BYTES) + ); + } + + #[test] + fn record_boundaries_reapply_limits_to_literal_dtos() { + let mut limits = tiny_limits(); + limits.max_updates_per_run = 8; + let store = InMemoryDiagnosticStore::new(limits).expect("store"); + let scope = scope("tenant", "user", "thread", TurnRunId::new()); + let model_call_id = DiagnosticModelCallId::new(); + let activity_id = ironclaw_host_api::turn::CapabilityActivityId::new(); + let oversized_label = + BoundedDiagnosticText::reconstructed_prompt("l".repeat(DIAGNOSTIC_LABEL_MAX_BYTES + 1)); + let oversized_summary = BoundedDiagnosticText::reconstructed_prompt( + "s".repeat(DIAGNOSTIC_SUMMARY_MAX_BYTES + 1), + ); + let oversized_tool_text = + BoundedDiagnosticText::reconstructed_prompt("x".repeat(TOOL_ARGUMENTS_MAX_BYTES + 1)); + + store + .record_model_call( + scope.clone(), + ModelCallDiagnostic { + call_id: model_call_id, + iteration: 1, + requested_model: oversized_label.clone(), + effective_model: Some(oversized_label.clone()), + started_at: Utc::now(), + completed_at: None, + duration_ms: None, + status: InspectorModelCallStatus::Failed, + usage: None, + failure_summary: Some(oversized_summary.clone()), + }, + ) + .expect("model call"); + store + .record_tool_execution( + scope.clone(), + ToolExecutionDiagnostic { + activity_id, + model_call_id: Some(model_call_id), + capability_name: oversized_label.clone(), + arguments: Some(oversized_tool_text.clone()), + result: Some(oversized_tool_text.clone()), + status: ToolExecutionStatus::Failed, + duration_ms: None, + output_bytes: Some(1), + failure_category: Some(oversized_label.clone()), + failure_summary: Some(oversized_summary.clone()), + }, + ) + .expect("tool execution"); + store + .record_activity( + scope.clone(), + DiagnosticActivityEvent { + occurred_at: Utc::now(), + kind: DiagnosticActivityKind::ToolFailed, + iteration: Some(1), + activity_id: Some(activity_id), + model_call_id: Some(model_call_id), + summary: Some(oversized_summary), + }, + ) + .expect("activity"); + store + .record_stats( + scope.clone(), + SessionDiagnosticStats { + calls_per_model: (0..=MAX_MODELS_IN_STATS) + .map(|_| DiagnosticModelCount { + model: oversized_label.clone(), + calls: 1, + }) + .collect(), + ..SessionDiagnosticStats::default() + }, + ) + .expect("stats"); + + let snapshot = store.snapshot(&scope).expect("snapshot").expect("run"); + let model_call = snapshot.model_calls.first().expect("model call"); + assert!(model_call.requested_model.content().len() <= DIAGNOSTIC_LABEL_MAX_BYTES); + assert!( + model_call + .effective_model + .as_ref() + .is_some_and(|model| model.content().len() <= DIAGNOSTIC_LABEL_MAX_BYTES) + ); + assert!( + model_call + .failure_summary + .as_ref() + .is_some_and(|summary| summary.content().len() <= DIAGNOSTIC_SUMMARY_MAX_BYTES) + ); + + let tool = snapshot.tool_executions.first().expect("tool execution"); + assert!(tool.capability_name.content().len() <= DIAGNOSTIC_LABEL_MAX_BYTES); + assert!( + tool.arguments + .as_ref() + .is_some_and(|arguments| arguments.content().len() <= TOOL_ARGUMENTS_MAX_BYTES) + ); + assert!( + tool.result + .as_ref() + .is_some_and(|result| result.content().len() <= TOOL_RESULT_MAX_BYTES) + ); + assert_eq!( + tool.output_bytes, + tool.result + .as_ref() + .map(BoundedDiagnosticText::original_bytes) + ); + assert!( + tool.failure_category + .as_ref() + .is_some_and(|category| category.content().len() <= DIAGNOSTIC_LABEL_MAX_BYTES) + ); + assert!( + tool.failure_summary + .as_ref() + .is_some_and(|summary| summary.content().len() <= DIAGNOSTIC_SUMMARY_MAX_BYTES) + ); + + let event = &snapshot.activity.first().expect("activity").event; + assert!( + event + .summary + .as_ref() + .is_some_and(|summary| summary.content().len() <= DIAGNOSTIC_SUMMARY_MAX_BYTES) + ); + assert!(snapshot.stats.calls_per_model_truncated); + assert_eq!(snapshot.stats.calls_per_model.len(), MAX_MODELS_IN_STATS); + assert!( + snapshot + .stats + .calls_per_model + .iter() + .all(|count| count.model.content().len() <= DIAGNOSTIC_LABEL_MAX_BYTES) + ); + } + + #[test] + fn session_eviction_is_deterministic_and_write_lru() { + let store = InMemoryDiagnosticStore::new(tiny_limits()).expect("store"); + let first = scope("tenant", "user-1", "thread", TurnRunId::new()); + let second = scope("tenant", "user-2", "thread", TurnRunId::new()); + let third = scope("tenant", "user-3", "thread", TurnRunId::new()); + store + .record_activity(first.clone(), activity("first")) + .expect("record first"); + store + .record_activity(second.clone(), activity("second")) + .expect("record second"); + store + .record_activity(first.clone(), activity("touch first")) + .expect("touch first"); + store + .record_activity(third.clone(), activity("third")) + .expect("record third"); + + assert!(store.snapshot(&first).expect("first").is_some()); + assert!(store.snapshot(&second).expect("second").is_none()); + assert!(store.snapshot(&third).expect("third").is_some()); + } + + #[test] + fn run_eviction_is_deterministic_and_write_lru() { + let store = InMemoryDiagnosticStore::new(tiny_limits()).expect("store"); + let first = scope("tenant", "user", "thread", TurnRunId::new()); + let second = scope("tenant", "user", "thread", TurnRunId::new()); + let third = scope("tenant", "user", "thread", TurnRunId::new()); + store + .record_activity(first.clone(), activity("first")) + .expect("first"); + store + .record_activity(second.clone(), activity("second")) + .expect("second"); + store + .record_activity(first.clone(), activity("touch")) + .expect("touch"); + store + .record_activity(third.clone(), activity("third")) + .expect("third"); + + assert!(store.snapshot(&first).expect("first").is_some()); + assert!(store.snapshot(&second).expect("second").is_none()); + assert!(store.snapshot(&third).expect("third").is_some()); + } + + #[test] + fn activity_and_update_history_are_bounded_and_ordered() { + let store = InMemoryDiagnosticStore::new(tiny_limits()).expect("store"); + let scope = scope("tenant", "user", "thread", TurnRunId::new()); + let mut stream_id = None; + for index in 1..=3 { + let cursor = store + .record_activity(scope.clone(), activity(&format!("entry-{index}"))) + .expect("record"); + stream_id = Some(cursor.stream_id); + assert_eq!(cursor.sequence.as_u64(), index); + } + let stream_id = stream_id.expect("stream id"); + + let snapshot = store.snapshot(&scope).expect("snapshot").expect("present"); + let sequences: Vec<_> = snapshot + .activity + .iter() + .map(|entry| entry.sequence.as_u64()) + .collect(); + assert_eq!(sequences, vec![2, 3]); + + let batch = store + .updates_after( + &scope, + Some(DiagnosticCursor::new(stream_id, DiagnosticSequence::ZERO)), + ) + .expect("updates"); + assert!(batch.rebase_required); + assert_eq!( + batch + .updates + .iter() + .map(|update| update.sequence.as_u64()) + .collect::>(), + vec![2, 3] + ); + assert_eq!( + batch.retention_floor, + Some(DiagnosticCursor::new(stream_id, DiagnosticSequence::new(2),)) + ); + assert_eq!( + batch.latest_cursor, + Some(DiagnosticCursor::new(stream_id, DiagnosticSequence::new(3),)) + ); + } + + #[test] + fn fresh_reader_requires_rebase_after_the_stream_prefix_is_evicted() { + let store = InMemoryDiagnosticStore::new(tiny_limits()).expect("store"); + let scope = scope("tenant", "user", "thread", TurnRunId::new()); + for index in 1..=2 { + store + .record_activity(scope.clone(), activity(&format!("entry-{index}"))) + .expect("record"); + } + let complete = store.updates_after(&scope, None).expect("complete updates"); + assert!(!complete.rebase_required); + + store + .record_activity(scope.clone(), activity("entry-3")) + .expect("record"); + + let batch = store.updates_after(&scope, None).expect("updates"); + + assert!(batch.rebase_required); + assert_eq!( + batch + .updates + .iter() + .map(|update| update.sequence.as_u64()) + .collect::>(), + vec![2, 3] + ); + } + + #[test] + fn recreating_an_evicted_run_changes_the_stream_generation() { + let store = InMemoryDiagnosticStore::new(tiny_limits()).expect("store"); + let first = scope("tenant", "user", "thread", TurnRunId::new()); + let second = scope("tenant", "user", "thread", TurnRunId::new()); + let third = scope("tenant", "user", "thread", TurnRunId::new()); + let original = store + .record_activity(first.clone(), activity("first")) + .expect("first"); + store + .record_activity(second, activity("second")) + .expect("second"); + store + .record_activity(third, activity("third")) + .expect("third"); + let recreated = store + .record_activity(first.clone(), activity("recreated")) + .expect("recreated"); + assert_ne!(original.stream_id, recreated.stream_id); + + let batch = store + .updates_after(&first, Some(original)) + .expect("rebase batch"); + assert!(batch.rebase_required); + assert_eq!(batch.updates.len(), 1); + assert_eq!(batch.latest_cursor, Some(recreated)); + } + + #[test] + fn cursor_for_an_evicted_run_requires_rebase() { + let store = InMemoryDiagnosticStore::new(tiny_limits()).expect("store"); + let first = scope("tenant", "user", "thread", TurnRunId::new()); + let second = scope("tenant", "user", "thread", TurnRunId::new()); + let third = scope("tenant", "user", "thread", TurnRunId::new()); + let original = store + .record_activity(first.clone(), activity("first")) + .expect("first"); + store + .record_activity(second, activity("second")) + .expect("second"); + store + .record_activity(third, activity("third")) + .expect("third"); + + let batch = store + .updates_after(&first, Some(original)) + .expect("evicted run batch"); + + assert!(batch.rebase_required); + assert!(batch.updates.is_empty()); + assert_eq!(batch.retention_floor, None); + assert_eq!(batch.latest_cursor, None); + } + + #[test] + fn future_cursor_in_the_current_stream_requires_rebase() { + let store = InMemoryDiagnosticStore::new(tiny_limits()).expect("store"); + let scope = scope("tenant", "user", "thread", TurnRunId::new()); + let latest = store + .record_activity(scope.clone(), activity("first")) + .expect("record"); + let future = DiagnosticCursor::new( + latest.stream_id, + DiagnosticSequence::new(latest.sequence.as_u64() + 1), + ); + + let batch = store + .updates_after(&scope, Some(future)) + .expect("future cursor batch"); + + assert!(batch.rebase_required); + assert!(batch.updates.is_empty()); + assert_eq!(batch.latest_cursor, Some(latest)); + } + + #[tokio::test] + async fn unrelated_scope_saturation_does_not_lag_scoped_subscription() { + let mut limits = tiny_limits(); + limits.live_update_capacity = 2; + let store = InMemoryDiagnosticStore::new(limits).expect("store"); + let allowed = scope("tenant", "user", "thread-a", TurnRunId::new()); + let other = scope("tenant", "user", "thread-b", TurnRunId::new()); + let mut subscription = store + .subscribe(allowed.clone()) + .expect("scoped subscription"); + let mut other_subscription = store + .subscribe(other.clone()) + .expect("other scoped subscription"); + for index in 0..=limits.live_update_capacity { + store + .record_activity(other.clone(), activity(&format!("other-{index}"))) + .expect("other"); + } + store + .record_activity(allowed.clone(), activity("allowed")) + .expect("allowed"); + + assert!(matches!( + other_subscription.recv().await, + Err(DiagnosticStoreError::SubscriberLagged(_)) + )); + let update = subscription.recv().await.expect("matching update"); + assert_eq!(update.scope, allowed); + assert_eq!(update.sequence, DiagnosticSequence::new(1)); + } + + #[tokio::test] + async fn live_scope_limit_evicts_the_least_recently_subscribed_scope() { + let mut limits = tiny_limits(); + limits.max_live_update_scopes = 2; + let store = InMemoryDiagnosticStore::new(limits).expect("store"); + let first_scope = scope("tenant", "user", "thread-a", TurnRunId::new()); + let second_scope = scope("tenant", "user", "thread-b", TurnRunId::new()); + let third_scope = scope("tenant", "user", "thread-c", TurnRunId::new()); + let mut first = store + .subscribe(first_scope.clone()) + .expect("first subscription"); + let mut second = store.subscribe(second_scope).expect("second subscription"); + let _refreshed_first = store + .subscribe(first_scope.clone()) + .expect("refresh first subscription"); + let mut third = store + .subscribe(third_scope.clone()) + .expect("third subscription"); + + assert_eq!( + second.recv().await, + Err(DiagnosticStoreError::SubscriptionClosed) + ); + store + .record_activity(first_scope.clone(), activity("first")) + .expect("record first"); + store + .record_activity(third_scope.clone(), activity("third")) + .expect("record third"); + + assert_eq!(first.recv().await.expect("first update").scope, first_scope); + assert_eq!(third.recv().await.expect("third update").scope, third_scope); + assert_eq!( + store.state.lock().expect("state").live_updates.len(), + limits.max_live_update_scopes + ); + } + + #[tokio::test] + async fn concurrent_writers_deliver_live_updates_in_sequence_order() { + const WRITER_COUNT: usize = 32; + + let mut limits = tiny_limits(); + limits.live_update_capacity = WRITER_COUNT; + let store = Arc::new(InMemoryDiagnosticStore::new(limits).expect("store")); + let scope = scope("tenant", "user", "thread", TurnRunId::new()); + let mut subscription = store.subscribe(scope.clone()).expect("scoped subscription"); + let barrier = Arc::new(Barrier::new(WRITER_COUNT)); + + std::thread::scope(|threads| { + let handles = (0..WRITER_COUNT) + .map(|index| { + let store = Arc::clone(&store); + let scope = scope.clone(); + let barrier = Arc::clone(&barrier); + threads.spawn(move || { + barrier.wait(); + store + .record_activity(scope, activity(&format!("writer-{index}"))) + .expect("record") + }) + }) + .collect::>(); + + for handle in handles { + handle.join().expect("writer thread"); + } + }); + + for expected in 1..=WRITER_COUNT as u64 { + let update = subscription.recv().await.expect("ordered update"); + assert_eq!(update.sequence, DiagnosticSequence::new(expected)); + } + } + + #[tokio::test] + async fn subscription_preserves_sequence_for_its_scope() { + let store = InMemoryDiagnosticStore::new(tiny_limits()).expect("store"); + let allowed = scope("tenant", "user", "thread-a", TurnRunId::new()); + let mut subscription = store + .subscribe(allowed.clone()) + .expect("scoped subscription"); + store + .record_activity(allowed.clone(), activity("allowed")) + .expect("allowed"); + let update = subscription.recv().await.expect("matching update"); + assert_eq!(update.scope, allowed); + assert_eq!(update.sequence, DiagnosticSequence::new(1)); + } + + #[tokio::test] + async fn stats_are_visible_when_the_update_is_delivered() { + let store = InMemoryDiagnosticStore::new(tiny_limits()).expect("store"); + let scope = scope("tenant", "user", "thread", TurnRunId::new()); + let mut subscription = store.subscribe(scope.clone()).expect("scoped subscription"); + let stats = SessionDiagnosticStats { + total_model_calls: 7, + ..SessionDiagnosticStats::default() + }; + store + .record_stats(scope.clone(), stats) + .expect("record stats"); + let update = subscription.recv().await.expect("stats update"); + assert!(matches!(update.update, DiagnosticUpdateKind::Stats(_))); + let snapshot = store.snapshot(&scope).expect("snapshot").expect("present"); + assert_eq!(snapshot.stats.total_model_calls, 7); + } + + #[test] + fn stats_remain_scoped_to_the_run_that_published_them() { + let store = InMemoryDiagnosticStore::new(tiny_limits()).expect("store"); + let first = scope("tenant", "user", "thread", TurnRunId::new()); + let second = scope("tenant", "user", "thread", TurnRunId::new()); + store + .record_stats( + first.clone(), + SessionDiagnosticStats { + total_model_calls: 3, + ..SessionDiagnosticStats::default() + }, + ) + .expect("first stats"); + store + .record_stats( + second.clone(), + SessionDiagnosticStats { + total_model_calls: 8, + ..SessionDiagnosticStats::default() + }, + ) + .expect("second stats"); + + assert_eq!( + store + .snapshot(&first) + .expect("first snapshot") + .expect("first run") + .stats + .total_model_calls, + 3 + ); + assert_eq!( + store + .snapshot(&second) + .expect("second snapshot") + .expect("second run") + .stats + .total_model_calls, + 8 + ); + } + + #[test] + fn poisoned_state_returns_a_redacted_error() { + let store = InMemoryDiagnosticStore::new(tiny_limits()).expect("store"); + let _ = catch_unwind(AssertUnwindSafe(|| { + let _guard = store.state.lock().expect("lock before poison"); + panic!("poison store lock"); + })); + let error = store + .snapshot(&scope("tenant", "user", "thread", TurnRunId::new())) + .expect_err("poison must fail closed"); + assert_eq!(error, DiagnosticStoreError::StateUnavailable); + assert!(!error.to_string().contains("poison")); + } +} diff --git a/crates/product/ironclaw_assistant/src/lib.rs b/crates/product/ironclaw_assistant/src/lib.rs index 28a825049a3..037e0d5ea72 100644 --- a/crates/product/ironclaw_assistant/src/lib.rs +++ b/crates/product/ironclaw_assistant/src/lib.rs @@ -51,6 +51,7 @@ mod filesystem_ledger; mod gate_state; mod in_memory_ledger; mod inbound_turn; +pub mod inspector_store; mod ledger; mod lifecycle; mod outbound_delivery;