1use std::collections::BTreeSet;
4
5use runifold_core::{
6 Budget, BudgetEvent, LifecycleEvent, RunError, RunEvent, RunEventKind, RunId, Usage,
7};
8use serde::{Deserialize, Serialize};
9use serde_json::Value;
10use thiserror::Error;
11
12pub const MAX_EVENT_PAGE_SIZE: usize = 1_000;
14
15#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
17#[serde(transparent)]
18pub struct RunEventCursor(u64);
19
20impl RunEventCursor {
21 #[must_use]
23 pub const fn after(sequence: u64) -> Self {
24 Self(sequence)
25 }
26
27 #[must_use]
29 pub const fn sequence(self) -> u64 {
30 self.0
31 }
32}
33
34#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
36#[serde(transparent)]
37pub struct RunEventPageSize(usize);
38
39impl RunEventPageSize {
40 pub const fn new(value: usize) -> Result<Self, RunEventQueryError> {
46 if value == 0 || value > MAX_EVENT_PAGE_SIZE {
47 Err(RunEventQueryError::InvalidPageSize { value })
48 } else {
49 Ok(Self(value))
50 }
51 }
52
53 #[must_use]
55 pub const fn get(self) -> usize {
56 self.0
57 }
58}
59
60#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
62pub struct RunEventPage {
63 pub events: Vec<RunEvent>,
65 pub next: Option<RunEventCursor>,
67}
68
69pub trait RunEventSource: Send + Sync {
71 fn event_page(
77 &self,
78 run_id: RunId,
79 after: Option<RunEventCursor>,
80 limit: RunEventPageSize,
81 ) -> Result<RunEventPage, RunEventSourceError>;
82}
83
84#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
86#[serde(rename_all = "snake_case")]
87#[non_exhaustive]
88pub enum RunEventSourceErrorKind {
89 Storage,
91 CorruptData,
93}
94
95#[derive(Clone, Debug, Error, Deserialize, Eq, PartialEq, Serialize)]
97#[error("run event source {kind:?}: {message}")]
98pub struct RunEventSourceError {
99 pub kind: RunEventSourceErrorKind,
101 pub message: String,
103}
104
105impl RunEventSourceError {
106 #[must_use]
108 pub fn storage(message: impl Into<String>) -> Self {
109 Self {
110 kind: RunEventSourceErrorKind::Storage,
111 message: message.into(),
112 }
113 }
114
115 #[must_use]
117 pub fn corrupt_data(message: impl Into<String>) -> Self {
118 Self {
119 kind: RunEventSourceErrorKind::CorruptData,
120 message: message.into(),
121 }
122 }
123}
124
125#[derive(Clone, Copy, Debug, Error, Eq, PartialEq)]
127#[non_exhaustive]
128pub enum RunEventQueryError {
129 #[error("event page size {value} must be between 1 and {MAX_EVENT_PAGE_SIZE}")]
131 InvalidPageSize {
132 value: usize,
134 },
135}
136
137#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
139#[serde(rename_all = "snake_case")]
140#[non_exhaustive]
141pub enum RunStatus {
142 Running,
144 Completed,
146 Failed,
148 Cancelled,
150}
151
152#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
154pub struct RunInspection {
155 pub schema_version: u32,
157 pub run_id: RunId,
159 pub status: RunStatus,
161 pub event_count: usize,
163 pub last_sequence: u64,
165 pub usage: Usage,
167 pub error: Option<RunError>,
169 pub domain_events: Vec<String>,
171}
172
173impl RunInspection {
174 pub const SCHEMA_VERSION: u32 = 1;
176
177 pub fn inspect(events: &[RunEvent]) -> Result<Self, InspectionError> {
184 let first = events.first().ok_or(InspectionError::Empty)?;
185 let run_id = first.meta.run_id;
186 let mut seen = BTreeSet::new();
187 let mut usage = Usage::default();
188 let mut terminal = None;
189 let mut error = None;
190 let mut domain_events = Vec::new();
191 for (index, event) in events.iter().enumerate() {
192 if event.meta.run_id != run_id {
193 return Err(InspectionError::MixedRun { index });
194 }
195 let expected = u64::try_from(index).unwrap_or(u64::MAX);
196 if event.meta.sequence != expected {
197 return Err(InspectionError::Sequence {
198 index,
199 expected,
200 actual: event.meta.sequence,
201 });
202 }
203 if event
204 .meta
205 .caused_by
206 .is_some_and(|cause| !seen.contains(&cause))
207 {
208 return Err(InspectionError::UnknownCause { index });
209 }
210 seen.insert(event.meta.event_id);
211 match &event.kind {
212 RunEventKind::Budget(BudgetEvent::Updated { usage: current }) => usage = *current,
213 RunEventKind::Domain(domain) => {
214 domain_events.push(format!("{}.{}", domain.namespace, domain.name));
215 }
216 RunEventKind::Lifecycle(LifecycleEvent::Completed { .. }) => {
217 set_terminal(&mut terminal, RunStatus::Completed, index)?;
218 }
219 RunEventKind::Lifecycle(LifecycleEvent::Failed { error: failure }) => {
220 set_terminal(&mut terminal, RunStatus::Failed, index)?;
221 error = Some(failure.clone());
222 }
223 RunEventKind::Lifecycle(LifecycleEvent::Cancelled) => {
224 set_terminal(&mut terminal, RunStatus::Cancelled, index)?;
225 }
226 _ => {}
227 }
228 }
229 Ok(Self {
230 schema_version: Self::SCHEMA_VERSION,
231 run_id,
232 status: terminal.unwrap_or(RunStatus::Running),
233 event_count: events.len(),
234 last_sequence: events.last().map_or(0, |event| event.meta.sequence),
235 usage,
236 error,
237 domain_events,
238 })
239 }
240}
241
242fn set_terminal(
243 terminal: &mut Option<RunStatus>,
244 status: RunStatus,
245 index: usize,
246) -> Result<(), InspectionError> {
247 if terminal.replace(status).is_some() {
248 Err(InspectionError::MultipleTerminalEvents { index })
249 } else {
250 Ok(())
251 }
252}
253
254#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
256pub struct BudgetDimension {
257 pub name: String,
259 pub limit: Option<u64>,
261 pub used: u64,
263 pub remaining: Option<u64>,
265 pub exceeded: bool,
267}
268
269#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
271pub struct BudgetExplanation {
272 pub dimensions: Vec<BudgetDimension>,
274}
275
276impl BudgetExplanation {
277 #[must_use]
279 pub fn new(budget: Budget, usage: Usage) -> Self {
280 let duration_limit = budget
281 .duration
282 .map(|value| u64::try_from(value.as_micros()).unwrap_or(u64::MAX));
283 Self {
284 dimensions: vec![
285 dimension("tokens", budget.tokens, usage.tokens),
286 dimension("cost_microusd", budget.cost_microusd, usage.cost_microusd),
287 dimension("duration_micros", duration_limit, usage.duration_micros),
288 dimension("turns", budget.turns, usage.turns),
289 dimension("tool_calls", budget.tool_calls, usage.tool_calls),
290 dimension("delegations", budget.delegations, usage.delegations),
291 ],
292 }
293 }
294}
295
296fn dimension(name: &str, limit: Option<u64>, used: u64) -> BudgetDimension {
297 BudgetDimension {
298 name: name.into(),
299 limit,
300 used,
301 remaining: limit.map(|limit| limit.saturating_sub(used)),
302 exceeded: limit.is_some_and(|limit| used > limit),
303 }
304}
305
306#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
308#[serde(rename_all = "snake_case")]
309pub enum CheckpointChangeKind {
310 Added,
312 Removed,
314 Changed,
316}
317
318#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
320pub struct CheckpointChange {
321 pub path: String,
323 pub kind: CheckpointChangeKind,
325}
326
327#[must_use]
329pub fn diff_checkpoints(before: &Value, after: &Value) -> Vec<CheckpointChange> {
330 let mut changes = Vec::new();
331 diff_at("", before, after, &mut changes);
332 changes.truncate(1_024);
333 changes
334}
335
336fn diff_at(path: &str, before: &Value, after: &Value, changes: &mut Vec<CheckpointChange>) {
337 if changes.len() >= 1_024 || before == after {
338 return;
339 }
340 match (before, after) {
341 (Value::Object(before), Value::Object(after)) => {
342 let keys = before.keys().chain(after.keys()).collect::<BTreeSet<_>>();
343 for key in keys {
344 let child = format!("{path}/{}", key.replace('~', "~0").replace('/', "~1"));
345 match (before.get(key), after.get(key)) {
346 (Some(left), Some(right)) => diff_at(&child, left, right, changes),
347 (None, Some(_)) => changes.push(CheckpointChange {
348 path: child,
349 kind: CheckpointChangeKind::Added,
350 }),
351 (Some(_), None) => changes.push(CheckpointChange {
352 path: child,
353 kind: CheckpointChangeKind::Removed,
354 }),
355 (None, None) => {}
356 }
357 }
358 }
359 _ => changes.push(CheckpointChange {
360 path: if path.is_empty() {
361 "/".into()
362 } else {
363 path.into()
364 },
365 kind: CheckpointChangeKind::Changed,
366 }),
367 }
368}
369
370#[derive(Clone, Debug, Error, Eq, PartialEq)]
372#[non_exhaustive]
373pub enum InspectionError {
374 #[error("run history is empty")]
376 Empty,
377 #[error("event {index} belongs to a different run")]
379 MixedRun {
380 index: usize,
382 },
383 #[error("event {index} has sequence {actual}; expected {expected}")]
385 Sequence {
386 index: usize,
388 expected: u64,
390 actual: u64,
392 },
393 #[error("event {index} references an unknown or future cause")]
395 UnknownCause {
396 index: usize,
398 },
399 #[error("event {index} adds a second terminal lifecycle state")]
401 MultipleTerminalEvents {
402 index: usize,
404 },
405}
406
407#[cfg(test)]
408mod tests {
409 use runifold_core::{Budget, RunEvent, Usage};
410
411 use super::{
412 BudgetExplanation, CheckpointChangeKind, RunEventPageSize, RunEventQueryError,
413 RunInspection, RunStatus, diff_checkpoints,
414 };
415
416 #[test]
417 fn inspection_validates_and_summarizes_canonical_history() {
418 let scenario = runifold_test_fixture();
419 let inspection = RunInspection::inspect(&scenario).unwrap();
420
421 assert_eq!(inspection.status, RunStatus::Completed);
422 assert_eq!(inspection.event_count, 2);
423 }
424
425 #[test]
426 fn budget_explanation_is_saturating_and_marks_excess() {
427 let explanation = BudgetExplanation::new(
428 Budget {
429 tokens: Some(10),
430 ..Budget::default()
431 },
432 Usage {
433 tokens: 12,
434 ..Usage::default()
435 },
436 );
437 let tokens = &explanation.dimensions[0];
438 assert_eq!(tokens.remaining, Some(0));
439 assert!(tokens.exceeded);
440 }
441
442 #[test]
443 fn checkpoint_diff_reports_paths_without_values() {
444 let changes = diff_checkpoints(
445 &serde_json::json!({"secret": "old", "keep": 1}),
446 &serde_json::json!({"secret": "new", "add": true}),
447 );
448 assert!(changes.iter().any(|change| {
449 change.path == "/secret" && change.kind == CheckpointChangeKind::Changed
450 }));
451 assert!(!serde_json::to_string(&changes).unwrap().contains("old"));
452 }
453
454 #[test]
455 fn event_page_size_enforces_public_query_bounds() {
456 assert_eq!(RunEventPageSize::new(1).unwrap().get(), 1);
457 assert!(matches!(
458 RunEventPageSize::new(0),
459 Err(RunEventQueryError::InvalidPageSize { value: 0 })
460 ));
461 assert!(RunEventPageSize::new(1_001).is_err());
462 }
463
464 fn runifold_test_fixture() -> Vec<RunEvent> {
465 use runifold_core::{
466 BudgetTracker, CapabilitySet, EventFactory, LifecycleEvent, RunContext, RunEventKind,
467 };
468 let run = RunContext::root(BudgetTracker::new(Budget::default()), CapabilitySet::new());
469 let factory = EventFactory::new(run.run_id(), None);
470 let started = factory.emit(RunEventKind::Lifecycle(LifecycleEvent::Started), None);
471 let completed = factory.emit(
472 RunEventKind::Lifecycle(LifecycleEvent::Completed {
473 output: serde_json::json!({"ok": true}),
474 }),
475 Some(started.meta.event_id),
476 );
477 vec![started, completed]
478 }
479}