From 8b4e40a595c8c6ad247fd3aca60cb52d1c4f10ae Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Wed, 5 Aug 2026 18:59:59 +0800 Subject: [PATCH 01/11] feat(inspector): add bounded diagnostic session store --- .../ironclaw_product/src/inspector_store.rs | 852 ++++++++++++++++++ crates/ironclaw_product/src/lib.rs | 1 + .../src/inspector.rs | 742 +++++++++++++++ crates/ironclaw_product_contracts/src/lib.rs | 1 + 4 files changed, 1596 insertions(+) create mode 100644 crates/ironclaw_product/src/inspector_store.rs create mode 100644 crates/ironclaw_product_contracts/src/inspector.rs diff --git a/crates/ironclaw_product/src/inspector_store.rs b/crates/ironclaw_product/src/inspector_store.rs new file mode 100644 index 00000000000..02850ae596b --- /dev/null +++ b/crates/ironclaw_product/src/inspector_store.rs @@ -0,0 +1,852 @@ +//! 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_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 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, + live_update_capacity: DEFAULT_MAX_RETAINED_UPDATES_PER_RUN, + } + } +} + +impl DiagnosticStoreLimits { + fn validate(self) -> Result { + let values = [ + ( + "max_sessions", + self.max_sessions, + DEFAULT_MAX_TRACKED_SESSIONS, + ), + ( + "max_runs_per_session", + self.max_runs_per_session, + DEFAULT_MAX_RETAINED_RUNS_PER_SESSION, + ), + ( + "max_model_calls_per_run", + self.max_model_calls_per_run, + DEFAULT_MAX_MODEL_CALLS_PER_RUN, + ), + ( + "max_tool_executions_per_run", + self.max_tool_executions_per_run, + DEFAULT_MAX_TOOL_EXECUTIONS_PER_RUN, + ), + ( + "max_activity_entries_per_run", + self.max_activity_entries_per_run, + DEFAULT_MAX_ACTIVITY_ENTRIES, + ), + ( + "max_updates_per_run", + self.max_updates_per_run, + DEFAULT_MAX_RETAINED_UPDATES_PER_RUN, + ), + ( + "live_update_capacity", + self.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, + 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(), + updates: VecDeque::new(), + latest_sequence: DiagnosticSequence::ZERO, + } + } +} + +#[derive(Debug, Default)] +struct DiagnosticSessionState { + runs: HashMap, + run_order: VecDeque, + stats: SessionDiagnosticStats, +} + +#[derive(Debug, Default)] +struct DiagnosticStoreState { + sessions: HashMap, + session_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 session(&self, scope: &DiagnosticScope) -> Option<&DiagnosticSessionState> { + self.sessions.get(&DiagnosticSessionKey::from(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, + updates: broadcast::Sender>, +} + +impl InMemoryDiagnosticStore { + pub fn new(limits: DiagnosticStoreLimits) -> Result { + let limits = limits.validate()?; + let (updates, _) = broadcast::channel(limits.live_update_capacity); + Ok(Self { + limits, + state: Mutex::new(DiagnosticStoreState::default()), + updates, + }) + } + + pub fn record_prompt( + &self, + scope: DiagnosticScope, + prompt: PromptDiagnostic, + ) -> Result { + let update = DiagnosticUpdateKind::PromptUpdated { + component_count: prompt.components.len(), + 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 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 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 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(); + let session_key = DiagnosticSessionKey::from(&scope); + 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); + let cursor = DiagnosticCursor::new(run.stream_id, sequence); + let envelope = DiagnosticUpdateEnvelope { + scope, + stream_id: run.stream_id, + sequence, + emitted_at: Utc::now(), + update: DiagnosticUpdateKind::Stats(stats.clone()), + }; + run.latest_sequence = sequence; + push_bounded( + &mut run.updates, + envelope.clone(), + self.limits.max_updates_per_run, + ); + let session = state + .sessions + .get_mut(&session_key) + .ok_or(DiagnosticStoreError::Invariant)?; + session.stats = stats; + let _ = self.updates.send(Arc::new(envelope)); + Ok(cursor) + } + + 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); + }; + let stats = state + .session(scope) + .map(|session| session.stats.clone()) + .unwrap_or_default(); + 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, + 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: false, + }); + }; + 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), Some(floor)) => { + after.sequence.as_u64().saturating_add(1) < floor.sequence.as_u64() + } + _ => 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) -> DiagnosticSubscription { + DiagnosticSubscription { + scope, + receiver: self.updates.subscribe(), + } + } + + 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, + ); + let _ = self.updates.send(Arc::new(envelope)); + Ok(cursor) + } +} + +impl Default for InMemoryDiagnosticStore { + fn default() -> Self { + let limits = DiagnosticStoreLimits::default(); + let (updates, _) = broadcast::channel(limits.live_update_capacity); + Self { + limits, + state: Mutex::new(DiagnosticStoreState::default()), + updates, + } + } +} + +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 { + scope: DiagnosticScope, + receiver: broadcast::Receiver>, +} + +impl DiagnosticSubscription { + pub async fn recv(&mut self) -> Result, DiagnosticStoreError> { + loop { + match self.receiver.recv().await { + Ok(update) if update.scope == self.scope => return Ok(update), + Ok(_) => continue, + Err(broadcast::error::RecvError::Lagged(count)) => { + return Err(DiagnosticStoreError::SubscriberLagged(count)); + } + Err(broadcast::error::RecvError::Closed) => { + return Err(DiagnosticStoreError::SubscriptionClosed); + } + } + } + } +} + +#[cfg(test)] +mod tests { + use std::panic::{AssertUnwindSafe, catch_unwind}; + + use chrono::Utc; + use ironclaw_host_api::{ + ids::{TenantId, ThreadId, UserId}, + turn::TurnRunId, + }; + use ironclaw_product_contracts::inspector::{ + DiagnosticActivityEvent, DiagnosticActivityKind, DiagnosticModelCallId, DiagnosticScope, + InspectorModelCallStatus, ModelCallDiagnostic, PromptDiagnostic, 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, + 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 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 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 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 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)); + } + + #[tokio::test] + async fn subscription_filters_scope_and_preserves_sequence() { + let store = InMemoryDiagnosticStore::new(tiny_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()); + store + .record_activity(other, activity("other")) + .expect("other"); + store + .record_activity(allowed, activity("allowed")) + .expect("allowed"); + let update = subscription.recv().await.expect("matching update"); + 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()); + 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 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/ironclaw_product/src/lib.rs b/crates/ironclaw_product/src/lib.rs index da6fbfc82c8..7209ff6b318 100644 --- a/crates/ironclaw_product/src/lib.rs +++ b/crates/ironclaw_product/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; diff --git a/crates/ironclaw_product_contracts/src/inspector.rs b/crates/ironclaw_product_contracts/src/inspector.rs new file mode 100644 index 00000000000..eb428d24a27 --- /dev/null +++ b/crates/ironclaw_product_contracts/src/inspector.rs @@ -0,0 +1,742 @@ +//! Operator-only diagnostic vocabulary for the Web Debug Inspector. +//! +//! These output-only DTOs cross the product boundary through the dedicated +//! inspection surface. 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::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_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)] +#[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)] +#[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)] +#[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)] +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, +} + +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 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, + } + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] +#[serde(rename_all = "snake_case")] +pub enum PromptComponentKind { + System, + Identity, + Instruction, + Skill, + Capability, + Conversation, + Other, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +pub struct PromptComponentDiagnostic { + pub kind: PromptComponentKind, + pub label: BoundedDiagnosticText, + 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 original_component_count = components.len(); + let mut remaining = PROMPT_COMPONENT_TOTAL_MAX_BYTES; + let mut bounded_components = + Vec::with_capacity(original_component_count.min(MAX_PROMPT_COMPONENTS)); + let mut components_truncated = original_component_count > MAX_PROMPT_COMPONENTS; + for mut component in components.into_iter().take(MAX_PROMPT_COMPONENTS) { + if remaining == 0 { + components_truncated = true; + break; + } + let retained = component.content.content().len().min(remaining); + component.content = component.content.rebound(retained); + components_truncated |= component.content.truncated(); + remaining = remaining.saturating_sub(component.content.content().len()); + bounded_components.push(component); + } + + 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: bounded_components, + components_truncated, + 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, + } + } + + 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)] +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)] +#[serde(rename_all = "snake_case")] +pub enum InspectorModelCallStatus { + Started, + Succeeded, + Failed, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +pub struct ModelCallDiagnostic { + pub call_id: DiagnosticModelCallId, + pub iteration: u32, + pub requested_model: BoundedDiagnosticText, + pub effective_model: Option, + pub started_at: DateTime, + pub completed_at: Option>, + pub duration_ms: Option, + pub status: InspectorModelCallStatus, + pub usage: Option, + 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), + } + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] +#[serde(rename_all = "snake_case")] +pub enum ToolExecutionStatus { + Started, + Succeeded, + Failed, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +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 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 result_truncated(&self) -> bool { + self.result + .as_ref() + .is_some_and(BoundedDiagnosticText::truncated) + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] +#[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)] +pub struct DiagnosticActivityEvent { + pub occurred_at: DateTime, + pub kind: DiagnosticActivityKind, + pub iteration: Option, + pub activity_id: Option, + pub model_call_id: Option, + 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), + } + } +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +pub struct DiagnosticActivityEntry { + pub sequence: DiagnosticSequence, + pub event: DiagnosticActivityEvent, +} + +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)] +pub struct DiagnosticMetricTotal { + pub known_total: u64, + pub unavailable_samples: u64, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +pub struct DiagnosticModelCount { + pub model: BoundedDiagnosticText, + pub calls: u64, +} + +impl DiagnosticModelCount { + pub fn new(model: impl Into, calls: u64) -> Self { + Self { + model: BoundedDiagnosticText::label(model), + calls, + } + } +} + +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)] +pub struct SessionDiagnosticStats { + pub total_model_calls: u64, + 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 + } +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +#[serde(tag = "type", content = "data", rename_all = "snake_case")] +pub enum DiagnosticUpdateKind { + PromptUpdated { + component_count: usize, + total_estimated_tokens: Option, + truncated: bool, + }, + ModelCall(ModelCallDiagnostic), + ToolExecutionUpdated { + activity_id: CapabilityActivityId, + model_call_id: Option, + 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 super::*; + + #[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!(std::str::from_utf8(bounded.content().as_bytes()).is_ok()); + assert_eq!(bounded.original_bytes(), value.len() as u64); + } + + #[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 tool_result_uses_the_fifty_kibibyte_contract() { + 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((TOOL_RESULT_MAX_BYTES + 1) as u64), + 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/ironclaw_product_contracts/src/lib.rs b/crates/ironclaw_product_contracts/src/lib.rs index effca45b1f9..2e6e1c38878 100644 --- a/crates/ironclaw_product_contracts/src/lib.rs +++ b/crates/ironclaw_product_contracts/src/lib.rs @@ -44,6 +44,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; From e1f0c11d72cf842db4369b977095a3085864207b Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Wed, 5 Aug 2026 19:57:14 +0800 Subject: [PATCH 02/11] fix(inspector): address diagnostics store review feedback --- .../ironclaw_product/src/inspector_store.rs | 287 +++++++++++++----- .../src/inspector.rs | 133 +++++++- 2 files changed, 331 insertions(+), 89 deletions(-) diff --git a/crates/ironclaw_product/src/inspector_store.rs b/crates/ironclaw_product/src/inspector_store.rs index 02850ae596b..61a7fddc559 100644 --- a/crates/ironclaw_product/src/inspector_store.rs +++ b/crates/ironclaw_product/src/inspector_store.rs @@ -50,40 +50,47 @@ impl Default for DiagnosticStoreLimits { 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, + live_update_capacity, + } = self; let values = [ - ( - "max_sessions", - self.max_sessions, - DEFAULT_MAX_TRACKED_SESSIONS, - ), + ("max_sessions", max_sessions, DEFAULT_MAX_TRACKED_SESSIONS), ( "max_runs_per_session", - self.max_runs_per_session, + max_runs_per_session, DEFAULT_MAX_RETAINED_RUNS_PER_SESSION, ), ( "max_model_calls_per_run", - self.max_model_calls_per_run, + max_model_calls_per_run, DEFAULT_MAX_MODEL_CALLS_PER_RUN, ), ( "max_tool_executions_per_run", - self.max_tool_executions_per_run, + max_tool_executions_per_run, DEFAULT_MAX_TOOL_EXECUTIONS_PER_RUN, ), ( "max_activity_entries_per_run", - self.max_activity_entries_per_run, + max_activity_entries_per_run, DEFAULT_MAX_ACTIVITY_ENTRIES, ), ( "max_updates_per_run", - self.max_updates_per_run, + max_updates_per_run, DEFAULT_MAX_RETAINED_UPDATES_PER_RUN, ), ( "live_update_capacity", - self.live_update_capacity, + live_update_capacity, DEFAULT_MAX_RETAINED_UPDATES_PER_RUN, ), ]; @@ -172,6 +179,7 @@ struct DiagnosticSessionState { struct DiagnosticStoreState { sessions: HashMap, session_order: VecDeque, + live_updates: HashMap>>, } impl DiagnosticStoreState { @@ -225,6 +233,32 @@ impl DiagnosticStoreState { fn session(&self, scope: &DiagnosticScope) -> Option<&DiagnosticSessionState> { self.sessions.get(&DiagnosticSessionKey::from(scope)) } + + fn subscribe( + &mut self, + scope: DiagnosticScope, + capacity: usize, + ) -> broadcast::Receiver> { + self.prune_inactive_live_updates(); + let sender = self.live_updates.entry(scope).or_insert_with(|| { + let (sender, _) = broadcast::channel(capacity); + sender + }); + sender.subscribe() + } + + fn live_update_sender( + &mut self, + scope: &DiagnosticScope, + ) -> Option>> { + self.prune_inactive_live_updates(); + self.live_updates.get(scope).cloned() + } + + fn prune_inactive_live_updates(&mut self) { + self.live_updates + .retain(|_, sender| sender.receiver_count() > 0); + } } fn touch(order: &mut VecDeque, value: T) { @@ -238,17 +272,14 @@ fn touch(order: &mut VecDeque, value: T) { pub struct InMemoryDiagnosticStore { limits: DiagnosticStoreLimits, state: Mutex, - updates: broadcast::Sender>, } impl InMemoryDiagnosticStore { pub fn new(limits: DiagnosticStoreLimits) -> Result { let limits = limits.validate()?; - let (updates, _) = broadcast::channel(limits.live_update_capacity); Ok(Self { limits, state: Mutex::new(DiagnosticStoreState::default()), - updates, }) } @@ -258,7 +289,7 @@ impl InMemoryDiagnosticStore { prompt: PromptDiagnostic, ) -> Result { let update = DiagnosticUpdateKind::PromptUpdated { - component_count: prompt.components.len(), + component_count: u32::try_from(prompt.components.len()).unwrap_or(u32::MAX), total_estimated_tokens: prompt.total_estimated_tokens, truncated: prompt.any_content_truncated(), }; @@ -328,39 +359,12 @@ impl InMemoryDiagnosticStore { stats: SessionDiagnosticStats, ) -> Result { let stats = stats.into_bounded(); - let session_key = DiagnosticSessionKey::from(&scope); - 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); - let cursor = DiagnosticCursor::new(run.stream_id, sequence); - let envelope = DiagnosticUpdateEnvelope { + self.record_with_session( scope, - stream_id: run.stream_id, - sequence, - emitted_at: Utc::now(), - update: DiagnosticUpdateKind::Stats(stats.clone()), - }; - run.latest_sequence = sequence; - push_bounded( - &mut run.updates, - envelope.clone(), - self.limits.max_updates_per_run, - ); - let session = state - .sessions - .get_mut(&session_key) - .ok_or(DiagnosticStoreError::Invariant)?; - session.stats = stats; - let _ = self.updates.send(Arc::new(envelope)); - Ok(cursor) + DiagnosticUpdateKind::Stats(stats.clone()), + |_, _| {}, + move |session| session.stats = stats, + ) } pub fn snapshot( @@ -410,6 +414,7 @@ impl InMemoryDiagnosticStore { 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() } @@ -433,11 +438,16 @@ impl InMemoryDiagnosticStore { }) } - pub fn subscribe(&self, scope: DiagnosticScope) -> DiagnosticSubscription { - DiagnosticSubscription { - scope, - receiver: self.updates.subscribe(), - } + 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); + Ok(DiagnosticSubscription { receiver }) } fn record( @@ -446,6 +456,17 @@ impl InMemoryDiagnosticStore { update: DiagnosticUpdateKind, mutate: impl FnOnce(&mut DiagnosticRunState, DiagnosticSequence), ) -> Result { + self.record_with_session(scope, update, mutate, |_| {}) + } + + fn record_with_session( + &self, + scope: DiagnosticScope, + update: DiagnosticUpdateKind, + mutate_run: impl FnOnce(&mut DiagnosticRunState, DiagnosticSequence), + mutate_session: impl FnOnce(&mut DiagnosticSessionState), + ) -> Result { + let session_key = DiagnosticSessionKey::from(&scope); let mut state = self .state .lock() @@ -457,7 +478,7 @@ impl InMemoryDiagnosticStore { .checked_add(1) .ok_or(DiagnosticStoreError::SequenceExhausted)?; let sequence = DiagnosticSequence::new(next); - mutate(run, sequence); + mutate_run(run, sequence); let cursor = DiagnosticCursor::new(run.stream_id, sequence); let envelope = DiagnosticUpdateEnvelope { scope, @@ -472,7 +493,16 @@ impl InMemoryDiagnosticStore { envelope.clone(), self.limits.max_updates_per_run, ); - let _ = self.updates.send(Arc::new(envelope)); + let session = state + .sessions + .get_mut(&session_key) + .ok_or(DiagnosticStoreError::Invariant)?; + mutate_session(session); + let live_sender = state.live_update_sender(&envelope.scope); + drop(state); + if let Some(sender) = live_sender { + let _ = sender.send(Arc::new(envelope)); + } Ok(cursor) } } @@ -480,11 +510,9 @@ impl InMemoryDiagnosticStore { impl Default for InMemoryDiagnosticStore { fn default() -> Self { let limits = DiagnosticStoreLimits::default(); - let (updates, _) = broadcast::channel(limits.live_update_capacity); Self { limits, state: Mutex::new(DiagnosticStoreState::default()), - updates, } } } @@ -509,22 +537,18 @@ fn replace_or_push( } pub struct DiagnosticSubscription { - scope: DiagnosticScope, receiver: broadcast::Receiver>, } impl DiagnosticSubscription { pub async fn recv(&mut self) -> Result, DiagnosticStoreError> { - loop { - match self.receiver.recv().await { - Ok(update) if update.scope == self.scope => return Ok(update), - Ok(_) => continue, - Err(broadcast::error::RecvError::Lagged(count)) => { - return Err(DiagnosticStoreError::SubscriberLagged(count)); - } - Err(broadcast::error::RecvError::Closed) => { - return Err(DiagnosticStoreError::SubscriptionClosed); - } + 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) } } } @@ -589,6 +613,69 @@ mod tests { ); } + #[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 + }, + ), + ( + "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(); @@ -602,6 +689,14 @@ mod tests { ); } + #[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"); @@ -802,19 +897,63 @@ mod tests { assert_eq!(batch.latest_cursor, Some(recreated)); } - #[tokio::test] - async fn subscription_filters_scope_and_preserves_sequence() { + #[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()); + let mut subscription = store + .subscribe(allowed.clone()) + .expect("scoped subscription"); + for index in 0..=limits.live_update_capacity { + store + .record_activity(other.clone(), activity(&format!("other-{index}"))) + .expect("other"); + } store - .record_activity(other, activity("other")) - .expect("other"); + .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 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, activity("allowed")) + .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)); } @@ -822,7 +961,7 @@ mod tests { 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()); + let mut subscription = store.subscribe(scope.clone()).expect("scoped subscription"); let stats = SessionDiagnosticStats { total_model_calls: 7, ..SessionDiagnosticStats::default() diff --git a/crates/ironclaw_product_contracts/src/inspector.rs b/crates/ironclaw_product_contracts/src/inspector.rs index eb428d24a27..64596edea4c 100644 --- a/crates/ironclaw_product_contracts/src/inspector.rs +++ b/crates/ironclaw_product_contracts/src/inspector.rs @@ -1,16 +1,17 @@ //! Operator-only diagnostic vocabulary for the Web Debug Inspector. //! -//! These output-only DTOs cross the product boundary through the dedicated -//! inspection surface. They are deliberately separate from product projection -//! events: raw prompt components and tool details must never enter the normal -//! product stream. +//! 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::Serialize; +use serde::{Deserialize, Serialize}; use uuid::Uuid; pub const PROMPT_COMPONENT_CONTENT_MAX_BYTES: usize = 64 * 1024; @@ -62,7 +63,7 @@ impl std::fmt::Display for DiagnosticModelCallId { } } -#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)] #[serde(transparent)] pub struct DiagnosticStreamId(Uuid); @@ -88,7 +89,9 @@ impl std::fmt::Display for DiagnosticStreamId { } } -#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize)] +#[derive( + Debug, Clone, Copy, Default, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize, +)] #[serde(transparent)] pub struct DiagnosticSequence(u64); @@ -104,7 +107,7 @@ impl DiagnosticSequence { } } -#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)] pub struct DiagnosticCursor { pub stream_id: DiagnosticStreamId, pub sequence: DiagnosticSequence, @@ -310,18 +313,16 @@ impl PromptDiagnostic { let mut remaining = PROMPT_COMPONENT_TOTAL_MAX_BYTES; let mut bounded_components = Vec::with_capacity(original_component_count.min(MAX_PROMPT_COMPONENTS)); - let mut components_truncated = original_component_count > MAX_PROMPT_COMPONENTS; for mut component in components.into_iter().take(MAX_PROMPT_COMPONENTS) { if remaining == 0 { - components_truncated = true; break; } let retained = component.content.content().len().min(remaining); component.content = component.content.rebound(retained); - components_truncated |= component.content.truncated(); remaining = remaining.saturating_sub(component.content.content().len()); bounded_components.push(component); } + let components_truncated = bounded_components.len() < original_component_count; let active_skills_truncated = active_skills.len() > MAX_ACTIVE_SKILLS; let active_skills = active_skills @@ -598,7 +599,7 @@ impl SessionDiagnosticStats { #[serde(tag = "type", content = "data", rename_all = "snake_case")] pub enum DiagnosticUpdateKind { PromptUpdated { - component_count: usize, + component_count: u32, total_estimated_tokens: Option, truncated: bool, }, @@ -661,10 +662,24 @@ mod tests { let bounded = BoundedDiagnosticText::tool_result(value.clone()); assert!(bounded.truncated()); assert!(bounded.content().len() <= TOOL_RESULT_MAX_BYTES); - assert!(std::str::from_utf8(bounded.content().as_bytes()).is_ok()); + 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 bounded_text_debug_never_exposes_content() { let bounded = BoundedDiagnosticText::tool_result("super-secret-value"); @@ -709,7 +724,95 @@ mod tests { } #[test] - fn tool_result_uses_the_fifty_kibibyte_contract() { + 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(), @@ -719,7 +822,7 @@ mod tests { Some("x".repeat(TOOL_RESULT_MAX_BYTES + 1)), ToolExecutionStatus::Succeeded, None, - Some((TOOL_RESULT_MAX_BYTES + 1) as u64), + Some(7), None, None, ); From bacdbfb4c4e537a67695ef10713cfba30f5368c7 Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Wed, 5 Aug 2026 20:15:18 +0800 Subject: [PATCH 03/11] fix(inspector): preserve ordered scoped live updates --- .../ironclaw_product/src/inspector_store.rs | 66 +++++++++++++++---- 1 file changed, 55 insertions(+), 11 deletions(-) diff --git a/crates/ironclaw_product/src/inspector_store.rs b/crates/ironclaw_product/src/inspector_store.rs index 61a7fddc559..060f092c6d0 100644 --- a/crates/ironclaw_product/src/inspector_store.rs +++ b/crates/ironclaw_product/src/inspector_store.rs @@ -247,12 +247,15 @@ impl DiagnosticStoreState { sender.subscribe() } - fn live_update_sender( - &mut self, - scope: &DiagnosticScope, - ) -> Option>> { - self.prune_inactive_live_updates(); - self.live_updates.get(scope).cloned() + 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); + } } fn prune_inactive_live_updates(&mut self) { @@ -498,11 +501,7 @@ impl InMemoryDiagnosticStore { .get_mut(&session_key) .ok_or(DiagnosticStoreError::Invariant)?; mutate_session(session); - let live_sender = state.live_update_sender(&envelope.scope); - drop(state); - if let Some(sender) = live_sender { - let _ = sender.send(Arc::new(envelope)); - } + state.send_live_update(envelope); Ok(cursor) } } @@ -557,6 +556,7 @@ impl DiagnosticSubscription { #[cfg(test)] mod tests { use std::panic::{AssertUnwindSafe, catch_unwind}; + use std::sync::Barrier; use chrono::Utc; use ironclaw_host_api::{ @@ -928,6 +928,9 @@ mod tests { 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}"))) @@ -937,11 +940,52 @@ mod tests { .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 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"); From 9ecb1e7b551823cb3160ab924249c62446eaa0de Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Wed, 5 Aug 2026 20:53:08 +0800 Subject: [PATCH 04/11] fix(inspector): bound retained live update scopes --- .../ironclaw_product/src/inspector_store.rs | 96 +++++++++++++++++-- .../src/inspector.rs | 2 + 2 files changed, 91 insertions(+), 7 deletions(-) diff --git a/crates/ironclaw_product/src/inspector_store.rs b/crates/ironclaw_product/src/inspector_store.rs index 060f092c6d0..cc987ac466d 100644 --- a/crates/ironclaw_product/src/inspector_store.rs +++ b/crates/ironclaw_product/src/inspector_store.rs @@ -12,7 +12,7 @@ use ironclaw_host_api::{ turn::TurnRunId, }; use ironclaw_product_contracts::inspector::{ - DEFAULT_MAX_ACTIVITY_ENTRIES, DEFAULT_MAX_MODEL_CALLS_PER_RUN, + 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, @@ -31,6 +31,7 @@ pub struct DiagnosticStoreLimits { 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, } @@ -43,6 +44,7 @@ impl Default for DiagnosticStoreLimits { 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, } } @@ -59,6 +61,7 @@ impl DiagnosticStoreLimits { max_tool_executions_per_run, max_activity_entries_per_run, max_updates_per_run, + max_live_update_scopes, live_update_capacity, } = self; let values = [ @@ -88,6 +91,11 @@ impl DiagnosticStoreLimits { 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, @@ -180,6 +188,7 @@ struct DiagnosticStoreState { sessions: HashMap, session_order: VecDeque, live_updates: HashMap>>, + live_update_order: VecDeque, } impl DiagnosticStoreState { @@ -238,13 +247,27 @@ impl DiagnosticStoreState { &mut self, scope: DiagnosticScope, capacity: usize, - ) -> broadcast::Receiver> { + max_scopes: usize, + ) -> Result>, DiagnosticStoreError> { self.prune_inactive_live_updates(); - let sender = self.live_updates.entry(scope).or_insert_with(|| { + 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); - sender - }); - sender.subscribe() + 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) { @@ -255,12 +278,21 @@ impl DiagnosticStoreState { .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)); } } @@ -449,7 +481,11 @@ impl InMemoryDiagnosticStore { .state .lock() .map_err(|_| DiagnosticStoreError::StateUnavailable)?; - let receiver = state.subscribe(scope, self.limits.live_update_capacity); + let receiver = state.subscribe( + scope, + self.limits.live_update_capacity, + self.limits.max_live_update_scopes, + )?; Ok(DiagnosticSubscription { receiver }) } @@ -599,6 +635,7 @@ mod tests { 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, } } @@ -659,6 +696,13 @@ mod tests { ..valid }, ), + ( + "max_live_update_scopes", + DiagnosticStoreLimits { + max_live_update_scopes: 0, + ..valid + }, + ), ( "live_update_capacity", DiagnosticStoreLimits { @@ -949,6 +993,44 @@ mod tests { 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; diff --git a/crates/ironclaw_product_contracts/src/inspector.rs b/crates/ironclaw_product_contracts/src/inspector.rs index 64596edea4c..727afd2c96b 100644 --- a/crates/ironclaw_product_contracts/src/inspector.rs +++ b/crates/ironclaw_product_contracts/src/inspector.rs @@ -29,6 +29,8 @@ pub const MAX_MODELS_IN_STATS: usize = 64; 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; From 0d9e914d860eaf8ccbb556a3aa30623a75363dc3 Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Wed, 5 Aug 2026 21:39:42 +0800 Subject: [PATCH 05/11] fix(inspector): keep diagnostics aligned to run scope --- .../ironclaw_assistant/src/inspector_store.rs | 108 +++++++++++++----- 1 file changed, 77 insertions(+), 31 deletions(-) diff --git a/crates/ironclaw_assistant/src/inspector_store.rs b/crates/ironclaw_assistant/src/inspector_store.rs index cc987ac466d..9a76a11acc5 100644 --- a/crates/ironclaw_assistant/src/inspector_store.rs +++ b/crates/ironclaw_assistant/src/inspector_store.rs @@ -158,6 +158,7 @@ struct DiagnosticRunState { model_calls: VecDeque, tool_executions: VecDeque, activity: VecDeque, + stats: SessionDiagnosticStats, updates: VecDeque, latest_sequence: DiagnosticSequence, } @@ -170,6 +171,7 @@ impl Default for DiagnosticRunState { model_calls: VecDeque::new(), tool_executions: VecDeque::new(), activity: VecDeque::new(), + stats: SessionDiagnosticStats::default(), updates: VecDeque::new(), latest_sequence: DiagnosticSequence::ZERO, } @@ -180,7 +182,6 @@ impl Default for DiagnosticRunState { struct DiagnosticSessionState { runs: HashMap, run_order: VecDeque, - stats: SessionDiagnosticStats, } #[derive(Debug, Default)] @@ -239,10 +240,6 @@ impl DiagnosticStoreState { .get(&scope.run_id) } - fn session(&self, scope: &DiagnosticScope) -> Option<&DiagnosticSessionState> { - self.sessions.get(&DiagnosticSessionKey::from(scope)) - } - fn subscribe( &mut self, scope: DiagnosticScope, @@ -394,11 +391,10 @@ impl InMemoryDiagnosticStore { stats: SessionDiagnosticStats, ) -> Result { let stats = stats.into_bounded(); - self.record_with_session( + self.record( scope, DiagnosticUpdateKind::Stats(stats.clone()), - |_, _| {}, - move |session| session.stats = stats, + move |run, _| run.stats = stats, ) } @@ -413,10 +409,6 @@ impl InMemoryDiagnosticStore { let Some(run) = state.run(scope) else { return Ok(None); }; - let stats = state - .session(scope) - .map(|session| session.stats.clone()) - .unwrap_or_default(); Ok(Some(DiagnosticSnapshot { scope: scope.clone(), stream_id: run.stream_id, @@ -424,7 +416,7 @@ impl InMemoryDiagnosticStore { model_calls: run.model_calls.iter().cloned().collect(), tool_executions: run.tool_executions.iter().cloned().collect(), activity: run.activity.iter().cloned().collect(), - stats, + stats: run.stats.clone(), latest_sequence: run.latest_sequence, })) } @@ -443,7 +435,7 @@ impl InMemoryDiagnosticStore { updates: Vec::new(), retention_floor: None, latest_cursor: None, - rebase_required: false, + rebase_required: after.is_some(), }); }; let retention_floor = run.updates.front().map(DiagnosticUpdateEnvelope::cursor); @@ -495,17 +487,6 @@ impl InMemoryDiagnosticStore { update: DiagnosticUpdateKind, mutate: impl FnOnce(&mut DiagnosticRunState, DiagnosticSequence), ) -> Result { - self.record_with_session(scope, update, mutate, |_| {}) - } - - fn record_with_session( - &self, - scope: DiagnosticScope, - update: DiagnosticUpdateKind, - mutate_run: impl FnOnce(&mut DiagnosticRunState, DiagnosticSequence), - mutate_session: impl FnOnce(&mut DiagnosticSessionState), - ) -> Result { - let session_key = DiagnosticSessionKey::from(&scope); let mut state = self .state .lock() @@ -517,7 +498,7 @@ impl InMemoryDiagnosticStore { .checked_add(1) .ok_or(DiagnosticStoreError::SequenceExhausted)?; let sequence = DiagnosticSequence::new(next); - mutate_run(run, sequence); + mutate(run, sequence); let cursor = DiagnosticCursor::new(run.stream_id, sequence); let envelope = DiagnosticUpdateEnvelope { scope, @@ -532,11 +513,6 @@ impl InMemoryDiagnosticStore { envelope.clone(), self.limits.max_updates_per_run, ); - let session = state - .sessions - .get_mut(&session_key) - .ok_or(DiagnosticStoreError::Invariant)?; - mutate_session(session); state.send_live_update(envelope); Ok(cursor) } @@ -941,6 +917,32 @@ mod tests { 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"); @@ -1101,6 +1103,50 @@ mod tests { 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"); From ad558b7fcca110817a8a8329395e85e4c3c0a7f9 Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Thu, 6 Aug 2026 10:03:37 +0800 Subject: [PATCH 06/11] fix(inspector): deserialize diagnostic wire enums --- .../src/inspector.rs | 163 ++++++++++++++++-- 1 file changed, 150 insertions(+), 13 deletions(-) diff --git a/crates/contracts/ironclaw_product_contracts/src/inspector.rs b/crates/contracts/ironclaw_product_contracts/src/inspector.rs index 727afd2c96b..d2d5f163565 100644 --- a/crates/contracts/ironclaw_product_contracts/src/inspector.rs +++ b/crates/contracts/ironclaw_product_contracts/src/inspector.rs @@ -35,7 +35,7 @@ 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)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)] #[serde(transparent)] pub struct DiagnosticModelCallId(Uuid); @@ -152,7 +152,7 @@ impl DiagnosticScope { /// /// Construction is limited to the purpose-specific constructors so callers /// cannot silently select an unbounded maximum. -#[derive(Clone, PartialEq, Eq, Serialize)] +#[derive(Clone, PartialEq, Eq, Serialize, Deserialize)] pub struct BoundedDiagnosticText { content: String, original_bytes: u64, @@ -240,7 +240,7 @@ impl BoundedDiagnosticText { } } -#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum PromptComponentKind { System, @@ -374,7 +374,7 @@ impl PromptDiagnostic { } } -#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)] +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] pub struct ModelTokenUsage { pub input_tokens: Option, pub output_tokens: Option, @@ -382,7 +382,7 @@ pub struct ModelTokenUsage { pub cache_creation_input_tokens: Option, } -#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum InspectorModelCallStatus { Started, @@ -390,7 +390,7 @@ pub enum InspectorModelCallStatus { Failed, } -#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub struct ModelCallDiagnostic { pub call_id: DiagnosticModelCallId, pub iteration: u32, @@ -434,7 +434,7 @@ impl ModelCallDiagnostic { } } -#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum ToolExecutionStatus { Started, @@ -497,7 +497,7 @@ impl ToolExecutionDiagnostic { } } -#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum DiagnosticActivityKind { TurnStarted, @@ -515,7 +515,7 @@ pub enum DiagnosticActivityKind { StreamResumed, } -#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub struct DiagnosticActivityEvent { pub occurred_at: DateTime, pub kind: DiagnosticActivityKind, @@ -551,13 +551,13 @@ pub struct DiagnosticActivityEntry { pub event: DiagnosticActivityEvent, } -#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)] +#[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)] +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub struct DiagnosticModelCount { pub model: BoundedDiagnosticText, pub calls: u64, @@ -572,7 +572,7 @@ impl DiagnosticModelCount { } } -#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)] +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] pub struct SessionDiagnosticStats { pub total_model_calls: u64, pub calls_per_model: Vec, @@ -597,7 +597,7 @@ impl SessionDiagnosticStats { } } -#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(tag = "type", content = "data", rename_all = "snake_case")] pub enum DiagnosticUpdateKind { PromptUpdated { @@ -656,8 +656,22 @@ pub struct DiagnosticSnapshot { #[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); + } + } + #[test] fn bounded_text_preserves_utf8_and_reports_original_size() { let value = "€".repeat(TOOL_RESULT_MAX_BYTES); @@ -682,6 +696,129 @@ mod tests { 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 bounded_text_debug_never_exposes_content() { let bounded = BoundedDiagnosticText::tool_result("super-secret-value"); From f315ac81cdebd024a275d6bfbd99165ff6e4322d Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Thu, 6 Aug 2026 10:12:38 +0800 Subject: [PATCH 07/11] fix(inspector): validate diagnostic text decoding --- .../src/inspector.rs | 185 +++++++++++++++++- 1 file changed, 184 insertions(+), 1 deletion(-) diff --git a/crates/contracts/ironclaw_product_contracts/src/inspector.rs b/crates/contracts/ironclaw_product_contracts/src/inspector.rs index d2d5f163565..70b44618257 100644 --- a/crates/contracts/ironclaw_product_contracts/src/inspector.rs +++ b/crates/contracts/ironclaw_product_contracts/src/inspector.rs @@ -152,13 +152,53 @@ impl DiagnosticScope { /// /// Construction is limited to the purpose-specific constructors so callers /// cannot silently select an unbounded maximum. -#[derive(Clone, PartialEq, Eq, Serialize, Deserialize)] +#[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 @@ -208,6 +248,13 @@ impl BoundedDiagnosticText { 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; @@ -240,6 +287,54 @@ impl BoundedDiagnosticText { } } +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_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 { @@ -394,13 +489,16 @@ pub enum InspectorModelCallStatus { 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, } @@ -522,6 +620,7 @@ pub struct DiagnosticActivityEvent { pub iteration: Option, pub activity_id: Option, pub model_call_id: Option, + #[serde(default, deserialize_with = "deserialize_optional_bounded_summary")] pub summary: Option, } @@ -559,6 +658,7 @@ pub struct DiagnosticMetricTotal { #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub struct DiagnosticModelCount { + #[serde(deserialize_with = "deserialize_bounded_label")] pub model: BoundedDiagnosticText, pub calls: u64, } @@ -575,6 +675,7 @@ impl DiagnosticModelCount { #[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, @@ -609,6 +710,7 @@ pub enum DiagnosticUpdateKind { ToolExecutionUpdated { activity_id: CapabilityActivityId, model_call_id: Option, + #[serde(deserialize_with = "deserialize_bounded_label")] capability_name: BoundedDiagnosticText, status: ToolExecutionStatus, duration_ms: Option, @@ -819,6 +921,87 @@ mod tests { } } + #[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"); From 9a05b5a57d9c006eecbc7c50c1fa686b514f7f85 Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Thu, 6 Aug 2026 10:23:00 +0800 Subject: [PATCH 08/11] fix(inspector): enforce owning diagnostic decode limits --- .../src/inspector.rs | 204 +++++++++++++++++- 1 file changed, 202 insertions(+), 2 deletions(-) diff --git a/crates/contracts/ironclaw_product_contracts/src/inspector.rs b/crates/contracts/ironclaw_product_contracts/src/inspector.rs index 70b44618257..d38408a0aa0 100644 --- a/crates/contracts/ironclaw_product_contracts/src/inspector.rs +++ b/crates/contracts/ironclaw_product_contracts/src/inspector.rs @@ -296,6 +296,41 @@ where .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> @@ -347,10 +382,12 @@ pub enum PromptComponentKind { Other, } -#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +#[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, } @@ -540,7 +577,30 @@ pub enum ToolExecutionStatus { Failed, } -#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +#[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, @@ -554,6 +614,30 @@ pub struct ToolExecutionDiagnostic { 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)] @@ -774,6 +858,14 @@ mod tests { } } + 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); @@ -921,6 +1013,114 @@ mod tests { } } + #[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); From 9bc3377f7ffbbf7f25b9e41461784e6a16a48d73 Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Thu, 6 Aug 2026 13:55:51 +0800 Subject: [PATCH 09/11] fix(inspector): enforce prompt retention boundaries --- .../src/inspector.rs | 66 ++++++--- .../ironclaw_assistant/src/inspector_store.rs | 125 +++++++++++++++++- 2 files changed, 172 insertions(+), 19 deletions(-) diff --git a/crates/contracts/ironclaw_product_contracts/src/inspector.rs b/crates/contracts/ironclaw_product_contracts/src/inspector.rs index d38408a0aa0..c8f2e88f4f0 100644 --- a/crates/contracts/ironclaw_product_contracts/src/inspector.rs +++ b/crates/contracts/ironclaw_product_contracts/src/inspector.rs @@ -443,21 +443,6 @@ impl PromptDiagnostic { effective_model: Option, context_limit: Option, ) -> Self { - let original_component_count = 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 components.into_iter().take(MAX_PROMPT_COMPONENTS) { - if remaining == 0 { - break; - } - 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); - } - let components_truncated = bounded_components.len() < original_component_count; - let active_skills_truncated = active_skills.len() > MAX_ACTIVE_SKILLS; let active_skills = active_skills .into_iter() @@ -467,8 +452,8 @@ impl PromptDiagnostic { Self { captured_at, - components: bounded_components, - components_truncated, + components, + components_truncated: false, reconstructed_prompt: BoundedDiagnosticText::reconstructed_prompt(reconstructed_prompt), total_estimated_tokens, message_count, @@ -481,6 +466,53 @@ impl PromptDiagnostic { 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 { diff --git a/crates/product/ironclaw_assistant/src/inspector_store.rs b/crates/product/ironclaw_assistant/src/inspector_store.rs index 9a76a11acc5..07e6a8f1e74 100644 --- a/crates/product/ironclaw_assistant/src/inspector_store.rs +++ b/crates/product/ironclaw_assistant/src/inspector_store.rs @@ -320,6 +320,7 @@ impl InMemoryDiagnosticStore { 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, @@ -445,6 +446,7 @@ impl InMemoryDiagnosticStore { (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 @@ -576,8 +578,11 @@ mod tests { turn::TurnRunId, }; use ironclaw_product_contracts::inspector::{ - DiagnosticActivityEvent, DiagnosticActivityKind, DiagnosticModelCallId, DiagnosticScope, - InspectorModelCallStatus, ModelCallDiagnostic, PromptDiagnostic, ToolExecutionDiagnostic, + BoundedDiagnosticText, DIAGNOSTIC_LABEL_MAX_BYTES, DiagnosticActivityEvent, + DiagnosticActivityKind, DiagnosticModelCallId, DiagnosticScope, InspectorModelCallStatus, + MAX_ACTIVE_SKILLS, MAX_PROMPT_COMPONENTS, ModelCallDiagnostic, + PROMPT_COMPONENT_CONTENT_MAX_BYTES, PROMPT_COMPONENT_TOTAL_MAX_BYTES, + PromptComponentDiagnostic, PromptComponentKind, PromptDiagnostic, ToolExecutionDiagnostic, ToolExecutionStatus, }; @@ -794,6 +799,93 @@ mod tests { 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 session_eviction_is_deterministic_and_write_lru() { let store = InMemoryDiagnosticStore::new(tiny_limits()).expect("store"); @@ -889,6 +981,35 @@ mod tests { ); } + #[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"); From 18fa717b95524184a43f57cde2219b77abec3af1 Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Thu, 6 Aug 2026 16:06:47 +0800 Subject: [PATCH 10/11] fix(inspector): bound all retained diagnostics --- .../src/inspector.rs | 48 ++++++ .../ironclaw_assistant/src/inspector_store.rs | 155 +++++++++++++++++- 2 files changed, 198 insertions(+), 5 deletions(-) diff --git a/crates/contracts/ironclaw_product_contracts/src/inspector.rs b/crates/contracts/ironclaw_product_contracts/src/inspector.rs index c8f2e88f4f0..602d7239e27 100644 --- a/crates/contracts/ironclaw_product_contracts/src/inspector.rs +++ b/crates/contracts/ironclaw_product_contracts/src/inspector.rs @@ -599,6 +599,17 @@ impl ModelCallDiagnostic { 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)] @@ -704,6 +715,26 @@ impl ToolExecutionDiagnostic { } } + 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() @@ -758,6 +789,13 @@ impl DiagnosticActivityEvent { 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)] @@ -786,6 +824,11 @@ impl DiagnosticModelCount { 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)] @@ -810,6 +853,11 @@ impl SessionDiagnosticStats { 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 } } diff --git a/crates/product/ironclaw_assistant/src/inspector_store.rs b/crates/product/ironclaw_assistant/src/inspector_store.rs index 07e6a8f1e74..e1af24cb883 100644 --- a/crates/product/ironclaw_assistant/src/inspector_store.rs +++ b/crates/product/ironclaw_assistant/src/inspector_store.rs @@ -336,6 +336,7 @@ impl InMemoryDiagnosticStore { 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, _| { @@ -353,6 +354,7 @@ impl InMemoryDiagnosticStore { 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, @@ -375,6 +377,7 @@ impl InMemoryDiagnosticStore { 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| { @@ -578,12 +581,13 @@ mod tests { turn::TurnRunId, }; use ironclaw_product_contracts::inspector::{ - BoundedDiagnosticText, DIAGNOSTIC_LABEL_MAX_BYTES, DiagnosticActivityEvent, - DiagnosticActivityKind, DiagnosticModelCallId, DiagnosticScope, InspectorModelCallStatus, - MAX_ACTIVE_SKILLS, MAX_PROMPT_COMPONENTS, ModelCallDiagnostic, + 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, ToolExecutionDiagnostic, - ToolExecutionStatus, + PromptComponentDiagnostic, PromptComponentKind, PromptDiagnostic, TOOL_ARGUMENTS_MAX_BYTES, + TOOL_RESULT_MAX_BYTES, ToolExecutionDiagnostic, ToolExecutionStatus, }; use super::*; @@ -886,6 +890,147 @@ mod tests { ); } + #[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"); From 4a658096267ee83e80f6e7f979bfd39c07ea7711 Mon Sep 17 00:00:00 2001 From: italic-jinxin <106428113+italic-jinxin@users.noreply.github.com> Date: Thu, 6 Aug 2026 17:15:16 +0800 Subject: [PATCH 11/11] test(architecture): recapture product contracts ceiling --- .../tests/reborn_dependency_boundaries.rs | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) 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), ];