diff --git a/CLAUDE.md b/CLAUDE.md index bf8754eb4..6eb14e043 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -134,8 +134,9 @@ aura/ ### Shared Event Types (`aura-events`) - Lightweight crate defining `AuraStreamEvent` and `OrchestrationStreamEvent` enums +- `run::SessionEvent` is a session's own envelope: its `SessionId`, the `RunId` when the event belongs to a run, a `SequenceNumber` dense per session, a `Timestamp`, and either an `agent::AgentEvent` or a `LifecycleEvent` (started, observer attached/detached, claims exhausted, liveness decided, parked, finished, cancelled, failed). `RunId` is a UUID (v7 when minted); there is one per run. No producer emits it yet - Both `Serialize + Deserialize` — used by the web server (producer) and CLI (consumer) -- No agent, MCP, or provider dependencies — only `serde` and `serde_json` +- No agent, MCP, or provider dependencies — only `serde`, `serde_json`, and `uuid` - `ProgressToken` type uses a local wire-compatible definition by default; enables `rmcp-types` feature for direct rmcp interop (used by the `aura` crate) ## Environment Setup diff --git a/Cargo.lock b/Cargo.lock index ab2093cac..19753da9f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -327,6 +327,7 @@ dependencies = [ "rmcp", "serde", "serde_json", + "uuid", ] [[package]] diff --git a/crates/aura-cli/src/api/types.rs b/crates/aura-cli/src/api/types.rs index 33c03f40a..f15c80dcf 100644 --- a/crates/aura-cli/src/api/types.rs +++ b/crates/aura-cli/src/api/types.rs @@ -222,7 +222,7 @@ pub struct ShellCallDetail { pub command_name: String, pub full_command: String, pub result: String, - #[serde(with = "duration_millis")] + #[serde(with = "aura_events::duration_ms")] pub duration: Duration, } @@ -233,7 +233,7 @@ pub enum DisplayEvent { ToolCall { tool_name: String, arguments: BTreeMap, - #[serde(with = "duration_millis")] + #[serde(with = "aura_events::duration_ms")] duration: Duration, result: Option, }, @@ -270,7 +270,7 @@ pub enum DisplayEvent { diff_text: String, lines_added: usize, lines_removed: usize, - #[serde(with = "duration_millis")] + #[serde(with = "aura_events::duration_ms")] duration: Duration, }, // Bullet colors are derived at render time via `task_color_for(key)` @@ -349,20 +349,6 @@ pub fn snake_to_pascal_case(s: &str) -> String { .collect() } -pub mod duration_millis { - use serde::{Deserialize, Deserializer, Serialize, Serializer}; - use std::time::Duration; - - pub fn serialize(d: &Duration, s: S) -> Result { - d.as_millis().serialize(s) - } - - pub fn deserialize<'de, D: Deserializer<'de>>(d: D) -> Result { - let millis = u64::deserialize(d)?; - Ok(Duration::from_millis(millis)) - } -} - #[cfg(test)] mod tests { use super::*; @@ -674,21 +660,20 @@ mod tests { } // ----------------------------------------------------------------------- - // duration_millis serde module + // ShellCallDetail // ----------------------------------------------------------------------- #[test] - fn duration_millis_roundtrip() { - // Test via ShellCallDetail which uses #[serde(with = "duration_millis")] + fn shell_call_duration_is_milliseconds() { let detail = ShellCallDetail { command_name: "ls".to_string(), full_command: "ls -la".to_string(), result: "output".to_string(), duration: Duration::from_millis(1234), }; - let json = serde_json::to_string(&detail).unwrap(); - assert!(json.contains("1234")); - let parsed: ShellCallDetail = serde_json::from_str(&json).unwrap(); + let json = serde_json::to_value(&detail).unwrap(); + assert_eq!(json["duration"], 1234); + let parsed: ShellCallDetail = serde_json::from_value(json).unwrap(); assert_eq!(parsed.duration, Duration::from_millis(1234)); } } diff --git a/crates/aura-events/Cargo.toml b/crates/aura-events/Cargo.toml index 209961b97..5a313d1c1 100644 --- a/crates/aura-events/Cargo.toml +++ b/crates/aura-events/Cargo.toml @@ -10,6 +10,7 @@ repository.workspace = true [dependencies] serde = { workspace = true } serde_json = { workspace = true } +uuid = { workspace = true } # Optional: re-export rmcp's ProgressToken for interop with the aura crate. # When rmcp is not available, we define compatible local types. diff --git a/crates/aura-events/src/agent.rs b/crates/aura-events/src/agent.rs index fd88ca52d..efefa046d 100644 --- a/crates/aura-events/src/agent.rs +++ b/crates/aura-events/src/agent.rs @@ -229,6 +229,11 @@ pub enum AgentEventPayload { Synthesizing { iteration: usize, }, + + /// Any other value on the wire, read by a version of this crate that does + /// not know it. It cannot be written back. + #[serde(other, skip_serializing)] + Unknown, } #[cfg(test)] @@ -414,4 +419,17 @@ mod tests { assert_eq!(content, "the answer"); assert_eq!(usage.total_tokens.get(), 15); } + + /// A reader older than its producer reads a type it does not know as + /// `Unknown` rather than failing the event, and cannot write it back. + #[test] + fn an_unknown_type_reads_as_unknown_and_cannot_be_serialized() { + let event: AgentEvent = serde_json::from_value(json!({ + "agent": serde_json::to_value(AgentContext::single_agent()).unwrap(), + "payload": { "type": "telemetry", "cpu": 0.5 } + })) + .unwrap(); + assert!(matches!(event.payload, AgentEventPayload::Unknown)); + assert!(serde_json::to_value(&event).is_err()); + } } diff --git a/crates/aura-events/src/duration_ms.rs b/crates/aura-events/src/duration_ms.rs new file mode 100644 index 000000000..451473552 --- /dev/null +++ b/crates/aura-events/src/duration_ms.rs @@ -0,0 +1,97 @@ +//! Serializes a [`Duration`] as whole milliseconds, the unit every duration on +//! the wire is expressed in. Use it as +//! `#[serde(with = "aura_events::duration_ms")]`, or `duration_ms::option` for +//! an `Option`. + +use std::time::Duration; + +use serde::{Deserialize, Deserializer, Serialize, Serializer}; + +pub fn serialize(duration: &Duration, serializer: S) -> Result { + u64::try_from(duration.as_millis()) + .unwrap_or(u64::MAX) + .serialize(serializer) +} + +pub fn deserialize<'de, D: Deserializer<'de>>(deserializer: D) -> Result { + u64::deserialize(deserializer).map(Duration::from_millis) +} + +pub mod option { + use std::time::Duration; + + use serde::{Deserialize, Deserializer, Serialize, Serializer}; + + pub fn serialize( + duration: &Option, + serializer: S, + ) -> Result { + duration + .map(|d| u64::try_from(d.as_millis()).unwrap_or(u64::MAX)) + .serialize(serializer) + } + + pub fn deserialize<'de, D: Deserializer<'de>>( + deserializer: D, + ) -> Result, D::Error> { + Ok(Option::::deserialize(deserializer)?.map(Duration::from_millis)) + } +} + +#[cfg(test)] +mod tests { + use std::time::Duration; + + use serde::{Deserialize, Serialize}; + use serde_json::json; + + #[derive(Debug, PartialEq, Serialize, Deserialize)] + struct Timed { + #[serde(with = "super")] + elapsed: Duration, + #[serde(with = "super::option")] + limit: Option, + } + + #[test] + fn a_duration_is_whole_milliseconds() { + let timed = Timed { + elapsed: Duration::from_micros(1_500_999), + limit: Some(Duration::from_millis(30_000)), + }; + let json = serde_json::to_value(&timed).unwrap(); + assert_eq!(json, json!({ "elapsed": 1500, "limit": 30_000 })); + assert_eq!( + serde_json::from_value::(json).unwrap(), + Timed { + elapsed: Duration::from_millis(1500), + limit: Some(Duration::from_millis(30_000)), + } + ); + } + + #[test] + fn an_absent_duration_is_null() { + let timed = Timed { + elapsed: Duration::ZERO, + limit: None, + }; + let json = serde_json::to_value(&timed).unwrap(); + assert_eq!(json, json!({ "elapsed": 0, "limit": null })); + assert_eq!(serde_json::from_value::(json).unwrap(), timed); + } + + /// A duration too long for a u64 of milliseconds saturates rather than + /// failing to serialize. + #[test] + fn a_duration_past_u64_milliseconds_saturates() { + let timed = Timed { + elapsed: Duration::MAX, + limit: Some(Duration::MAX), + }; + assert_eq!( + serde_json::to_value(&timed).unwrap(), + json!({ "elapsed": u64::MAX, "limit": u64::MAX }) + ); + } +} diff --git a/crates/aura-events/src/lib.rs b/crates/aura-events/src/lib.rs index e39dedc94..e1279875a 100644 --- a/crates/aura-events/src/lib.rs +++ b/crates/aura-events/src/lib.rs @@ -10,13 +10,17 @@ //! //! - [`AuraStreamEvent`] — Base aura events (tool lifecycle, usage, reasoning, progress) //! - [`OrchestrationStreamEvent`] — Multi-agent orchestration events +//! - [`agent::AgentEvent`] — What a running agent emits, before any wire projection +//! - [`run::SessionEvent`] — One session's ordered stream: its agents' events and its runs' lifecycle //! -//! Both enums derive `Serialize + Deserialize` so they can be used for +//! All derive `Serialize + Deserialize` so they can be used for //! producing SSE (server) and parsing SSE (client) with the same types. pub mod agent; +pub mod duration_ms; pub mod event_names; pub mod orchestration; +pub mod run; use serde::{Deserialize, Serialize}; use std::collections::BTreeMap; @@ -135,6 +139,191 @@ macro_rules! string_newtype { }; } +/// Generates a string newtype that cannot hold the empty string: a fallible +/// constructor, `TryFrom` the string types, and deserialization that refuses +/// `""`. +macro_rules! nonempty_string_newtype { + ($(#[$meta:meta])* $name:ident) => { + $(#[$meta])* + #[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)] + #[serde(try_from = "String", into = "String")] + pub struct $name(String); + + impl $name { + pub fn new(value: impl Into) -> Result { + let value = value.into(); + if value.is_empty() { + Err($crate::EmptyId { + id: stringify!($name), + }) + } else { + Ok(Self(value)) + } + } + + pub fn as_str(&self) -> &str { + &self.0 + } + + pub fn into_string(self) -> String { + self.0 + } + } + + impl TryFrom for $name { + type Error = $crate::EmptyId; + + fn try_from(value: String) -> Result { + Self::new(value) + } + } + + impl TryFrom<&str> for $name { + type Error = $crate::EmptyId; + + fn try_from(value: &str) -> Result { + Self::new(value) + } + } + + impl From<$name> for String { + fn from(id: $name) -> Self { + id.0 + } + } + + impl std::fmt::Display for $name { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.write_str(&self.0) + } + } + + impl AsRef for $name { + fn as_ref(&self) -> &str { + &self.0 + } + } + + impl PartialEq for $name { + fn eq(&self, other: &str) -> bool { + self.0 == other + } + } + + impl PartialEq<&str> for $name { + fn eq(&self, other: &&str) -> bool { + self.0 == *other + } + } + }; +} + +// Path-based so submodules can invoke it regardless of declaration order. +pub(crate) use nonempty_string_newtype; + +/// An identifier given as the empty string. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct EmptyId { + /// The identifier's type name. + pub id: &'static str, +} + +impl std::fmt::Display for EmptyId { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!(f, "{} cannot be empty", self.id) + } +} + +impl std::error::Error for EmptyId {} + +string_newtype! { + /// Identifier for a session: an agent's identity over time, and the key + /// its [`SessionEvent`](run::SessionEvent) stream is ordered under. + SessionId +} + +/// Identifier for one run: a unit of work within a session, from a prompt to +/// a terminal state. +/// +/// This is the epic's "unique agent id for that run". It is not +/// [`AgentContext::agent_id`], which is [`CONVERSATION_AGENT_ID`] or a worker +/// name and attributes an event *within* a run. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)] +#[serde(try_from = "uuid::Uuid", into = "uuid::Uuid")] +pub struct RunId(uuid::Uuid); + +impl RunId { + /// A fresh id. Version 7, so ids order by when they were minted. + pub fn mint() -> Self { + Self(uuid::Uuid::now_v7()) + } + + pub fn as_uuid(&self) -> &uuid::Uuid { + &self.0 + } +} + +/// Why a value is not a run id. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum InvalidRunId { + /// Not a UUID. + Malformed(uuid::Error), + /// The nil UUID. + Nil, +} + +impl std::fmt::Display for InvalidRunId { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Self::Malformed(err) => err.fmt(f), + Self::Nil => f.write_str("the nil UUID names no run"), + } + } +} + +impl std::error::Error for InvalidRunId { + fn source(&self) -> Option<&(dyn std::error::Error + 'static)> { + match self { + Self::Malformed(err) => Some(err), + Self::Nil => None, + } + } +} + +impl TryFrom for RunId { + type Error = InvalidRunId; + + fn try_from(uuid: uuid::Uuid) -> Result { + if uuid.is_nil() { + Err(InvalidRunId::Nil) + } else { + Ok(Self(uuid)) + } + } +} + +impl From for uuid::Uuid { + fn from(id: RunId) -> Self { + id.0 + } +} + +impl std::str::FromStr for RunId { + type Err = InvalidRunId; + + fn from_str(s: &str) -> Result { + uuid::Uuid::from_str(s) + .map_err(InvalidRunId::Malformed)? + .try_into() + } +} + +impl std::fmt::Display for RunId { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + self.0.fmt(f) + } +} + string_newtype! { /// The id a model assigns to one tool call. ToolCallId @@ -237,6 +426,98 @@ impl std::fmt::Display for Progress { } } +/// An instant as Unix time in milliseconds. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)] +#[serde(transparent)] +pub struct Timestamp(u64); + +impl Timestamp { + pub fn from_unix_millis(millis: u64) -> Self { + Self(millis) + } + + pub fn unix_millis(self) -> u64 { + self.0 + } + + /// The current instant. A clock set before the Unix epoch reads as the + /// epoch itself. + pub fn now() -> Self { + let since_epoch = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default(); + Self(u64::try_from(since_epoch.as_millis()).unwrap_or(u64::MAX)) + } +} + +impl std::fmt::Display for Timestamp { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + self.0.fmt(f) + } +} + +/// A position in a sequence that starts at 1 and counts every item. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)] +#[serde(try_from = "u64", into = "u64")] +pub struct SequenceNumber(u64); + +impl SequenceNumber { + /// The first position. + pub const FIRST: Self = Self(1); + + pub fn get(self) -> u64 { + self.0 + } + + /// The position after this one. + pub fn next(self) -> Self { + Self(self.0.saturating_add(1)) + } + + /// Whether this is the position immediately after `previous`. `false` + /// means the stream did not advance by exactly one: a position was + /// missed, this one was delivered again, or it arrived out of order. + pub fn follows(self, previous: Self) -> bool { + self.0 == previous.0.saturating_add(1) + } +} + +impl std::fmt::Display for SequenceNumber { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + self.0.fmt(f) + } +} + +/// A sequence number of zero, which no sequence uses. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct ZeroSequenceNumber; + +impl std::fmt::Display for ZeroSequenceNumber { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.write_str("sequence numbers start at 1") + } +} + +impl std::error::Error for ZeroSequenceNumber {} + +impl TryFrom for SequenceNumber { + type Error = ZeroSequenceNumber; + + fn try_from(value: u64) -> Result { + if value == 0 { + Err(ZeroSequenceNumber) + } else { + Ok(Self(value)) + } + } +} + +impl From for u64 { + fn from(seq: SequenceNumber) -> Self { + seq.0 + } +} + /// The agent id of an orchestrated run's coordinator. Workers name it as their /// parent, so it defines the agent hierarchy and every producer spells it the /// same way. @@ -1403,6 +1684,58 @@ mod tests { other => panic!("expected ApprovalCompleted, got {:?}", other), } } + + /// A session's sequence starts at 1, so a zero is refused rather than + /// read as an event before the first. + #[test] + fn a_zero_sequence_number_is_refused() { + assert_eq!(SequenceNumber::try_from(0), Err(ZeroSequenceNumber)); + assert!(serde_json::from_value::(serde_json::json!(0)).is_err()); + assert_eq!( + serde_json::from_value::(serde_json::json!(1)).unwrap(), + SequenceNumber::FIRST + ); + } + + /// Dense numbering is what lets a consumer tell a missed event from a + /// quiet stream: the next number is always exactly one more. + #[test] + fn sequence_numbers_are_dense() { + let first = SequenceNumber::FIRST; + assert_eq!(first.get(), 1); + + let second = first.next(); + assert!(second.follows(first)); + assert!(!second.next().follows(first), "a skipped number is a gap"); + assert!(!first.follows(second), "order matters"); + assert!(!first.follows(first), "a repeat is not a successor"); + } + + #[test] + fn a_sequence_number_serializes_as_a_bare_integer() { + let json = serde_json::to_value(SequenceNumber::FIRST.next()).unwrap(); + assert_eq!(json, serde_json::json!(2)); + assert_eq!( + serde_json::from_value::(json).unwrap(), + SequenceNumber(2) + ); + } + + #[test] + fn a_timestamp_is_unix_milliseconds() { + let at = Timestamp::from_unix_millis(1_700_000_000_000); + assert_eq!( + serde_json::to_value(at).unwrap(), + serde_json::json!(1_700_000_000_000_u64) + ); + assert_eq!(at.unix_millis(), 1_700_000_000_000); + + let now = Timestamp::now(); + assert!( + now > at, + "now ({now}) should be after November 2023 ({at}) on any sane clock" + ); + } } // --------------------------------------------------------------------------- diff --git a/crates/aura-events/src/run.rs b/crates/aura-events/src/run.rs new file mode 100644 index 000000000..94f667b22 --- /dev/null +++ b/crates/aura-events/src/run.rs @@ -0,0 +1,1083 @@ +//! A session's event stream: what its agents did and what happened to its runs. +//! +//! [`AgentEvent`] says what an agent did. It says nothing about the run it +//! belongs to — that it started, who is watching, that it finished or was +//! cancelled and why. A [`SessionEvent`] is the envelope that carries both, so +//! a projection reads one ordered sequence and never correlates two. +//! +//! # Agent, session, run +//! +//! An agent is a config entry. A session is that agent's identity over time, +//! and the key its stream is ordered under. A run is one unit of work in a +//! session: a prompt or headless start, through every inference and tool turn, +//! to one terminal state. One stream per session means an agent's whole +//! history loads as one stream rather than as a composition of runs, and a run +//! boundary is a position in it. +//! +//! # Correlation lives on the envelope +//! +//! The session and run ids that [`AgentEvent`] deliberately omits are here, +//! applied once by whatever owns the session rather than by every observer. +//! Every event names its session. An event names a run when it belongs to +//! one; an observer attaching to an idle session belongs to none. +//! +//! A run has one id, a [`RunId`] minted once when the run starts. +//! +//! # Sequence numbers are dense +//! +//! Every event of a session carries a [`SequenceNumber`], starting at +//! [`SequenceNumber::FIRST`] and increasing by exactly one per event across +//! all of the session's runs. A consumer that sees a gap knows it missed +//! something rather than silently reading a partial stream. The number is +//! minted in one place — by whatever appends to the session's stream — never +//! by a relay or projection. +//! +//! Order by `seq`, never by `at`. The [`Timestamp`] is wall-clock time, for +//! display. It is Unix milliseconds, a number, where other instants on the wire, +//! such as an approval's `expires_at`, are RFC 3339 strings: it is stamped on +//! every event, and a consumer only displays it. +//! +//! # Tags and formats +//! +//! The payload is adjacently tagged: `kind` names the half, and the event sits +//! under `event`, so no field of any event can collide with the envelope's tag. +//! A lifecycle event is internally tagged by `type`, as +//! [`AgentEventPayload`](crate::agent::AgentEventPayload) is, so a lifecycle +//! field named `type` — or `reason`, which [`RunCancelReason`] flattens into +//! [`LifecycleEvent::Cancelled`] — would collide; the roundtrip test over +//! every variant is what catches one. +//! +//! A reader older than its producer meets variants it does not know. Every +//! tagged enum here reads any other tag as its `Unknown` variant, so the event +//! still parses and its `seq` still counts. [`SessionEventPayload`] is the +//! exception: serde cannot fall back on an adjacent tag that carries content, +//! and its two kinds are not expected to grow. An `Unknown` cannot be written +//! back, so a relay forwards the bytes it received, never a re-encoded event. +//! +//! Internal tagging and `#[serde(flatten)]` need a self-describing format. +//! These types round-trip through JSON or MessagePack, not `bincode` or +//! `postcard`, and a store that persists them inherits that. +//! +//! # Lifecycle beside activity +//! +//! [`AgentEventPayload::RunParked`](crate::agent::AgentEventPayload::RunParked) +//! is an agent event: the orchestrator emits it from inside the run with its +//! iteration state. [`LifecycleEvent::Parked`] is the run owner's +//! acknowledgement and carries the checkpoint reference. Both appear in the +//! stream; a projection decides which it surfaces. +//! +//! [`LifecycleEvent::Started`] carries the prompt, not the history. History +//! is the run's input, too large to replay to every late observer. + +use std::time::Duration; + +use serde::{Deserialize, Serialize}; + +use crate::agent::AgentEvent; +use crate::{nonempty_string_newtype, RunId, SequenceNumber, SessionId, Timestamp, TokenUsage}; + +nonempty_string_newtype! { + /// Identifier for one observer of a session, unique among its observers. + ObserverId +} + +nonempty_string_newtype! { + /// Locates a parked run's checkpoint, as the park machinery names it. + CheckpointRef +} + +/// What an observer's subscription does to the live run's lifetime. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +#[non_exhaustive] +pub enum ObserverKind { + /// Reads the stream without a claim on any run. + Collecting, + /// Holds the live run's right to continue. + Claiming, + /// Any other value on the wire, read by a version of this crate that does + /// not know it. It cannot be written back. + #[serde(other, skip_serializing)] + Unknown, +} + +/// One observer of a session. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct Observer { + pub id: ObserverId, + pub kind: ObserverKind, + /// Whether a human is at this observer. + pub presence: bool, +} + +/// Why an observer stopped observing. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +#[non_exhaustive] +pub enum DetachCause { + /// The observer let go. + Released, + /// A lease held across a process boundary lapsed unrenewed. + Expired, + /// Another claimant took over. + Displaced, + /// Any other value on the wire, read by a version of this crate that does + /// not know it. It cannot be written back. + #[serde(other, skip_serializing)] + Unknown, +} + +/// What a run does once nothing claims it. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +#[non_exhaustive] +pub enum LivenessPolicy { + /// End the run. + Cancel, + /// Run to completion unclaimed. + Continue, + /// Checkpoint and end the run resumable. + Park, + /// Any other value on the wire, read by a version of this crate that does + /// not know it. It cannot be written back. + #[serde(other, skip_serializing)] + Unknown, +} + +/// How a run reacts to going unclaimed. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +pub struct Liveness { + pub policy: LivenessPolicy, + /// How long the run goes unclaimed before `policy` acts. + #[serde(rename = "grace_ms", with = "crate::duration_ms")] + pub grace: Duration, +} + +impl Default for Liveness { + /// Cancel at once: a run nobody claims spends provider turns for no one. + fn default() -> Self { + Self { + policy: LivenessPolicy::Cancel, + grace: Duration::ZERO, + } + } +} + +/// Why a run was ended. +#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[serde(tag = "reason", rename_all = "snake_case")] +#[non_exhaustive] +pub enum RunCancelReason { + /// The run outlived the bound it was started with. + Deadline { + #[serde(rename = "after_ms", with = "crate::duration_ms")] + after: Duration, + }, + /// Something outside the run cancelled it. + External, + /// The model called a tool the caller executes. + ClientTool, + /// Nothing claimed the run, and its liveness policy ended it. + Unclaimed, + /// The process stopped serving runs. + Shutdown, + /// Any other value on the wire, read by a version of this crate that does + /// not know it. It cannot be written back. + #[serde(other, skip_serializing)] + Unknown, +} + +/// What happened to a session's run, or to the session itself, as distinct +/// from what its agent did. +/// +/// `#[non_exhaustive]` because the vocabulary grows as the session's owner +/// learns to say more. +#[derive(Clone, Debug, Serialize, Deserialize)] +#[serde(tag = "type", rename_all = "snake_case")] +#[non_exhaustive] +pub enum LifecycleEvent { + Started { + /// The configured agent running, matching + /// [`AgentInfo::id`](crate::AgentInfo::id). + agent: String, + prompt: String, + #[serde( + rename = "timeout_ms", + default, + skip_serializing_if = "Option::is_none", + with = "crate::duration_ms::option" + )] + timeout: Option, + #[serde(default, deserialize_with = "null_as_default")] + liveness: Liveness, + /// The run this one continues, if any. + #[serde(default, skip_serializing_if = "Option::is_none")] + continues: Option, + }, + + ObserverAttached { + observer: Observer, + }, + + ObserverDetached { + observer: Observer, + because: DetachCause, + }, + + /// The last claiming observer detached. + ClaimsExhausted, + + LivenessDecided { + policy: LivenessPolicy, + }, + + /// The run stopped resumable. + Parked { + checkpoint: CheckpointRef, + /// As [`Finished`](Self::Finished)'s. + #[serde(flatten)] + usage: TokenUsage, + }, + + /// The run completed its work. + Finished { + /// Cumulative provider-billed tokens across every turn. + #[serde(flatten)] + usage: TokenUsage, + }, + + /// The run was ended before completing. + Cancelled { + #[serde(flatten)] + reason: RunCancelReason, + /// The caller's own words. + #[serde(default, skip_serializing_if = "Option::is_none")] + message: Option, + /// As [`Finished`](Self::Finished)'s. + #[serde(flatten)] + usage: TokenUsage, + }, + + /// The run ended in an error. + Failed { + error: String, + /// As [`Finished`](Self::Finished)'s. + #[serde(flatten)] + usage: TokenUsage, + }, + + /// Any other value on the wire, read by a version of this crate that does + /// not know it. It cannot be written back. + #[serde(other, skip_serializing)] + Unknown, +} + +/// Reads `null` as the field's default, as a missing field already is. +fn null_as_default<'de, D, T>(deserializer: D) -> Result +where + D: serde::Deserializer<'de>, + T: Default + Deserialize<'de>, +{ + Ok(Option::::deserialize(deserializer)?.unwrap_or_default()) +} + +impl LifecycleEvent { + /// Whether this ends its run: a park, a finish, a cancel, or a failure. + /// The session's stream goes on past it. + pub fn ends_run(&self) -> bool { + matches!( + self, + Self::Parked { .. } + | Self::Finished { .. } + | Self::Cancelled { .. } + | Self::Failed { .. } + ) + } +} + +/// Either half of what a session's stream carries. +#[expect( + clippy::large_enum_variant, + reason = "agent events are nearly every event on a session's stream; \ + boxing them would add an allocation to each to save memory \ + only on the rare lifecycle ones" +)] +#[derive(Clone, Debug, Serialize, Deserialize)] +#[serde(tag = "kind", content = "event", rename_all = "snake_case")] +#[non_exhaustive] +pub enum SessionEventPayload { + Agent(AgentEvent), + Lifecycle(LifecycleEvent), +} + +/// One event in a session's ordered stream. +#[derive(Clone, Debug, Serialize, Deserialize)] +pub struct SessionEvent { + pub session_id: SessionId, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub run_id: Option, + pub seq: SequenceNumber, + pub at: Timestamp, + pub payload: SessionEventPayload, +} + +impl SessionEvent { + /// See [`LifecycleEvent::ends_run`]. + pub fn ends_run(&self) -> bool { + match &self.payload { + SessionEventPayload::Lifecycle(lifecycle) => lifecycle.ends_run(), + _ => false, + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::agent::AgentEventPayload; + use crate::{AgentContext, TokenCount}; + use serde_json::json; + + const RUN: &str = "0191e8c0-1111-7000-8000-000000000001"; + + fn run_id() -> RunId { + RUN.parse().expect("a well-formed UUID") + } + + fn envelope(seq: u64, run: Option, payload: SessionEventPayload) -> SessionEvent { + SessionEvent { + session_id: SessionId::new("sess_1"), + run_id: run, + seq: SequenceNumber::try_from(seq).expect("a sequence starts at 1"), + at: Timestamp::from_unix_millis(1_700_000_000_000), + payload, + } + } + + fn lifecycle(seq: u64, event: LifecycleEvent) -> SessionEvent { + envelope(seq, Some(run_id()), SessionEventPayload::Lifecycle(event)) + } + + fn roundtrip(event: &SessionEvent) -> SessionEvent { + let json = serde_json::to_string(event).expect("session event should serialize"); + serde_json::from_str(&json) + .unwrap_or_else(|err| panic!("session event should deserialize: {err}\n{json}")) + } + + fn usage() -> TokenUsage { + TokenUsage { + prompt_tokens: TokenCount::new(10), + completion_tokens: TokenCount::new(5), + total_tokens: TokenCount::new(15), + } + } + + fn observer(kind: ObserverKind, presence: bool) -> Observer { + Observer { + id: ObserverId::new("obs_1").unwrap(), + kind, + presence, + } + } + + /// The correlation `AgentEvent` leaves off — session, run — sits on the + /// envelope, beside the ordering a consumer needs, and the agent's own + /// event is carried whole under `event` rather than reshaped. + #[test] + fn an_agent_event_travels_inside_the_envelope_untouched() { + let event = envelope( + 3, + Some(run_id()), + SessionEventPayload::Agent(AgentEvent::single_agent(AgentEventPayload::TextDelta { + content: "hi".to_string(), + })), + ); + let json = serde_json::to_value(&event).expect("should serialize"); + + assert_eq!( + json, + json!({ + "session_id": "sess_1", + "run_id": RUN, + "seq": 3, + "at": 1_700_000_000_000_u64, + "payload": { + "kind": "agent", + "event": { + "agent": { "agent_id": "main" }, + "payload": { "type": "text_delta", "content": "hi" } + } + } + }) + ); + + let SessionEventPayload::Agent(inner) = roundtrip(&event).payload else { + panic!("expected an agent event"); + }; + assert!(inner.agent.is_single_agent()); + assert!(matches!(inner.payload, AgentEventPayload::TextDelta { .. })); + } + + /// A lifecycle event shares the envelope, told apart by `kind` and then + /// by its own `type`, so a consumer switches on two tags and never on + /// field presence. + #[test] + fn a_lifecycle_event_nests_under_its_kind() { + let json = serde_json::to_value(lifecycle( + 1, + LifecycleEvent::Started { + agent: "sre".to_string(), + prompt: "why is checkout slow".to_string(), + timeout: Some(Duration::from_secs(300)), + liveness: Liveness { + policy: LivenessPolicy::Continue, + grace: Duration::from_secs(30), + }, + continues: None, + }, + )) + .expect("should serialize"); + + assert_eq!( + json["payload"], + json!({ + "kind": "lifecycle", + "event": { + "type": "started", + "agent": "sre", + "prompt": "why is checkout slow", + "timeout_ms": 300_000, + "liveness": { "policy": "continue", "grace_ms": 30_000 } + } + }) + ); + } + + /// An observer attaching to a session with no live run belongs to no run, + /// and the envelope says so by omission rather than with an empty id. + #[test] + fn an_event_outside_any_run_omits_the_run_id() { + let event = envelope( + 1, + None, + SessionEventPayload::Lifecycle(LifecycleEvent::ObserverAttached { + observer: observer(ObserverKind::Collecting, true), + }), + ); + + let json = serde_json::to_value(&event).unwrap(); + assert!(json.get("run_id").is_none()); + assert_eq!(json["session_id"], "sess_1"); + assert_eq!(roundtrip(&event).run_id, None); + } + + /// Every lifecycle variant, every cancel reason and every detach cause + /// survives a roundtrip. A field named after a tag — `type` on a lifecycle + /// variant, `reason` beside a flattened cancel reason — serializes as a + /// duplicate key and fails here. + #[test] + fn every_lifecycle_variant_survives_a_roundtrip() { + let variants = vec![ + LifecycleEvent::Started { + agent: "sre".to_string(), + prompt: "hi".to_string(), + timeout: None, + liveness: Liveness::default(), + continues: None, + }, + LifecycleEvent::ObserverAttached { + observer: observer(ObserverKind::Claiming, true), + }, + LifecycleEvent::ObserverDetached { + observer: observer(ObserverKind::Collecting, false), + because: DetachCause::Released, + }, + LifecycleEvent::ObserverDetached { + observer: observer(ObserverKind::Claiming, false), + because: DetachCause::Expired, + }, + LifecycleEvent::ObserverDetached { + observer: observer(ObserverKind::Claiming, true), + because: DetachCause::Displaced, + }, + LifecycleEvent::ClaimsExhausted, + LifecycleEvent::LivenessDecided { + policy: LivenessPolicy::Park, + }, + LifecycleEvent::Parked { + checkpoint: CheckpointRef::new("memory/sess_1/parked/run_1.json").unwrap(), + usage: usage(), + }, + LifecycleEvent::Finished { usage: usage() }, + LifecycleEvent::Cancelled { + reason: RunCancelReason::Deadline { + after: Duration::from_millis(1500), + }, + message: None, + usage: usage(), + }, + LifecycleEvent::Cancelled { + reason: RunCancelReason::External, + message: Some("operator stopped it".to_string()), + usage: usage(), + }, + LifecycleEvent::Cancelled { + reason: RunCancelReason::ClientTool, + message: None, + usage: usage(), + }, + LifecycleEvent::Cancelled { + reason: RunCancelReason::Unclaimed, + message: None, + usage: usage(), + }, + LifecycleEvent::Cancelled { + reason: RunCancelReason::Shutdown, + message: None, + usage: usage(), + }, + LifecycleEvent::Failed { + error: "provider returned 500".to_string(), + usage: usage(), + }, + ]; + + for (i, variant) in variants.into_iter().enumerate() { + let before = lifecycle(i as u64 + 1, variant); + let before_json = serde_json::to_value(&before).unwrap(); + let after_json = serde_json::to_value(roundtrip(&before)).unwrap(); + assert_eq!( + before_json, after_json, + "variant {i} changed across a roundtrip" + ); + } + } + + /// The observer's `kind` sits inside the observer, inside the event, so + /// it never meets the envelope's `kind`; the detach says why it left. + #[test] + fn a_detach_names_its_observer_and_why_it_left() { + let json = serde_json::to_value(lifecycle( + 2, + LifecycleEvent::ObserverDetached { + observer: observer(ObserverKind::Claiming, true), + because: DetachCause::Displaced, + }, + )) + .unwrap(); + + assert_eq!( + json["payload"], + json!({ + "kind": "lifecycle", + "event": { + "type": "observer_detached", + "observer": { "id": "obs_1", "kind": "claiming", "presence": true }, + "because": "displaced" + } + }) + ); + } + + /// A unit variant under an internal tag is just its tag, so the event with + /// nothing to say still parses. + #[test] + fn claims_exhausted_carries_only_its_tag() { + let json = serde_json::to_value(lifecycle(2, LifecycleEvent::ClaimsExhausted)).unwrap(); + assert_eq!( + json["payload"], + json!({ "kind": "lifecycle", "event": { "type": "claims_exhausted" } }) + ); + + let parsed: SessionEvent = serde_json::from_value(json).unwrap(); + assert!(matches!( + parsed.payload, + SessionEventPayload::Lifecycle(LifecycleEvent::ClaimsExhausted) + )); + } + + /// The cancel reason flattens under the event the way a tool outcome + /// flattens under `tool_complete`, and the caller's words ride beside it. + #[test] + fn a_cancel_reason_flattens_beside_the_callers_message() { + let json = serde_json::to_value(lifecycle( + 9, + LifecycleEvent::Cancelled { + reason: RunCancelReason::Deadline { + after: Duration::from_secs(300), + }, + message: Some("request timeout".to_string()), + usage: usage(), + }, + )) + .unwrap(); + + assert_eq!( + json["payload"]["event"], + json!({ + "type": "cancelled", + "reason": "deadline", + "after_ms": 300_000, + "message": "request timeout", + "prompt_tokens": 10, + "completion_tokens": 5, + "total_tokens": 15 + }) + ); + + for (reason, tag) in [ + (RunCancelReason::External, "external"), + (RunCancelReason::Unclaimed, "unclaimed"), + (RunCancelReason::Shutdown, "shutdown"), + ] { + let json = serde_json::to_value(lifecycle( + 9, + LifecycleEvent::Cancelled { + reason, + message: None, + usage: usage(), + }, + )) + .unwrap(); + assert_eq!( + json["payload"]["event"], + json!({ + "type": "cancelled", + "reason": tag, + "prompt_tokens": 10, + "completion_tokens": 5, + "total_tokens": 15 + }) + ); + } + } + + /// Durations go on the wire as whole milliseconds; anything finer is + /// dropped, and nothing rounds up. + #[test] + fn a_duration_is_whole_milliseconds_on_the_wire() { + let before = RunCancelReason::Deadline { + after: Duration::from_micros(1_500_999), + }; + let json = serde_json::to_value(before).unwrap(); + assert_eq!(json, json!({ "reason": "deadline", "after_ms": 1500 })); + + let after: RunCancelReason = serde_json::from_value(json).unwrap(); + assert_eq!( + after, + RunCancelReason::Deadline { + after: Duration::from_millis(1500) + } + ); + } + + /// A start that omits its liveness parses, and reads as the default: + /// cancel at once. + #[test] + fn a_start_without_liveness_reads_as_cancel_at_once() { + let parsed: SessionEvent = serde_json::from_value(json!({ + "session_id": "sess_1", + "run_id": RUN, + "seq": 1, + "at": 0, + "payload": { + "kind": "lifecycle", + "event": { "type": "started", "agent": "sre", "prompt": "hi" } + } + })) + .unwrap(); + + let SessionEventPayload::Lifecycle(LifecycleEvent::Started { liveness, .. }) = + parsed.payload + else { + panic!("expected a start"); + }; + assert_eq!(liveness, Liveness::default()); + assert_eq!(liveness.policy, LivenessPolicy::Cancel); + assert_eq!(liveness.grace, Duration::ZERO); + } + + /// A run that continues another names it; one that does not carries no + /// `continues` on the wire, so an older payload reads as a fresh run. + #[test] + fn a_continuing_start_names_the_run_it_follows() { + let parked = run_id(); + let start = LifecycleEvent::Started { + agent: "sre".to_string(), + prompt: "hi".to_string(), + timeout: None, + liveness: Liveness::default(), + continues: Some(parked), + }; + let json = serde_json::to_value(&start).unwrap(); + assert_eq!(json["continues"], RUN); + + let LifecycleEvent::Started { continues, .. } = serde_json::from_value(json).unwrap() + else { + panic!("expected a start"); + }; + assert_eq!(continues, Some(parked)); + + let fresh = serde_json::to_value(LifecycleEvent::Started { + agent: "sre".to_string(), + prompt: "hi".to_string(), + timeout: None, + liveness: Liveness::default(), + continues: None, + }) + .unwrap(); + assert!(fresh.get("continues").is_none()); + } + + /// Every terminal event carries the run's cumulative usage at the top + /// level, as `Finished` does, so cost reads the same way off any of them. + #[test] + fn every_terminal_event_carries_usage_the_same_way() { + let terminal = [ + LifecycleEvent::Parked { + checkpoint: CheckpointRef::new("c").unwrap(), + usage: usage(), + }, + LifecycleEvent::Finished { usage: usage() }, + LifecycleEvent::Cancelled { + reason: RunCancelReason::External, + message: None, + usage: usage(), + }, + LifecycleEvent::Failed { + error: "boom".to_string(), + usage: usage(), + }, + ]; + for event in terminal { + assert!(event.ends_run()); + let json = serde_json::to_value(&event).unwrap(); + assert_eq!(json["prompt_tokens"], 10, "{json}"); + assert_eq!(json["completion_tokens"], 5, "{json}"); + assert_eq!(json["total_tokens"], 15, "{json}"); + } + } + + /// A producer that writes an absent liveness as `null` gets the same + /// default as one that omits it. + #[test] + fn a_start_with_a_null_liveness_reads_as_cancel_at_once() { + let parsed: LifecycleEvent = serde_json::from_value(json!({ + "type": "started", + "agent": "sre", + "prompt": "hi", + "timeout_ms": null, + "liveness": null + })) + .unwrap(); + + let LifecycleEvent::Started { + liveness, timeout, .. + } = parsed + else { + panic!("expected a start"); + }; + assert_eq!(liveness, Liveness::default()); + assert_eq!(timeout, None); + } + + #[test] + fn ids_serialize_as_bare_strings() { + assert_eq!(serde_json::to_value(run_id()).unwrap(), json!(RUN)); + assert_eq!( + serde_json::to_value(SessionId::new("sess_1")).unwrap(), + json!("sess_1") + ); + assert_eq!( + serde_json::to_value(ObserverId::new("obs_1").unwrap()).unwrap(), + json!("obs_1") + ); + assert_eq!( + serde_json::to_value(CheckpointRef::new("parked/run_1.json").unwrap()).unwrap(), + json!("parked/run_1.json") + ); + } + + /// A run id is a UUID on the way in as well as out, so an id minted as + /// anything else is refused at the boundary rather than carried along. + #[test] + fn a_run_id_that_is_not_a_uuid_is_refused() { + assert!(serde_json::from_value::(json!("req_1")).is_err()); + assert!("req_1".parse::().is_err()); + assert_eq!( + serde_json::from_value::(json!(RUN)).unwrap(), + run_id() + ); + } + + /// The nil UUID is the UUID that names nothing, so it is refused on every + /// path a run id can be built by. + #[test] + fn the_nil_uuid_is_not_a_run_id() { + const NIL: &str = "00000000-0000-0000-0000-000000000000"; + + assert_eq!(NIL.parse::(), Err(crate::InvalidRunId::Nil)); + assert_eq!( + RunId::try_from(uuid::Uuid::nil()), + Err(crate::InvalidRunId::Nil) + ); + let err = serde_json::from_value::(json!(NIL)).unwrap_err(); + assert!(err.to_string().contains("nil UUID"), "got: {err}"); + } + + /// An observer or checkpoint named by the empty string names nothing, so + /// neither the constructor nor deserialization will build one. + #[test] + fn an_empty_observer_or_checkpoint_id_is_refused() { + let err = ObserverId::new("").unwrap_err(); + assert_eq!(err.to_string(), "ObserverId cannot be empty"); + assert!(CheckpointRef::new(String::new()).is_err()); + assert!(ObserverId::try_from("").is_err()); + + assert!(serde_json::from_value::(json!("")).is_err()); + assert!(serde_json::from_value::(json!("")).is_err()); + assert_eq!( + serde_json::from_value::(json!("obs_1")).unwrap(), + "obs_1" + ); + } + + /// The refusal reaches the envelope: an attach naming an empty observer + /// does not parse into an event. + #[test] + fn an_event_naming_an_empty_observer_does_not_parse() { + let parsed = serde_json::from_value::(json!({ + "session_id": "sess_1", + "seq": 1, + "at": 0, + "payload": { + "kind": "lifecycle", + "event": { + "type": "observer_attached", + "observer": { "id": "", "kind": "collecting", "presence": false } + } + } + })); + assert!(parsed.is_err()); + } + + /// Version 7 puts the mint time in the leading bits, which is what lets + /// ids order by when their runs started. + #[test] + fn a_minted_run_id_is_version_7() { + assert_eq!(RunId::mint().as_uuid().get_version_num(), 7); + assert_ne!(RunId::mint(), RunId::mint()); + } + + /// A consumer that closes a run's view on its last event must agree with + /// the vocabulary about which events those are. + #[test] + fn only_a_terminal_lifecycle_event_ends_a_run() { + let ending = [ + LifecycleEvent::Parked { + checkpoint: CheckpointRef::new("c").unwrap(), + usage: usage(), + }, + LifecycleEvent::Finished { usage: usage() }, + LifecycleEvent::Cancelled { + reason: RunCancelReason::External, + message: None, + usage: usage(), + }, + LifecycleEvent::Failed { + error: "boom".to_string(), + usage: usage(), + }, + ]; + for event in ending { + assert!(event.ends_run(), "{event:?} should end its run"); + assert!(lifecycle(1, event).ends_run()); + } + + let ongoing = [ + LifecycleEvent::Started { + agent: "sre".to_string(), + prompt: "hi".to_string(), + timeout: None, + liveness: Liveness::default(), + continues: None, + }, + LifecycleEvent::ObserverAttached { + observer: observer(ObserverKind::Collecting, false), + }, + LifecycleEvent::ObserverDetached { + observer: observer(ObserverKind::Collecting, false), + because: DetachCause::Released, + }, + LifecycleEvent::ClaimsExhausted, + LifecycleEvent::LivenessDecided { + policy: LivenessPolicy::Continue, + }, + ]; + for event in ongoing { + assert!(!event.ends_run(), "{event:?} should not end its run"); + } + + let activity = envelope( + 1, + Some(run_id()), + SessionEventPayload::Agent(AgentEvent::new( + AgentContext::single_agent(), + AgentEventPayload::TextDelta { + content: "done".to_string(), + }, + )), + ); + assert!( + !activity.ends_run(), + "an agent saying it is done is not the run ending" + ); + } + + /// The agent's `RunParked` and the lifecycle `Parked` both appear in a + /// parked run's stream, and neither is mistaken for the other. + #[test] + fn the_agents_run_parked_and_the_lifecycle_parked_are_distinct() { + let from_agent = envelope( + 7, + Some(run_id()), + SessionEventPayload::Agent(AgentEvent::new( + AgentContext::coordinator(), + AgentEventPayload::RunParked { + run_id: RUN.to_string(), + decision_ids: vec!["d1".to_string()], + expires_at: "2026-09-24T14:00:00Z".to_string(), + iteration: 2, + }, + )), + ); + let from_owner = lifecycle( + 8, + LifecycleEvent::Parked { + checkpoint: CheckpointRef::new("memory/sess_1/parked/run_1.json").unwrap(), + usage: usage(), + }, + ); + + let agent_json = serde_json::to_value(&from_agent).unwrap(); + let owner_json = serde_json::to_value(&from_owner).unwrap(); + assert_eq!(agent_json["payload"]["kind"], "agent"); + assert_eq!( + agent_json["payload"]["event"]["payload"]["type"], + "run_parked" + ); + assert_eq!(owner_json["payload"]["kind"], "lifecycle"); + assert_eq!(owner_json["payload"]["event"]["type"], "parked"); + + assert!(!from_agent.ends_run()); + assert!(from_owner.ends_run()); + } + + /// A reader older than its producer keeps the stream: a variant it does + /// not know reads as `Unknown`, and the event's `seq` still counts. + #[test] + fn an_unknown_variant_reads_as_unknown_and_keeps_the_envelope() { + let event: SessionEvent = serde_json::from_value(json!({ + "session_id": "sess_1", + "run_id": RUN, + "seq": 7, + "at": 0, + "payload": { + "kind": "lifecycle", + "event": { "type": "suspended", "until_ms": 5 } + } + })) + .unwrap(); + assert_eq!(event.seq.get(), 7); + assert!(matches!( + event.payload, + SessionEventPayload::Lifecycle(LifecycleEvent::Unknown) + )); + assert!( + !event.ends_run(), + "an unknown event is not known to end its run" + ); + + let cancelled: LifecycleEvent = serde_json::from_value(json!({ + "type": "cancelled", + "reason": "preempted", + "prompt_tokens": 10, + "completion_tokens": 5, + "total_tokens": 15 + })) + .unwrap(); + assert!(matches!( + cancelled, + LifecycleEvent::Cancelled { + reason: RunCancelReason::Unknown, + message: None, + .. + } + )); + + let started: LifecycleEvent = serde_json::from_value(json!({ + "type": "started", + "agent": "sre", + "prompt": "hi", + "liveness": { "policy": "hibernate", "grace_ms": 1 } + })) + .unwrap(); + assert!(matches!( + started, + LifecycleEvent::Started { + liveness: Liveness { + policy: LivenessPolicy::Unknown, + .. + }, + .. + } + )); + + let detached: LifecycleEvent = serde_json::from_value(json!({ + "type": "observer_detached", + "observer": { "id": "obs_1", "kind": "mirroring", "presence": false }, + "because": "preempted" + })) + .unwrap(); + assert!(matches!( + detached, + LifecycleEvent::ObserverDetached { + observer: Observer { + kind: ObserverKind::Unknown, + .. + }, + because: DetachCause::Unknown + } + )); + } + + /// An `Unknown` is read, never written: a relay forwards the bytes it + /// received rather than re-encoding an event it could not read. + #[test] + fn an_unknown_variant_cannot_be_serialized() { + assert!(serde_json::to_value(LifecycleEvent::Unknown).is_err()); + assert!(serde_json::to_value(RunCancelReason::Unknown).is_err()); + assert!(serde_json::to_value(LivenessPolicy::Unknown).is_err()); + assert!(serde_json::to_value(ObserverKind::Unknown).is_err()); + assert!(serde_json::to_value(DetachCause::Unknown).is_err()); + + let event = envelope( + 7, + Some(run_id()), + SessionEventPayload::Lifecycle(LifecycleEvent::Unknown), + ); + assert!(serde_json::to_value(&event).is_err()); + } + + /// The payload's `kind` is the one tag with no fallback. + #[test] + fn an_unknown_payload_kind_is_refused() { + let parsed = serde_json::from_value::(json!({ + "session_id": "sess_1", + "seq": 1, + "at": 0, + "payload": { "kind": "telemetry", "event": {} } + })); + assert!(parsed.is_err()); + } +} diff --git a/crates/aura/src/config.rs b/crates/aura/src/config.rs index 3d9550074..7e19c726d 100644 --- a/crates/aura/src/config.rs +++ b/crates/aura/src/config.rs @@ -11,7 +11,6 @@ use crate::hitl::HitlRuntime; use crate::scratchpad::ScratchpadToolsConfig; use crate::tool_wrapper::{ToolCallContext, ToolWrapper}; use aura_config::GlobPattern; -use serde::{Deserialize, Serialize}; use std::sync::Arc; // Re-export the pure config types from `aura-config` so existing @@ -40,28 +39,7 @@ pub enum WorkerSkills { Override(Vec), } -/// Identifier for a chat session — the conversational context an agent runs in. -/// -/// Threaded from the web server's `chat_session_id`, but meaningful for any -/// agent run, including library use without the web layer; not every run has -/// one. An opaque, branded string. Serializes as the bare string. -#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)] -#[serde(transparent)] -pub struct SessionId(String); - -impl SessionId { - /// Wrap a session-id string. Accepts a `&str` or an owned `String`. - #[must_use] - pub fn new(value: impl Into) -> Self { - Self(value.into()) - } - - /// Borrow the underlying id as a string slice. - #[must_use] - pub fn as_str(&self) -> &str { - &self.0 - } -} +pub use aura_events::SessionId; /// Runtime build context for constructing agents. /// diff --git a/crates/aura/src/hooks.rs b/crates/aura/src/hooks.rs index 2e2e13100..2f34f7e51 100644 --- a/crates/aura/src/hooks.rs +++ b/crates/aura/src/hooks.rs @@ -9,21 +9,12 @@ use std::sync::Arc; use async_trait::async_trait; +pub use aura_events::run::RunCancelReason; + pub struct ToolCall<'a> { pub name: &'a str, } -/// Why a hook ended a run. -#[derive(Clone, Copy, Debug, PartialEq, Eq)] -pub enum RunCancelReason { - /// The run outlived the bound it was started with. - Deadline { after: std::time::Duration }, - /// Something outside the run cancelled it. - External, - /// The model called a tool the caller executes, so the run yields to it. - ClientTool, -} - /// Something watching one run. /// /// A hook belongs to the run whose [`Hooks`] holds it and is registered when diff --git a/crates/aura/src/streaming_request_hook.rs b/crates/aura/src/streaming_request_hook.rs index bfe05f29b..44c9909a3 100644 --- a/crates/aura/src/streaming_request_hook.rs +++ b/crates/aura/src/streaming_request_hook.rs @@ -476,6 +476,7 @@ impl StreamingRequestHook { RunCancelReason::ClientTool => { tracing::info!("Run yielding to a client tool during {}", context) } + other => tracing::info!("Run cancelled ({:?}) during {}", other, context), } cancel_sig.cancel(); }