From a153df9c43d0cb2faf1bbfb17bb10dc25227b884 Mon Sep 17 00:00:00 2001 From: Justin Gross Date: Thu, 24 Sep 2026 16:44:48 -0400 Subject: [PATCH 01/10] feat(events): give a session an event envelope carrying run lifecycle MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit AgentEvent says what an agent did and nothing about the run it belongs to. Observers reconstruct the run's story from the transport — a stream closing, an A2A status update, a [DONE] — and no two agree. aura_events::run adds SessionEvent: the session id, the run id when the event belongs to a run, a sequence number dense per session across its runs, a timestamp, and either an AgentEvent or a LifecycleEvent (started, observer attached and detached, claims exhausted, liveness decided, parked, finished, cancelled, failed). One stream per session means an agent's history loads as one stream, and a run boundary is a position in it rather than a second stream. The payload is adjacently tagged, kind beside event, so no field of an event can collide with the envelope's tag. Lifecycle events stay internally tagged by type, like AgentEventPayload; the roundtrip test over every variant is what catches a field named after a tag. Internal tagging and flatten need a self-describing format, which the module doc records. RunId is a UUID, v7 when minted, so it agrees with the orchestration run id HITL parses and the park owner key; aura-events takes uuid for it. SessionId and RunCancelReason move into aura-events, re-exported from aura::config and aura::hooks. RunCancelReason is non_exhaustive and gains Unclaimed and Shutdown, so a liveness or shutdown cancel is not reported as External. ObserverDetached says why with a DetachCause, and Started carries the run's Liveness, read as cancel-at-once when an older payload omits it. No producer emits it yet; this is the type and its tests. Ref: GH-778 Ref: GH-578 --- CLAUDE.md | 3 +- Cargo.lock | 1 + crates/aura-events/Cargo.toml | 1 + crates/aura-events/src/lib.rs | 61 +- crates/aura-events/src/run.rs | 903 ++++++++++++++++++++++ crates/aura/src/config.rs | 24 +- crates/aura/src/hooks.rs | 13 +- crates/aura/src/streaming_request_hook.rs | 1 + 8 files changed, 971 insertions(+), 36 deletions(-) create mode 100644 crates/aura-events/src/run.rs 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-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/lib.rs b/crates/aura-events/src/lib.rs index e39dedc94..7bc5b26b3 100644 --- a/crates/aura-events/src/lib.rs +++ b/crates/aura-events/src/lib.rs @@ -10,13 +10,16 @@ //! //! - [`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 event_names; pub mod orchestration; +pub mod run; use serde::{Deserialize, Serialize}; use std::collections::BTreeMap; @@ -135,6 +138,62 @@ macro_rules! string_newtype { }; } +// Path-based so submodules can invoke it regardless of declaration order. +pub(crate) use string_newtype; + +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(transparent)] +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 + } +} + +impl From for RunId { + fn from(uuid: uuid::Uuid) -> Self { + Self(uuid) + } +} + +impl From for uuid::Uuid { + fn from(id: RunId) -> Self { + id.0 + } +} + +impl std::str::FromStr for RunId { + type Err = uuid::Error; + + fn from_str(s: &str) -> Result { + uuid::Uuid::from_str(s).map(Self) + } +} + +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 diff --git a/crates/aura-events/src/run.rs b/crates/aura-events/src/run.rs new file mode 100644 index 000000000..d842bd012 --- /dev/null +++ b/crates/aura-events/src/run.rs @@ -0,0 +1,903 @@ +//! 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. The envelope +//! never carries a second id for the same run. +//! +//! # 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. +//! +//! # 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. +//! +//! 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::{string_newtype, RunId, SessionId, TokenUsage}; + +/// 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) + } +} + +/// An event's position in its session's stream. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)] +#[serde(transparent)] +pub struct SequenceNumber(u64); + +impl SequenceNumber { + /// The number of a session's first event. + pub const FIRST: Self = Self(1); + + pub fn get(self) -> u64 { + self.0 + } + + /// The number the event after this one carries. + pub fn next(self) -> Self { + Self(self.0.saturating_add(1)) + } + + /// Whether this is the event immediately after `previous`. `false` means + /// the consumer missed at least one event between them. + 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) + } +} + +string_newtype! { + /// Identifier for one observer of a session, unique among its observers. + ObserverId +} + +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")] +pub enum ObserverKind { + /// Reads the stream without a claim on any run. + Collecting, + /// Holds the live run's right to continue. + Claiming, +} + +/// 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")] +pub enum DetachCause { + /// The observer let go. + Released, + /// A lease held across a process boundary lapsed unrenewed. + Expired, + /// Another claimant took over. + Displaced, +} + +/// 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, +} + +/// 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 = "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 = "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, +} + +/// 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 = "duration_ms::option" + )] + timeout: Option, + #[serde(default)] + liveness: Liveness, + }, + + ObserverAttached { + observer: Observer, + }, + + ObserverDetached { + observer: Observer, + because: DetachCause, + }, + + /// The last claiming observer detached. + ClaimsExhausted, + + LivenessDecided { + policy: LivenessPolicy, + }, + + /// The run stopped resumable. + Parked { + checkpoint: CheckpointRef, + }, + + /// 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, + }, + + /// The run ended in an error. + Failed { + error: String, + }, +} + +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, + } + } +} + +/// Serializes a [`Duration`] as whole milliseconds, the unit every other +/// duration in this crate is expressed in on the wire. +mod duration_ms { + 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 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(seq), + 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"), + 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), + }, + }, + )) + .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(), + }, + 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"), + }, + LifecycleEvent::Finished { usage: usage() }, + LifecycleEvent::Cancelled { + reason: RunCancelReason::Deadline { + after: Duration::from_millis(1500), + }, + message: None, + }, + LifecycleEvent::Cancelled { + reason: RunCancelReason::External, + message: Some("operator stopped it".to_string()), + }, + LifecycleEvent::Cancelled { + reason: RunCancelReason::ClientTool, + message: None, + }, + LifecycleEvent::Cancelled { + reason: RunCancelReason::Unclaimed, + message: None, + }, + LifecycleEvent::Cancelled { + reason: RunCancelReason::Shutdown, + message: None, + }, + LifecycleEvent::Failed { + error: "provider returned 500".to_string(), + }, + ]; + + 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()), + }, + )) + .unwrap(); + + assert_eq!( + json["payload"]["event"], + json!({ + "type": "cancelled", + "reason": "deadline", + "after_ms": 300_000, + "message": "request timeout" + }) + ); + + 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, + }, + )) + .unwrap(); + assert_eq!( + json["payload"]["event"], + json!({ "type": "cancelled", "reason": tag }) + ); + } + } + + /// 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); + } + + #[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(), + json!("obs_1") + ); + assert_eq!( + serde_json::to_value(CheckpointRef::new("parked/run_1.json")).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() + ); + } + + /// 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()); + } + + /// Dense numbering is what lets a consumer tell a missed event from a + /// quiet session: 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, 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(), + 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" + ); + } + + /// 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"), + }, + LifecycleEvent::Finished { usage: usage() }, + LifecycleEvent::Cancelled { + reason: RunCancelReason::External, + message: None, + }, + LifecycleEvent::Failed { + error: "boom".to_string(), + }, + ]; + 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(), + }, + 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"), + }, + ); + + 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()); + } +} 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(); } From 58337b412d1286955c0c7c0af13155f71ab95df5 Mon Sep 17 00:00:00 2001 From: Justin Gross Date: Thu, 8 Oct 2026 18:26:55 -0400 Subject: [PATCH 02/10] feat(events): refuse empty and nil identifiers An identifier that names nothing should not be constructible. The ones this PR introduces now refuse their empty value on every path they can be built by: - ObserverId and CheckpointRef come from nonempty_string_newtype!, whose constructor returns a Result, with TryFrom the string types and deserialization that refuses "". - RunId refuses the nil UUID through TryFrom, FromStr and deserialization, and its parse error is InvalidRunId. - SequenceNumber refuses 0, since a session's sequence starts at 1. The string ids that predate this PR (ToolCallId, ToolName, ToolNamespace, SessionId) keep string_newtype!; hardening them reaches well past the envelope. Ref: GH-778 --- crates/aura-events/src/lib.rs | 143 ++++++++++++++++++++++++++++++++-- crates/aura-events/src/run.rs | 113 ++++++++++++++++++++++++--- 2 files changed, 239 insertions(+), 17 deletions(-) diff --git a/crates/aura-events/src/lib.rs b/crates/aura-events/src/lib.rs index 7bc5b26b3..9aa672b2e 100644 --- a/crates/aura-events/src/lib.rs +++ b/crates/aura-events/src/lib.rs @@ -138,8 +138,102 @@ 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 string_newtype; +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 @@ -154,7 +248,7 @@ string_newtype! { /// [`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(transparent)] +#[serde(try_from = "uuid::Uuid", into = "uuid::Uuid")] pub struct RunId(uuid::Uuid); impl RunId { @@ -168,9 +262,42 @@ impl RunId { } } -impl From for RunId { - fn from(uuid: uuid::Uuid) -> Self { - Self(uuid) +/// 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)) + } } } @@ -181,10 +308,12 @@ impl From for uuid::Uuid { } impl std::str::FromStr for RunId { - type Err = uuid::Error; + type Err = InvalidRunId; fn from_str(s: &str) -> Result { - uuid::Uuid::from_str(s).map(Self) + uuid::Uuid::from_str(s) + .map_err(InvalidRunId::Malformed)? + .try_into() } } diff --git a/crates/aura-events/src/run.rs b/crates/aura-events/src/run.rs index d842bd012..634eebe24 100644 --- a/crates/aura-events/src/run.rs +++ b/crates/aura-events/src/run.rs @@ -66,7 +66,7 @@ use std::time::Duration; use serde::{Deserialize, Serialize}; use crate::agent::AgentEvent; -use crate::{string_newtype, RunId, SessionId, TokenUsage}; +use crate::{nonempty_string_newtype, RunId, SessionId, TokenUsage}; /// An instant as Unix time in milliseconds. #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)] @@ -100,7 +100,7 @@ impl std::fmt::Display for Timestamp { /// An event's position in its session's stream. #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)] -#[serde(transparent)] +#[serde(try_from = "u64", into = "u64")] pub struct SequenceNumber(u64); impl SequenceNumber { @@ -129,12 +129,42 @@ impl std::fmt::Display for SequenceNumber { } } -string_newtype! { +/// A sequence number of zero, which no event carries. +#[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("a session's 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 + } +} + +nonempty_string_newtype! { /// Identifier for one observer of a session, unique among its observers. ObserverId } -string_newtype! { +nonempty_string_newtype! { /// Locates a parked run's checkpoint, as the park machinery names it. CheckpointRef } @@ -422,7 +452,7 @@ mod tests { fn observer(kind: ObserverKind, presence: bool) -> Observer { Observer { - id: ObserverId::new("obs_1"), + id: ObserverId::new("obs_1").unwrap(), kind, presence, } @@ -551,7 +581,7 @@ mod tests { policy: LivenessPolicy::Park, }, LifecycleEvent::Parked { - checkpoint: CheckpointRef::new("memory/sess_1/parked/run_1.json"), + checkpoint: CheckpointRef::new("memory/sess_1/parked/run_1.json").unwrap(), }, LifecycleEvent::Finished { usage: usage() }, LifecycleEvent::Cancelled { @@ -733,11 +763,11 @@ mod tests { json!("sess_1") ); assert_eq!( - serde_json::to_value(ObserverId::new("obs_1")).unwrap(), + 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(), + serde_json::to_value(CheckpointRef::new("parked/run_1.json").unwrap()).unwrap(), json!("parked/run_1.json") ); } @@ -754,6 +784,69 @@ mod tests { ); } + /// 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()); + } + + /// 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::(json!(0)).is_err()); + assert_eq!( + serde_json::from_value::(json!(1)).unwrap(), + SequenceNumber::FIRST + ); + } + /// Version 7 puts the mint time in the leading bits, which is what lets /// ids order by when their runs started. #[test] @@ -808,7 +901,7 @@ mod tests { fn only_a_terminal_lifecycle_event_ends_a_run() { let ending = [ LifecycleEvent::Parked { - checkpoint: CheckpointRef::new("c"), + checkpoint: CheckpointRef::new("c").unwrap(), }, LifecycleEvent::Finished { usage: usage() }, LifecycleEvent::Cancelled { @@ -883,7 +976,7 @@ mod tests { let from_owner = lifecycle( 8, LifecycleEvent::Parked { - checkpoint: CheckpointRef::new("memory/sess_1/parked/run_1.json"), + checkpoint: CheckpointRef::new("memory/sess_1/parked/run_1.json").unwrap(), }, ); From 1e7fd932e2c1051d025986a68a2c7d85ade241e4 Mon Sep 17 00:00:00 2001 From: Justin Gross Date: Fri, 9 Oct 2026 09:59:40 -0400 Subject: [PATCH 03/10] refactor(events): share timestamp, sequence number, and duration_ms Timestamp, SequenceNumber, and the duration_ms serde module are not specific to a session's envelope, so they move from the run module to the crate root, beside the other shared value types. duration_ms becomes a public module so other crates can serialize a Duration as whole milliseconds without keeping their own copy, and it gains tests of its own, including the Option form and saturation past u64. SequenceNumber's docs describe any dense sequence starting at 1; the per-session rules stay on SessionEvent and the run module doc. That doc now notes that `at` is Unix milliseconds while other instants on the wire are RFC 3339 strings. Ref: GH-778 --- crates/aura-events/src/duration_ms.rs | 97 +++++++++++++ crates/aura-events/src/lib.rs | 144 +++++++++++++++++++ crates/aura-events/src/run.rs | 196 ++------------------------ 3 files changed, 249 insertions(+), 188 deletions(-) create mode 100644 crates/aura-events/src/duration_ms.rs 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 9aa672b2e..9884eba95 100644 --- a/crates/aura-events/src/lib.rs +++ b/crates/aura-events/src/lib.rs @@ -17,6 +17,7 @@ //! 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; @@ -425,6 +426,97 @@ 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 at least one position between them was missed. + 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. @@ -1591,6 +1683,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 index 634eebe24..22c481981 100644 --- a/crates/aura-events/src/run.rs +++ b/crates/aura-events/src/run.rs @@ -34,7 +34,9 @@ //! by a relay or projection. //! //! Order by `seq`, never by `at`. The [`Timestamp`] is wall-clock time, for -//! display. +//! 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 //! @@ -66,98 +68,7 @@ use std::time::Duration; use serde::{Deserialize, Serialize}; use crate::agent::AgentEvent; -use crate::{nonempty_string_newtype, RunId, SessionId, TokenUsage}; - -/// 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) - } -} - -/// An event's position in its session's stream. -#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)] -#[serde(try_from = "u64", into = "u64")] -pub struct SequenceNumber(u64); - -impl SequenceNumber { - /// The number of a session's first event. - pub const FIRST: Self = Self(1); - - pub fn get(self) -> u64 { - self.0 - } - - /// The number the event after this one carries. - pub fn next(self) -> Self { - Self(self.0.saturating_add(1)) - } - - /// Whether this is the event immediately after `previous`. `false` means - /// the consumer missed at least one event between them. - 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 event carries. -#[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("a session's 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 - } -} +use crate::{nonempty_string_newtype, RunId, SequenceNumber, SessionId, Timestamp, TokenUsage}; nonempty_string_newtype! { /// Identifier for one observer of a session, unique among its observers. @@ -218,7 +129,7 @@ pub enum LivenessPolicy { pub struct Liveness { pub policy: LivenessPolicy, /// How long the run goes unclaimed before `policy` acts. - #[serde(rename = "grace_ms", with = "duration_ms")] + #[serde(rename = "grace_ms", with = "crate::duration_ms")] pub grace: Duration, } @@ -239,7 +150,7 @@ impl Default for Liveness { pub enum RunCancelReason { /// The run outlived the bound it was started with. Deadline { - #[serde(rename = "after_ms", with = "duration_ms")] + #[serde(rename = "after_ms", with = "crate::duration_ms")] after: Duration, }, /// Something outside the run cancelled it. @@ -270,7 +181,7 @@ pub enum LifecycleEvent { rename = "timeout_ms", default, skip_serializing_if = "Option::is_none", - with = "duration_ms::option" + with = "crate::duration_ms::option" )] timeout: Option, #[serde(default)] @@ -370,45 +281,6 @@ impl SessionEvent { } } -/// Serializes a [`Duration`] as whole milliseconds, the unit every other -/// duration in this crate is expressed in on the wire. -mod duration_ms { - 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 super::*; @@ -426,7 +298,7 @@ mod tests { SessionEvent { session_id: SessionId::new("sess_1"), run_id: run, - seq: SequenceNumber(seq), + seq: SequenceNumber::try_from(seq).expect("a sequence starts at 1"), at: Timestamp::from_unix_millis(1_700_000_000_000), payload, } @@ -835,18 +707,6 @@ mod tests { assert!(parsed.is_err()); } - /// 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::(json!(0)).is_err()); - assert_eq!( - serde_json::from_value::(json!(1)).unwrap(), - SequenceNumber::FIRST - ); - } - /// Version 7 puts the mint time in the leading bits, which is what lets /// ids order by when their runs started. #[test] @@ -855,46 +715,6 @@ mod tests { assert_ne!(RunId::mint(), RunId::mint()); } - /// Dense numbering is what lets a consumer tell a missed event from a - /// quiet session: 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, 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(), - 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" - ); - } - /// A consumer that closes a run's view on its last event must agree with /// the vocabulary about which events those are. #[test] From 5db77be51b183ebc34eb0765e3a2068621b78fd3 Mon Sep 17 00:00:00 2001 From: Justin Gross Date: Fri, 9 Oct 2026 10:01:26 -0400 Subject: [PATCH 04/10] refactor(cli): use the shared duration_ms serde module The CLI kept its own duration_millis module, the same whole-millisecond encoding aura-events now exports as duration_ms. Saved REPL events use the shared module instead. The format is unchanged for any duration that fits a u64 of milliseconds; a longer one saturates instead of being written as a number the reader cannot parse back. Ref: GH-778 --- crates/aura-cli/src/api/types.rs | 31 ++++++++----------------------- 1 file changed, 8 insertions(+), 23 deletions(-) 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)); } } From 2a2527634e187482263a7bfee0ccdc949f766fa4 Mon Sep 17 00:00:00 2001 From: Justin Gross Date: Fri, 9 Oct 2026 17:17:30 -0400 Subject: [PATCH 05/10] feat(events): read an unknown variant as unknown rather than failing A reader older than its producer failed the whole SessionEvent on a variant it did not know, losing the event's seq, so the next event looked like a gap. The lifecycle enums, the observer and detach enums, the cancel reason, the liveness policy, and AgentEventPayload now read any other tag as an Unknown variant that cannot be written back, so a relay forwards the bytes it received rather than a re-encoded event. ObserverKind and DetachCause become non_exhaustive to match. SessionEventPayload keeps a closed kind: serde cannot fall back on an adjacent tag that carries content, and its two kinds are not expected to grow. Ref: GH-778 --- crates/aura-events/src/agent.rs | 18 +++++ crates/aura-events/src/run.rs | 134 ++++++++++++++++++++++++++++++++ 2 files changed, 152 insertions(+) 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/run.rs b/crates/aura-events/src/run.rs index 22c481981..48cba30c6 100644 --- a/crates/aura-events/src/run.rs +++ b/crates/aura-events/src/run.rs @@ -48,6 +48,13 @@ //! [`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. @@ -83,11 +90,16 @@ nonempty_string_newtype! { /// 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. @@ -102,6 +114,7 @@ pub struct Observer { /// 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, @@ -109,6 +122,10 @@ pub enum DetachCause { 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. @@ -122,6 +139,10 @@ pub enum LivenessPolicy { 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. @@ -161,6 +182,10 @@ pub enum RunCancelReason { 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 @@ -229,6 +254,11 @@ pub enum LifecycleEvent { Failed { error: String, }, + + /// 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, } impl LifecycleEvent { @@ -813,4 +843,108 @@ mod tests { 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" + })) + .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()); + } } From 778eaa41dd2675bf6a14d3ca8a68e73ef04bded4 Mon Sep 17 00:00:00 2001 From: Justin Gross Date: Fri, 9 Oct 2026 17:23:42 -0400 Subject: [PATCH 06/10] fix(events): read a null liveness as the default Started.liveness defaulted only when the field was missing; an explicit null failed the event, so a producer writing absent optionals as null broke the older-payload fallback. A null now reads as the default, as a missing field does. Ref: GH-778 --- crates/aura-events/src/run.rs | 34 +++++++++++++++++++++++++++++++++- 1 file changed, 33 insertions(+), 1 deletion(-) diff --git a/crates/aura-events/src/run.rs b/crates/aura-events/src/run.rs index 48cba30c6..74ea39357 100644 --- a/crates/aura-events/src/run.rs +++ b/crates/aura-events/src/run.rs @@ -209,7 +209,7 @@ pub enum LifecycleEvent { with = "crate::duration_ms::option" )] timeout: Option, - #[serde(default)] + #[serde(default, deserialize_with = "null_as_default")] liveness: Liveness, }, @@ -261,6 +261,15 @@ pub enum LifecycleEvent { 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. @@ -657,6 +666,29 @@ mod tests { assert_eq!(liveness.grace, Duration::ZERO); } + /// 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)); From 53e19a1eca3ccc0e86f5d1249900ec762b81dc98 Mon Sep 17 00:00:00 2001 From: Justin Gross Date: Fri, 9 Oct 2026 18:01:11 -0400 Subject: [PATCH 07/10] doc(events): say what a false from follows() covers `follows` answers whether a position is exactly one past the previous one, so a repeat and an out-of-order arrival return false as a gap does. The doc named only the gap, which would have a consumer resync on a redelivery. Ref: GH-778 --- crates/aura-events/src/lib.rs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/crates/aura-events/src/lib.rs b/crates/aura-events/src/lib.rs index 9884eba95..e1279875a 100644 --- a/crates/aura-events/src/lib.rs +++ b/crates/aura-events/src/lib.rs @@ -475,7 +475,8 @@ impl SequenceNumber { } /// Whether this is the position immediately after `previous`. `false` - /// means at least one position between them was missed. + /// 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) } From fa9af6fcfa9e1c43f126a734087dd4eac9ba9a9d Mon Sep 17 00:00:00 2001 From: Justin Gross Date: Fri, 9 Oct 2026 18:07:24 -0400 Subject: [PATCH 08/10] feat(events): let a start name the run it continues A park ends its run, and a resume is a new run that continues the parked one, in the shape GH-785's resume route gives it. Nothing on the stream linked the two, so a consumer that closed a run on Parked could not tell the resumed run's Started from a fresh one. Started now carries the run it continues, when it continues one; a start that does not carries no such field, so an older payload reads as a fresh run. Ref: GH-778 Ref: GH-785 --- crates/aura-events/src/run.rs | 38 +++++++++++++++++++++++++++++++++++ 1 file changed, 38 insertions(+) diff --git a/crates/aura-events/src/run.rs b/crates/aura-events/src/run.rs index 74ea39357..de061d51c 100644 --- a/crates/aura-events/src/run.rs +++ b/crates/aura-events/src/run.rs @@ -211,6 +211,9 @@ pub enum LifecycleEvent { 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 { @@ -422,6 +425,7 @@ mod tests { policy: LivenessPolicy::Continue, grace: Duration::from_secs(30), }, + continues: None, }, )) .expect("should serialize"); @@ -471,6 +475,7 @@ mod tests { prompt: "hi".to_string(), timeout: None, liveness: Liveness::default(), + continues: None, }, LifecycleEvent::ObserverAttached { observer: observer(ObserverKind::Claiming, true), @@ -666,6 +671,38 @@ mod tests { 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()); + } + /// A producer that writes an absent liveness as `null` gets the same /// default as one that omits it. #[test] @@ -805,6 +842,7 @@ mod tests { prompt: "hi".to_string(), timeout: None, liveness: Liveness::default(), + continues: None, }, LifecycleEvent::ObserverAttached { observer: observer(ObserverKind::Collecting, false), From 85f0a05dede2c2d0d6136f25b6a9e8570b47523e Mon Sep 17 00:00:00 2001 From: Justin Gross Date: Fri, 9 Oct 2026 18:29:32 -0400 Subject: [PATCH 09/10] feat(events): carry usage on every terminal lifecycle event Only Finished carried the run's cumulative provider-billed tokens, so a consumer reading cost from the terminal event got nothing for a run that was cancelled, failed, or parked after billing turns. Parked, Cancelled, and Failed now carry the same flattened usage as Finished. It is required rather than optional: a flattened Option reads a malformed usage as None and hides a producer bug, and a producer always has a cumulative figure, zero included. Ref: GH-778 --- crates/aura-events/src/run.rs | 72 +++++++++++++++++++++++++++++++++-- 1 file changed, 68 insertions(+), 4 deletions(-) diff --git a/crates/aura-events/src/run.rs b/crates/aura-events/src/run.rs index de061d51c..5043fe7c5 100644 --- a/crates/aura-events/src/run.rs +++ b/crates/aura-events/src/run.rs @@ -235,6 +235,9 @@ pub enum LifecycleEvent { /// The run stopped resumable. Parked { checkpoint: CheckpointRef, + /// As [`Finished`](Self::Finished)'s. + #[serde(flatten)] + usage: TokenUsage, }, /// The run completed its work. @@ -251,11 +254,17 @@ pub enum LifecycleEvent { /// 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 @@ -498,6 +507,7 @@ mod tests { }, LifecycleEvent::Parked { checkpoint: CheckpointRef::new("memory/sess_1/parked/run_1.json").unwrap(), + usage: usage(), }, LifecycleEvent::Finished { usage: usage() }, LifecycleEvent::Cancelled { @@ -505,25 +515,31 @@ mod tests { 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(), }, ]; @@ -592,6 +608,7 @@ mod tests { after: Duration::from_secs(300), }, message: Some("request timeout".to_string()), + usage: usage(), }, )) .unwrap(); @@ -602,7 +619,10 @@ mod tests { "type": "cancelled", "reason": "deadline", "after_ms": 300_000, - "message": "request timeout" + "message": "request timeout", + "prompt_tokens": 10, + "completion_tokens": 5, + "total_tokens": 15 }) ); @@ -616,12 +636,19 @@ mod tests { LifecycleEvent::Cancelled { reason, message: None, + usage: usage(), }, )) .unwrap(); assert_eq!( json["payload"]["event"], - json!({ "type": "cancelled", "reason": tag }) + json!({ + "type": "cancelled", + "reason": tag, + "prompt_tokens": 10, + "completion_tokens": 5, + "total_tokens": 15 + }) ); } } @@ -703,6 +730,35 @@ mod tests { 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] @@ -821,14 +877,17 @@ mod tests { 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 { @@ -897,6 +956,7 @@ mod tests { 8, LifecycleEvent::Parked { checkpoint: CheckpointRef::new("memory/sess_1/parked/run_1.json").unwrap(), + usage: usage(), }, ); @@ -941,14 +1001,18 @@ mod tests { let cancelled: LifecycleEvent = serde_json::from_value(json!({ "type": "cancelled", - "reason": "preempted" + "reason": "preempted", + "prompt_tokens": 10, + "completion_tokens": 5, + "total_tokens": 15 })) .unwrap(); assert!(matches!( cancelled, LifecycleEvent::Cancelled { reason: RunCancelReason::Unknown, - message: None + message: None, + .. } )); From 8c0a23f55c545fd780a47ec002df07a2440f2e74 Mon Sep 17 00:00:00 2001 From: Justin Gross Date: Fri, 9 Oct 2026 18:39:56 -0400 Subject: [PATCH 10/10] doc(events): drop the claim that no second run id travels RunParked, an agent event the envelope carries, names its run by orchestration persistence's own id until the runtime mints both, so the module doc overclaimed. A run still has one RunId on the envelope. Ref: GH-778 --- crates/aura-events/src/run.rs | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/crates/aura-events/src/run.rs b/crates/aura-events/src/run.rs index 5043fe7c5..94f667b22 100644 --- a/crates/aura-events/src/run.rs +++ b/crates/aura-events/src/run.rs @@ -21,8 +21,7 @@ //! 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. The envelope -//! never carries a second id for the same run. +//! A run has one id, a [`RunId`] minted once when the run starts. //! //! # Sequence numbers are dense //!