use std::collections::BTreeSet;
use runifold_core::{
Budget, BudgetEvent, LifecycleEvent, RunError, RunEvent, RunEventKind, RunId, Usage,
};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use thiserror::Error;
pub const MAX_EVENT_PAGE_SIZE: usize = 1_000;
#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(transparent)]
pub struct RunEventCursor(u64);
impl RunEventCursor {
#[must_use]
pub const fn after(sequence: u64) -> Self {
Self(sequence)
}
#[must_use]
pub const fn sequence(self) -> u64 {
self.0
}
}
#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(transparent)]
pub struct RunEventPageSize(usize);
impl RunEventPageSize {
pub const fn new(value: usize) -> Result<Self, RunEventQueryError> {
if value == 0 || value > MAX_EVENT_PAGE_SIZE {
Err(RunEventQueryError::InvalidPageSize { value })
} else {
Ok(Self(value))
}
}
#[must_use]
pub const fn get(self) -> usize {
self.0
}
}
#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
pub struct RunEventPage {
pub events: Vec<RunEvent>,
pub next: Option<RunEventCursor>,
}
pub trait RunEventSource: Send + Sync {
fn event_page(
&self,
run_id: RunId,
after: Option<RunEventCursor>,
limit: RunEventPageSize,
) -> Result<RunEventPage, RunEventSourceError>;
}
#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "snake_case")]
#[non_exhaustive]
pub enum RunEventSourceErrorKind {
Storage,
CorruptData,
}
#[derive(Clone, Debug, Error, Deserialize, Eq, PartialEq, Serialize)]
#[error("run event source {kind:?}: {message}")]
pub struct RunEventSourceError {
pub kind: RunEventSourceErrorKind,
pub message: String,
}
impl RunEventSourceError {
#[must_use]
pub fn storage(message: impl Into<String>) -> Self {
Self {
kind: RunEventSourceErrorKind::Storage,
message: message.into(),
}
}
#[must_use]
pub fn corrupt_data(message: impl Into<String>) -> Self {
Self {
kind: RunEventSourceErrorKind::CorruptData,
message: message.into(),
}
}
}
#[derive(Clone, Copy, Debug, Error, Eq, PartialEq)]
#[non_exhaustive]
pub enum RunEventQueryError {
#[error("event page size {value} must be between 1 and {MAX_EVENT_PAGE_SIZE}")]
InvalidPageSize {
value: usize,
},
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "snake_case")]
#[non_exhaustive]
pub enum RunStatus {
Running,
Completed,
Failed,
Cancelled,
}
#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
pub struct RunInspection {
pub schema_version: u32,
pub run_id: RunId,
pub status: RunStatus,
pub event_count: usize,
pub last_sequence: u64,
pub usage: Usage,
pub error: Option<RunError>,
pub domain_events: Vec<String>,
}
impl RunInspection {
pub const SCHEMA_VERSION: u32 = 1;
pub fn inspect(events: &[RunEvent]) -> Result<Self, InspectionError> {
let first = events.first().ok_or(InspectionError::Empty)?;
let run_id = first.meta.run_id;
let mut seen = BTreeSet::new();
let mut usage = Usage::default();
let mut terminal = None;
let mut error = None;
let mut domain_events = Vec::new();
for (index, event) in events.iter().enumerate() {
if event.meta.run_id != run_id {
return Err(InspectionError::MixedRun { index });
}
let expected = u64::try_from(index).unwrap_or(u64::MAX);
if event.meta.sequence != expected {
return Err(InspectionError::Sequence {
index,
expected,
actual: event.meta.sequence,
});
}
if event
.meta
.caused_by
.is_some_and(|cause| !seen.contains(&cause))
{
return Err(InspectionError::UnknownCause { index });
}
seen.insert(event.meta.event_id);
match &event.kind {
RunEventKind::Budget(BudgetEvent::Updated { usage: current }) => usage = *current,
RunEventKind::Domain(domain) => {
domain_events.push(format!("{}.{}", domain.namespace, domain.name));
}
RunEventKind::Lifecycle(LifecycleEvent::Completed { .. }) => {
set_terminal(&mut terminal, RunStatus::Completed, index)?;
}
RunEventKind::Lifecycle(LifecycleEvent::Failed { error: failure }) => {
set_terminal(&mut terminal, RunStatus::Failed, index)?;
error = Some(failure.clone());
}
RunEventKind::Lifecycle(LifecycleEvent::Cancelled) => {
set_terminal(&mut terminal, RunStatus::Cancelled, index)?;
}
_ => {}
}
}
Ok(Self {
schema_version: Self::SCHEMA_VERSION,
run_id,
status: terminal.unwrap_or(RunStatus::Running),
event_count: events.len(),
last_sequence: events.last().map_or(0, |event| event.meta.sequence),
usage,
error,
domain_events,
})
}
}
fn set_terminal(
terminal: &mut Option<RunStatus>,
status: RunStatus,
index: usize,
) -> Result<(), InspectionError> {
if terminal.replace(status).is_some() {
Err(InspectionError::MultipleTerminalEvents { index })
} else {
Ok(())
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct BudgetDimension {
pub name: String,
pub limit: Option<u64>,
pub used: u64,
pub remaining: Option<u64>,
pub exceeded: bool,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct BudgetExplanation {
pub dimensions: Vec<BudgetDimension>,
}
impl BudgetExplanation {
#[must_use]
pub fn new(budget: Budget, usage: Usage) -> Self {
let duration_limit = budget
.duration
.map(|value| u64::try_from(value.as_micros()).unwrap_or(u64::MAX));
Self {
dimensions: vec![
dimension("tokens", budget.tokens, usage.tokens),
dimension("cost_microusd", budget.cost_microusd, usage.cost_microusd),
dimension("duration_micros", duration_limit, usage.duration_micros),
dimension("turns", budget.turns, usage.turns),
dimension("tool_calls", budget.tool_calls, usage.tool_calls),
dimension("delegations", budget.delegations, usage.delegations),
],
}
}
}
fn dimension(name: &str, limit: Option<u64>, used: u64) -> BudgetDimension {
BudgetDimension {
name: name.into(),
limit,
used,
remaining: limit.map(|limit| limit.saturating_sub(used)),
exceeded: limit.is_some_and(|limit| used > limit),
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum CheckpointChangeKind {
Added,
Removed,
Changed,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct CheckpointChange {
pub path: String,
pub kind: CheckpointChangeKind,
}
#[must_use]
pub fn diff_checkpoints(before: &Value, after: &Value) -> Vec<CheckpointChange> {
let mut changes = Vec::new();
diff_at("", before, after, &mut changes);
changes.truncate(1_024);
changes
}
fn diff_at(path: &str, before: &Value, after: &Value, changes: &mut Vec<CheckpointChange>) {
if changes.len() >= 1_024 || before == after {
return;
}
match (before, after) {
(Value::Object(before), Value::Object(after)) => {
let keys = before.keys().chain(after.keys()).collect::<BTreeSet<_>>();
for key in keys {
let child = format!("{path}/{}", key.replace('~', "~0").replace('/', "~1"));
match (before.get(key), after.get(key)) {
(Some(left), Some(right)) => diff_at(&child, left, right, changes),
(None, Some(_)) => changes.push(CheckpointChange {
path: child,
kind: CheckpointChangeKind::Added,
}),
(Some(_), None) => changes.push(CheckpointChange {
path: child,
kind: CheckpointChangeKind::Removed,
}),
(None, None) => {}
}
}
}
_ => changes.push(CheckpointChange {
path: if path.is_empty() {
"/".into()
} else {
path.into()
},
kind: CheckpointChangeKind::Changed,
}),
}
}
#[derive(Clone, Debug, Error, Eq, PartialEq)]
#[non_exhaustive]
pub enum InspectionError {
#[error("run history is empty")]
Empty,
#[error("event {index} belongs to a different run")]
MixedRun {
index: usize,
},
#[error("event {index} has sequence {actual}; expected {expected}")]
Sequence {
index: usize,
expected: u64,
actual: u64,
},
#[error("event {index} references an unknown or future cause")]
UnknownCause {
index: usize,
},
#[error("event {index} adds a second terminal lifecycle state")]
MultipleTerminalEvents {
index: usize,
},
}
#[cfg(test)]
mod tests {
use runifold_core::{Budget, RunEvent, Usage};
use super::{
BudgetExplanation, CheckpointChangeKind, RunEventPageSize, RunEventQueryError,
RunInspection, RunStatus, diff_checkpoints,
};
#[test]
fn inspection_validates_and_summarizes_canonical_history() {
let scenario = runifold_test_fixture();
let inspection = RunInspection::inspect(&scenario).unwrap();
assert_eq!(inspection.status, RunStatus::Completed);
assert_eq!(inspection.event_count, 2);
}
#[test]
fn budget_explanation_is_saturating_and_marks_excess() {
let explanation = BudgetExplanation::new(
Budget {
tokens: Some(10),
..Budget::default()
},
Usage {
tokens: 12,
..Usage::default()
},
);
let tokens = &explanation.dimensions[0];
assert_eq!(tokens.remaining, Some(0));
assert!(tokens.exceeded);
}
#[test]
fn checkpoint_diff_reports_paths_without_values() {
let changes = diff_checkpoints(
&serde_json::json!({"secret": "old", "keep": 1}),
&serde_json::json!({"secret": "new", "add": true}),
);
assert!(changes.iter().any(|change| {
change.path == "/secret" && change.kind == CheckpointChangeKind::Changed
}));
assert!(!serde_json::to_string(&changes).unwrap().contains("old"));
}
#[test]
fn event_page_size_enforces_public_query_bounds() {
assert_eq!(RunEventPageSize::new(1).unwrap().get(), 1);
assert!(matches!(
RunEventPageSize::new(0),
Err(RunEventQueryError::InvalidPageSize { value: 0 })
));
assert!(RunEventPageSize::new(1_001).is_err());
}
fn runifold_test_fixture() -> Vec<RunEvent> {
use runifold_core::{
BudgetTracker, CapabilitySet, EventFactory, LifecycleEvent, RunContext, RunEventKind,
};
let run = RunContext::root(BudgetTracker::new(Budget::default()), CapabilitySet::new());
let factory = EventFactory::new(run.run_id(), None);
let started = factory.emit(RunEventKind::Lifecycle(LifecycleEvent::Started), None);
let completed = factory.emit(
RunEventKind::Lifecycle(LifecycleEvent::Completed {
output: serde_json::json!({"ok": true}),
}),
Some(started.meta.event_id),
);
vec![started, completed]
}
}