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..8f4c49ba1 --- /dev/null +++ b/crates/aura-events/src/run.rs @@ -0,0 +1,904 @@ +//! 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, + /// and a request-bound run has always ended with its request. + 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 { + /// Whether this ends its run. The session's stream goes on past it. + 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 written before runs carried a liveness still parses, as the + /// policy every run had then: 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-web-server/src/a2a/agent_executor.rs b/crates/aura-web-server/src/a2a/agent_executor.rs index 2fed3d215..a75dcc428 100644 --- a/crates/aura-web-server/src/a2a/agent_executor.rs +++ b/crates/aura-web-server/src/a2a/agent_executor.rs @@ -269,7 +269,11 @@ impl AgentExecutor for AuraAgentExecutor { metadata: None, })); - let request_id = format!("a2a_{}", task_id); + // The run's id, and in string form the request id everything + // request-keyed reads: MCP cancellation, approvals, the cancel map. + let run_id = aura::RunId::mint(); + let request_id = run_id.to_string(); + event!(Level::DEBUG, task_id, %run_id, "a2a task starts its run"); // Registered before the agent build and history fetch, both of which // await, so a cancelTask during those has a token to cancel. Its @@ -295,7 +299,7 @@ impl AgentExecutor for AuraAgentExecutor { Some(&req_headers), session_id, None, - Some(request_id.clone()), + Some(run_id), run_tools, ) .await diff --git a/crates/aura-web-server/src/handlers.rs b/crates/aura-web-server/src/handlers.rs index a5e1ccffd..1f3c6426a 100644 --- a/crates/aura-web-server/src/handlers.rs +++ b/crates/aura-web-server/src/handlers.rs @@ -176,7 +176,7 @@ async fn build_agent_for_request( req_headers: &HashMap, additional_tools: Vec>, client_tools: Option<&[ClientToolDefinition]>, - request_id: String, + run_id: aura::RunId, session_id: String, ) -> Result, PrepareError> { let client_tool_defs = @@ -186,7 +186,7 @@ async fn build_agent_for_request( Some(req_headers), additional_tools, client_tool_defs, - Some(request_id), + Some(run_id), Some(session_id), ) .await @@ -243,11 +243,12 @@ pub async fn prepare_request( // `finish_reason: "tool_calls"` when one fires. let has_client_tools = req.tools.is_some(); - // Generate the request id up front so the agent build (single-agent or - // orchestration) shares one value with the completion stream. The HITL gate - // and approval events stamp this id; previously it was minted later in - // `build_completion_config`, after the agent was already built. - let request_id = format!("req_{}", Uuid::new_v4().simple()); + // The run's id, minted before the agent build (single-agent or + // orchestration) so the build and the completion stream share one value. + // Its string form is the request id every request-keyed registry reads: + // the HITL gate's approvals, their sweep, and MCP cancellation. + let run_id = aura::RunId::mint(); + let request_id = run_id.to_string(); // Find the matching config: single-config passthrough > explicit model > DEFAULT_AGENT // Single-config servers accept any model field value (clients like LibreChat always send one). @@ -323,7 +324,7 @@ pub async fn prepare_request( Some(req_headers_map), Some(chat_session_id.to_string()), client_tools_vec.clone(), - Some(request_id.clone()), + Some(run_id), run_tools, ) .await @@ -351,7 +352,7 @@ pub async fn prepare_request( req_headers_map, additional_tools, client_tools, - request_id.clone(), + run_id, chat_session_id.to_string(), ) .await?; diff --git a/crates/aura-web-server/src/slack/runner.rs b/crates/aura-web-server/src/slack/runner.rs index abbf0d9c7..bd0b3bc36 100644 --- a/crates/aura-web-server/src/slack/runner.rs +++ b/crates/aura-web-server/src/slack/runner.rs @@ -264,7 +264,16 @@ impl SlackIngress { /// hold the bot turn that earned it. A conversation read here is /// checked here. async fn answer(&self, inbound: Inbound, prefetched: Option>) { - let request_id = format!("slack_{}_{}", inbound.channel, inbound.ts); + // The run's id, and in string form the request id everything + // request-keyed reads. The message it answers is logged beside it. + let run_id = aura::RunId::mint(); + let request_id = run_id.to_string(); + debug!( + request_id, + channel = %inbound.channel, + ts = %inbound.ts, + "slack message starts its run" + ); let earlier = match prefetched { Some(earlier) => earlier, None => { @@ -291,7 +300,7 @@ impl SlackIngress { warn!(request_id, error = %e, "could not react to slack message"); } - let reply = match self.run_agent(&inbound, &earlier, &request_id).await { + let reply = match self.run_agent(&inbound, &earlier, run_id).await { Ok(text) if text.trim().is_empty() => EMPTY_REPLY.to_owned(), Ok(text) => text, // The server is going down; a reply would race the shutdown and @@ -344,8 +353,10 @@ impl SlackIngress { &self, inbound: &Inbound, earlier: &[SlackMessage], - request_id: &str, + run_id: aura::RunId, ) -> Result { + let request_id = run_id.to_string(); + let request_id = request_id.as_str(); let history = thread_history(earlier, &self.identity, &inbound.ts); let session_id = match inbound.reply_thread() { Some(thread) => format!("slack:{}:{thread}", inbound.channel), @@ -359,13 +370,7 @@ impl SlackIngress { ); let agent = RigBuilder::new(config, self.state.pending_approvals.clone()) .with_hitl_hmac(self.state.hitl_webhook_hmac.clone()) - .build_streaming_agent_with_tools( - None, - Some(session_id), - None, - Some(request_id.to_owned()), - tools, - ) + .build_streaming_agent_with_tools(None, Some(session_id), None, Some(run_id), tools) .await .map_err(|e| RunError::Build(e.to_string()))?; diff --git a/crates/aura/src/builder.rs b/crates/aura/src/builder.rs index 92868e045..80521d3f4 100644 --- a/crates/aura/src/builder.rs +++ b/crates/aura/src/builder.rs @@ -17,6 +17,7 @@ use crate::{ vector_dynamic::DynamicVectorSearchTool, vector_store::VectorStoreManager, }; +use aura_events::RunId; use aura_events::agent::AgentEvent; use futures::StreamExt; use rig::client::CompletionClient; @@ -210,7 +211,7 @@ impl std::fmt::Debug for PreparedAgent { impl std::fmt::Debug for Agent { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.debug_struct("Agent") - .field("request_id", &self.request_id()) + .field("run_id", &self.run_id()) .field("prepared", &self.prepared) .finish() } @@ -967,8 +968,8 @@ impl PreparedAgent { &self.forwarded_headers } - /// Begin a run of this agent for the request `request_id`, whose headers - /// are `req_headers`. + /// Begin the run `run_id` of this agent, for a request whose headers are + /// `req_headers`. /// /// The run is a [`RunContext`] of its own, on a token of its own, with a /// fresh scratchpad budget and turn-limit counters starting from what @@ -982,7 +983,7 @@ impl PreparedAgent { /// calls as someone else. Prepare an agent for it instead. pub fn begin_run( self: &Arc, - request_id: impl Into, + run_id: RunId, req_headers: Option<&HashMap>, skill_recorder: Option>, ) -> Result { @@ -992,7 +993,7 @@ impl PreparedAgent { }); } let (run, events) = RunContext::channel_for_agent( - request_id.into(), + run_id, CancellationToken::new(), self.scratchpad_budget.as_ref().map(ContextBudget::fresh), self.turn_nudge.as_ref().map(|seed| seed.fresh()), @@ -1039,7 +1040,7 @@ impl PreparedAgent { let mut active = self.active.lock().unwrap_or_else(PoisonError::into_inner); if let Some(alive) = active.upgrade() { return Err(RunInProgress { - active: alive.run().id().to_string(), + active: alive.run().id(), }); } let lease = Arc::new(RunLease::new(Arc::clone(&run))); @@ -1536,7 +1537,7 @@ impl Agent { /// Prepare an agent from configuration and begin its single run. /// /// The run is for the request `config` was resolved for: its id is - /// `config.request_id` (empty when unset), its headers are the ones + /// `config.run_id` (a fresh one when unset), its headers are the ones /// `config.forwarded_headers` recorded, and `config.skill_recorder` /// records its skill-tool invocations. See [`PreparedAgent::prepare`] for /// `additional_tools` and `client_tools`; callers that want to reuse the @@ -1551,7 +1552,7 @@ impl Agent { Arc::new(PreparedAgent::prepare(config, additional_tools, client_tools).await?); let req_headers = config.forwarded_headers.as_request(); Ok(prepared.begin_run( - config.request_id.clone().unwrap_or_default(), + config.run_id.unwrap_or_else(RunId::mint), Some(&req_headers), config.skill_recorder.clone(), )?) @@ -1562,8 +1563,8 @@ impl Agent { self.lease.run() } - /// Request id of this run. - pub fn request_id(&self) -> &str { + /// This run's id. + pub fn run_id(&self) -> RunId { self.lease.run().id() } @@ -1994,7 +1995,7 @@ impl StreamingAgent for Agent { // The run began at `begin_run`, so its context is what this streams // under; the id is the run's. let run = Arc::clone(self.lease.run()); - if request_id != run.id().as_ref() { + if run.id().to_string() != request_id { tracing::debug!( run_id = %run.id(), request_id, @@ -2613,9 +2614,10 @@ mod tests { /// one `add_mcp_tool` stamped, carried through `pre_call`. #[tokio::test] async fn a_namespace_scoped_pattern_gates_a_tool_from_that_server() { - let request_id = "req_ns_gating_match"; let server = RecordingMcpServer::start().await; - let (run, mut rx) = crate::run_context::RunContext::channel(request_id); + let (run, mut rx) = crate::run_context::RunContext::channel( + crate::run_context::named_run_id("req_ns_gating_match"), + ); let agent = compose_gated_agent(&server, "github:*", "github", "list_repos", run).await; // The parked approval expires unanswered; the call's own outcome is @@ -2645,9 +2647,10 @@ mod tests { /// namespace were ignored and the bare name alone matched. #[tokio::test] async fn a_namespace_scoped_pattern_ignores_a_tool_from_another_server() { - let request_id = "req_ns_gating_miss"; let server = RecordingMcpServer::start().await; - let (run, mut rx) = crate::run_context::RunContext::channel(request_id); + let (run, mut rx) = crate::run_context::RunContext::channel( + crate::run_context::named_run_id("req_ns_gating_miss"), + ); let agent = compose_gated_agent(&server, "github:*", "k8s", "list_repos", run).await; agent @@ -3132,7 +3135,9 @@ mod tests { prepared: &Arc, request_id: &str, ) -> (Agent, tokio::task::JoinHandle<()>) { - let agent = prepared.begin_run(request_id, None, None).unwrap(); + let agent = prepared + .begin_run(crate::run_context::named_run_id(request_id), None, None) + .unwrap(); let mut events = agent.events.lock().unwrap().take().unwrap(); let call = tokio::spawn({ let prepared = Arc::clone(prepared); @@ -3158,7 +3163,7 @@ mod tests { }; let (first, call) = park(&prepared, "req_a").await; - registry.cancel_request_local("req_a"); + registry.cancel_request_local(&crate::run_context::named_run_id("req_a").to_string()); assert!( released(call).await, "ending req_a releases the approval it parked" @@ -3166,13 +3171,13 @@ mod tests { drop(first); let (_second, call) = park(&prepared, "req_b").await; - registry.cancel_request_local("req_a"); + registry.cancel_request_local(&crate::run_context::named_run_id("req_a").to_string()); tokio::time::sleep(Duration::from_millis(100)).await; assert!( !call.is_finished(), "ending req_a leaves req_b's approval parked" ); - registry.cancel_request_local("req_b"); + registry.cancel_request_local(&crate::run_context::named_run_id("req_b").to_string()); assert!( released(call).await, "ending req_b releases the approval it parked" @@ -3187,22 +3192,22 @@ mod tests { async fn a_prepared_agent_serves_one_run_at_a_time() { let prepared = prepared(); let first = prepared - .begin_run("req_a", None, None) + .begin_run(crate::run_context::named_run_id("req_a"), None, None) .expect("a fresh agent has no run"); let refused = prepared - .begin_run("req_b", None, None) + .begin_run(crate::run_context::named_run_id("req_b"), None, None) .expect_err("the slot is taken"); assert!( - matches!(refused, BeginRunError::RunInProgress(RunInProgress { ref active }) if active == "req_a"), + matches!(refused, BeginRunError::RunInProgress(RunInProgress { ref active }) if *active == crate::run_context::named_run_id("req_a")), "got: {refused:?}", ); drop(first); let second = prepared - .begin_run("req_b", None, None) + .begin_run(crate::run_context::named_run_id("req_b"), None, None) .expect("the slot is free again"); - assert_eq!(second.request_id(), "req_b"); + assert_eq!(second.run_id(), crate::run_context::named_run_id("req_b")); } /// Each run starts from the prepared seeds, and the slot the tools @@ -3212,14 +3217,12 @@ mod tests { let prepared = prepared(); let seed = prepared.scratchpad_budget.as_ref().unwrap(); - let first = prepared.begin_run("req_a", None, None).unwrap(); + let first = prepared + .begin_run(crate::run_context::named_run_id("req_a"), None, None) + .unwrap(); assert_eq!( - prepared - .run - .get() - .map(|run| run.id().to_string()) - .as_deref(), - Some("req_a") + prepared.run.get().map(|run| run.id()), + Some(crate::run_context::named_run_id("req_a")) ); prepared .run @@ -3241,14 +3244,12 @@ mod tests { drop(first); - let second = prepared.begin_run("req_b", None, None).unwrap(); + let second = prepared + .begin_run(crate::run_context::named_run_id("req_b"), None, None) + .unwrap(); assert_eq!( - prepared - .run - .get() - .map(|run| run.id().to_string()) - .as_deref(), - Some("req_b") + prepared.run.get().map(|run| run.id()), + Some(crate::run_context::named_run_id("req_b")) ); assert_eq!( second.scratchpad_budget().unwrap().scratchpad_usage().0, @@ -3274,21 +3275,34 @@ mod tests { )); let refused = prepared - .begin_run("req_bob", Some(&token("bob")), None) + .begin_run( + crate::run_context::named_run_id("req_bob"), + Some(&token("bob")), + None, + ) .expect_err("bob's token is not alice's"); assert!( matches!(refused, BeginRunError::ForwardedHeaderDiffers { ref header } if header == "x-user-token"), "got: {refused:?}", ); assert!( - prepared.begin_run("req_none", None, None).is_err(), + prepared + .begin_run(crate::run_context::named_run_id("req_none"), None, None) + .is_err(), "a request carrying no token is not alice's either", ); let served = prepared - .begin_run("req_alice", Some(&token("alice")), None) + .begin_run( + crate::run_context::named_run_id("req_alice"), + Some(&token("alice")), + None, + ) .expect("the same credentials are served"); - assert_eq!(served.request_id(), "req_alice"); + assert_eq!( + served.run_id(), + crate::run_context::named_run_id("req_alice") + ); } /// A stream still driving tools after its `Agent` is dropped keeps @@ -3296,24 +3310,26 @@ mod tests { #[tokio::test] async fn a_live_stream_keeps_the_run_bound_after_the_agent_drops() { let prepared = prepared(); - let agent = prepared.begin_run("req_a", None, None).unwrap(); + let agent = prepared + .begin_run(crate::run_context::named_run_id("req_a"), None, None) + .unwrap(); let stream = agent.stream_prompt("hello").await; drop(agent); assert_eq!( - prepared - .run - .get() - .map(|run| run.id().to_string()) - .as_deref(), - Some("req_a"), + prepared.run.get().map(|run| run.id()), + Some(crate::run_context::named_run_id("req_a")), "the stream holds the run", ); - assert!(prepared.begin_run("req_b", None, None).is_err()); + assert!( + prepared + .begin_run(crate::run_context::named_run_id("req_b"), None, None) + .is_err() + ); drop(stream); prepared - .begin_run("req_b", None, None) + .begin_run(crate::run_context::named_run_id("req_b"), None, None) .expect("the slot frees with the stream"); } } diff --git a/crates/aura/src/config.rs b/crates/aura/src/config.rs index 3d9550074..e51d49400 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::{RunId, SessionId}; /// Runtime build context for constructing agents. /// @@ -133,8 +111,8 @@ pub struct AgentRuntimeConfig { /// `None` disables approval gating. pub hitl: Option, - /// Request id (`req_…`) of the request this build serves. - pub request_id: Option, + /// The run this build serves. + pub run_id: Option, /// The request headers this build forwards; see [`ForwardedHeaders`]. pub forwarded_headers: ForwardedHeaders, @@ -175,7 +153,7 @@ impl Clone for AgentRuntimeConfig { scratchpad_tools_config: self.scratchpad_tools_config.clone(), orchestration_submit_result: self.orchestration_submit_result.clone(), hitl: self.hitl.clone(), - request_id: self.request_id.clone(), + run_id: self.run_id, forwarded_headers: self.forwarded_headers.clone(), instance_id: self.instance_id.clone(), hitl_request_approval_tool: self.hitl_request_approval_tool.clone(), @@ -220,7 +198,7 @@ impl std::fmt::Debug for AgentRuntimeConfig { .map(|_| ""), ) .field("hitl", &self.hitl.as_ref().map(|_| "")) - .field("request_id", &self.request_id) + .field("run_id", &self.run_id) .field("forwarded_headers", &self.forwarded_headers) .field("instance_id", &self.instance_id) .field( diff --git a/crates/aura/src/hitl/gate.rs b/crates/aura/src/hitl/gate.rs index c85493417..2dc35cc92 100644 --- a/crates/aura/src/hitl/gate.rs +++ b/crates/aura/src/hitl/gate.rs @@ -977,8 +977,10 @@ mod tests { )); let cancel = tokio_util::sync::CancellationToken::new(); - let (run, mut events) = - crate::run_context::RunContext::channel_on("req_run_cancel", cancel.clone()); + let (run, mut events) = crate::run_context::RunContext::channel_on( + crate::run_context::named_run_id("req_run_cancel"), + cancel.clone(), + ); gate.bind_run(run); let gated = Arc::clone(&gate); @@ -1215,7 +1217,8 @@ mod tests { "test-agent".to_string(), "test-instance-id".to_string(), ); - let (run, events) = crate::run_context::RunContext::channel(request_id); + let (run, events) = + crate::run_context::RunContext::channel(request_id.parse().expect("a run id")); gate.bind_run(run); ( WrappedTool::new(inner, Arc::new(gate) as Arc), @@ -1241,7 +1244,7 @@ mod tests { } fn unique_request_id() -> String { - format!("req_span_{}", uuid::Uuid::new_v4().simple()) + aura_events::RunId::mint().to_string() } /// The correlation the whole feature exists for: the id the approver diff --git a/crates/aura/src/hitl/route.rs b/crates/aura/src/hitl/route.rs index 56c7ecc7f..e30815c90 100644 --- a/crates/aura/src/hitl/route.rs +++ b/crates/aura/src/hitl/route.rs @@ -1221,8 +1221,9 @@ mod tests { #[tokio::test] async fn conversational_resolve_at_requested_event_succeeds() { - let request_id = format!("req_test_{}", uuid::Uuid::new_v4().simple()); - let (run, mut rx) = crate::run_context::RunContext::channel(request_id.as_str()); + let run_id = aura_events::RunId::mint(); + let request_id = run_id.to_string(); + let (run, mut rx) = crate::run_context::RunContext::channel(run_id); let (registry, route) = conv_route(Duration::from_secs(60)); let request = single_request( @@ -2219,8 +2220,9 @@ mod tests { #[tokio::test] async fn webhook_route_emits_requested_and_completed_on_channel_error() { - let request_id = format!("req_test_{}", uuid::Uuid::new_v4().simple()); - let (run, mut rx) = crate::run_context::RunContext::channel(request_id.as_str()); + let run_id = aura_events::RunId::mint(); + let request_id = run_id.to_string(); + let (run, mut rx) = crate::run_context::RunContext::channel(run_id); let route = super::DecisionRoute::Webhook { client: super::WebhookClient::new( super::build_webhook_client(), diff --git a/crates/aura/src/hitl/tool.rs b/crates/aura/src/hitl/tool.rs index 628e8346b..85f8d8e24 100644 --- a/crates/aura/src/hitl/tool.rs +++ b/crates/aura/src/hitl/tool.rs @@ -334,7 +334,7 @@ mod tests { route: &Arc, args: RequestApprovalArgs, ) -> ApprovalItem { - let request_id = format!("req_w2_{}", uuid::Uuid::new_v4().simple()); + let run_id = aura_events::RunId::mint(); let tool = RequestApprovalTool::new( route.clone(), AgentScope::Single { session_id: None }, @@ -343,7 +343,7 @@ mod tests { ); // The scope goes inside the spawn, because task-locals do not cross one. - let (run, mut rx) = crate::run_context::RunContext::channel(request_id.as_str()); + let (run, mut rx) = crate::run_context::RunContext::channel(run_id); let call_handle: tokio::task::JoinHandle> = tokio::spawn(crate::run_context::with_run(run, async move { tool.call(args).await @@ -499,8 +499,8 @@ mod tests { timeout: std::time::Duration::from_secs(60), }); - let request_id = format!("req_tool_span_{}", uuid::Uuid::new_v4().simple()); - let (run, mut events) = crate::run_context::RunContext::channel(request_id.as_str()); + let run_id = aura_events::RunId::mint(); + let (run, mut events) = crate::run_context::RunContext::channel(run_id); let tool = RequestApprovalTool::new( route, AgentScope::Single { session_id: None }, 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/mcp/client.rs b/crates/aura/src/mcp/client.rs index 9b1c6f6d7..edad45165 100644 --- a/crates/aura/src/mcp/client.rs +++ b/crates/aura/src/mcp/client.rs @@ -29,9 +29,9 @@ use crate::approver_headers::ApproverHeaders; use crate::mcp::progress::ProgressEnabledHandler; use crate::mcp::response::extract_tool_result; use crate::mcp::types::ToolNamespace; -use aura_events::AgentContext; use aura_events::ToolName; use aura_events::agent::{AgentEvent, AgentEventPayload}; +use aura_events::{AgentContext, RunId}; /// Custom HTTP client that captures the underlying HTTP status when a request /// fails. @@ -574,7 +574,7 @@ impl McpClient { /// per request and owning its manager. Sharing a manager across runs — warm /// MCP reuse, #578 — makes this answer last-writer-wins, and wants per-run /// tool instances or a rig-side change rather than another field here. - pub async fn run_id(&self) -> Option> { + pub async fn run_id(&self) -> Option { match crate::run_context::current_run_id() { Some(id) => Some(id), None => self @@ -582,7 +582,7 @@ impl McpClient { .read() .await .as_ref() - .map(|call| Arc::clone(call.run.id())), + .map(|call| call.run.id()), } } @@ -597,7 +597,7 @@ impl McpClient { .read() .await .as_ref() - .filter(|call| call.run.id().as_ref() == request_id) + .filter(|call| call.run.has_id(request_id)) .cloned() } @@ -659,7 +659,12 @@ impl McpClient { tool_name, http_request_id ); return self - .call_tool_tracked(tool_name, arguments, &http_request_id, approver_overrides) + .call_tool_tracked( + tool_name, + arguments, + &http_request_id.to_string(), + approver_overrides, + ) .await; } @@ -732,7 +737,7 @@ impl McpClient { // The transport task cannot read the run's task-local, so tie the token // to the run here, while still inside it. if let Some(run_id) = self.run_id().await { - self.own_progress_token(progress_token.clone(), &run_id) + self.own_progress_token(progress_token.clone(), &run_id.to_string()) .await; } // Held from here so a run cancelled mid-await, which drops this future @@ -817,7 +822,7 @@ impl McpClient { .context("Failed to send tool call request")?; if let Some(run_id) = self.run_id().await { - self.own_progress_token(handle.progress_token.clone(), &run_id) + self.own_progress_token(handle.progress_token.clone(), &run_id.to_string()) .await; } let _token_guard = ProgressTokenGuard { @@ -1062,7 +1067,7 @@ impl McpClient { // Drop this call's token ownership so straggler notifications stop // routing, then clear the binding. - owners_of(&self.token_owners).retain(|_, call| call.run.id().as_ref() != http_request_id); + owners_of(&self.token_owners).retain(|_, call| !call.run.has_id(http_request_id)); self.clear_current_call().await; // Forcefully close connection - server is ignoring cancellation anyway @@ -1087,9 +1092,9 @@ pub(crate) mod tests { use super::*; use crate::approver_headers::tests::captured_overrides; - fn owner(request_id: &str) -> CallContext { + fn owner(name: &str) -> CallContext { CallContext { - run: crate::run_context::RunContext::detached(request_id), + run: crate::run_context::RunContext::detached(crate::run_context::named_run_id(name)), agent: AgentContext::single_agent(), } } @@ -1627,30 +1632,32 @@ pub(crate) mod tests { async fn a_call_answers_only_for_the_request_that_named_it() { let (_server, client) = client_and_server(&requester_headers()).await; let worker = AgentContext::worker("log_worker", None, "coordinator"); + let req_1 = crate::run_context::named_run_id("req-1").to_string(); + let req_2 = crate::run_context::named_run_id("req-2").to_string(); assert!( - client.call_for("req-1").await.is_none(), + client.call_for(&req_1).await.is_none(), "an unbound client has no call to attribute work to" ); client .bind_call( - crate::run_context::RunContext::detached("req-1"), + crate::run_context::RunContext::detached(crate::run_context::named_run_id("req-1")), worker.clone(), ) .await; assert_eq!( - client.call_for("req-1").await.map(|call| call.agent), + client.call_for(&req_1).await.map(|call| call.agent), Some(worker) ); assert!( - client.call_for("req-2").await.is_none(), + client.call_for(&req_2).await.is_none(), "another request's id must not pick up this call" ); client.clear_current_call().await; assert!( - client.call_for("req-1").await.is_none(), + client.call_for(&req_1).await.is_none(), "clearing the call drops it with the request id" ); } @@ -1667,11 +1674,16 @@ pub(crate) mod tests { client .bind_call( - crate::run_context::RunContext::detached("req_bound"), + crate::run_context::RunContext::detached(crate::run_context::named_run_id( + "req_bound", + )), AgentContext::single_agent(), ) .await; - assert_eq!(client.run_id().await.as_deref(), Some("req_bound")); + assert_eq!( + client.run_id().await, + Some(crate::run_context::named_run_id("req_bound")) + ); } /// A scope still wins, so a call made inside one is attributed to that run @@ -1681,18 +1693,22 @@ pub(crate) mod tests { let (_server, client) = client_and_server(&requester_headers()).await; client .bind_call( - crate::run_context::RunContext::detached("req_bound"), + crate::run_context::RunContext::detached(crate::run_context::named_run_id( + "req_bound", + )), AgentContext::single_agent(), ) .await; let seen = crate::run_context::with_run( - crate::run_context::RunContext::detached("req_scoped"), + crate::run_context::RunContext::detached(crate::run_context::named_run_id( + "req_scoped", + )), async { client.run_id().await }, ) .await; - assert_eq!(seen.as_deref(), Some("req_scoped")); + assert_eq!(seen, Some(crate::run_context::named_run_id("req_scoped"))); } /// Being inside a run selects the tracked branch, so this is the same entry @@ -1712,7 +1728,9 @@ pub(crate) mod tests { .expect("the untracked call succeeds"); crate::run_context::with_run( - crate::run_context::RunContext::detached("http-req-1"), + crate::run_context::RunContext::detached(crate::run_context::named_run_id( + "http-req-1", + )), async { client .call_tool( diff --git a/crates/aura/src/mcp/progress.rs b/crates/aura/src/mcp/progress.rs index 888f937e5..6bc150e63 100644 --- a/crates/aura/src/mcp/progress.rs +++ b/crates/aura/src/mcp/progress.rs @@ -172,7 +172,7 @@ impl ClientHandler for ProgressEnabledHandler { let call = self.owner_of(¶ms.progress_token).await; if let Some(CallContext { run, agent }) = call { - let req_id = run.id().as_ref(); + let req_id = run.id(); let routed = run .emit(AgentEvent::new( @@ -231,9 +231,9 @@ mod tests { ProgressToken(NumberOrString::Number(n)) } - fn call(request_id: &str) -> CallContext { + fn call(name: &str) -> CallContext { CallContext { - run: crate::run_context::RunContext::detached(request_id), + run: crate::run_context::RunContext::detached(crate::run_context::named_run_id(name)), agent: aura_events::AgentContext::single_agent(), } } @@ -313,22 +313,16 @@ mod tests { // Nothing owns the token yet: the send has returned but the claim has not // landed. assert_eq!( - handler - .owner_of(&token(1)) - .await - .map(|call| call.run.id().to_string()), - Some("run_a".to_string()), + handler.owner_of(&token(1)).await.map(|call| call.run.id()), + Some(crate::run_context::named_run_id("run_a")), "an unclaimed token belongs to the call this client serves" ); // Once claimed, the token answers for itself. owners.lock().unwrap().insert(token(1), call("run_a_tool")); assert_eq!( - handler - .owner_of(&token(1)) - .await - .map(|call| call.run.id().to_string()), - Some("run_a_tool".to_string()), + handler.owner_of(&token(1)).await.map(|call| call.run.id()), + Some(crate::run_context::named_run_id("run_a_tool")), "a claim is more precise than the binding" ); } @@ -340,18 +334,12 @@ mod tests { let handler = handler_owning(&[(1, "run_a"), (2, "run_b")]); assert_eq!( - handler - .owner_of(&token(1)) - .await - .map(|c| c.run.id().to_string()), - Some("run_a".to_string()) + handler.owner_of(&token(1)).await.map(|c| c.run.id()), + Some(crate::run_context::named_run_id("run_a")) ); assert_eq!( - handler - .owner_of(&token(2)) - .await - .map(|c| c.run.id().to_string()), - Some("run_b".to_string()) + handler.owner_of(&token(2)).await.map(|c| c.run.id()), + Some(crate::run_context::named_run_id("run_b")) ); } @@ -366,20 +354,14 @@ mod tests { owners.lock().unwrap().insert(token(7), call("run_a_tool")); assert_eq!( - handler - .owner_of(&token(7)) - .await - .map(|c| c.run.id().to_string()), - Some("run_a_tool".to_string()) + handler.owner_of(&token(7)).await.map(|c| c.run.id()), + Some(crate::run_context::named_run_id("run_a_tool")) ); owners.lock().unwrap().remove(&token(7)); assert_eq!( - handler - .owner_of(&token(7)) - .await - .map(|c| c.run.id().to_string()), - Some("run_a".to_string()), + handler.owner_of(&token(7)).await.map(|c| c.run.id()), + Some(crate::run_context::named_run_id("run_a")), "the binding outlives the tokens of the calls it serves" ); diff --git a/crates/aura/src/orchestration/factory.rs b/crates/aura/src/orchestration/factory.rs index 9043b3735..17bb50c25 100644 --- a/crates/aura/src/orchestration/factory.rs +++ b/crates/aura/src/orchestration/factory.rs @@ -11,7 +11,7 @@ use async_trait::async_trait; use futures::stream::{self, BoxStream}; use tokio_util::sync::CancellationToken; -use crate::config::AgentRuntimeConfig; +use crate::config::{AgentRuntimeConfig, RunId}; use crate::provider_agent::{StreamError, StreamItem}; use crate::streaming::StreamingAgent; @@ -100,7 +100,7 @@ impl OrchestratorFactory { let (event_tx, event_rx) = tokio::sync::mpsc::channel::>(100); - let run_id = std::sync::Arc::clone(run.id()); + let run_id = run.id().to_string(); let RunTokens { cancel, finished } = tokens; let cancel_token_clone = cancel.clone(); // Marks the run finished on every exit path, which is what lets the @@ -174,7 +174,7 @@ impl OrchestratorFactory { tracing::info!("Orchestration cancelled"); if let Some(ref mcp_manager) = orchestrator.mcp_manager { let cancelled = mcp_manager - .cancel_and_close_all(run_id.as_ref(), "Client disconnected or timeout") + .cancel_and_close_all(&run_id, "Client disconnected or timeout") .await; if cancelled > 0 { tracing::info!("Cancelled {} MCP request(s) during orchestration shutdown", cancelled); @@ -218,6 +218,16 @@ impl StreamingAgent for OrchestratorFactory { options: crate::streaming::RunOptions, request_id: &str, ) -> crate::streaming::AgentRun { + // The run is the one this factory was built for, as an `Agent`'s is + // the one it began; a factory built for none starts a fresh run. + let run_id = self.agent_config.run_id.unwrap_or_else(RunId::mint); + if run_id.to_string() != request_id { + tracing::debug!( + %run_id, + request_id, + "streaming under the run's id rather than the one passed", + ); + } let (timeout, cancel) = options.into_parts(); let cancel_token = cancel.unwrap_or_default(); @@ -229,7 +239,7 @@ impl StreamingAgent for OrchestratorFactory { timeout, cancel_token.clone(), finished.clone(), - request_id.to_string(), + run_id.to_string(), ); finished }); @@ -239,7 +249,7 @@ impl StreamingAgent for OrchestratorFactory { // all orchestration LLM turns. let usage_state = crate::UsageState::new(); let (run, run_events) = - crate::run_context::RunContext::channel_on(request_id, cancel_token.clone()); + crate::run_context::RunContext::channel_on(run_id, cancel_token.clone()); let stream = self.spawn_orchestration_stream( query.to_string(), chat_history, diff --git a/crates/aura/src/orchestration/orchestrator.rs b/crates/aura/src/orchestration/orchestrator.rs index 2de6bf707..34bce1917 100644 --- a/crates/aura/src/orchestration/orchestrator.rs +++ b/crates/aura/src/orchestration/orchestrator.rs @@ -635,6 +635,9 @@ impl Orchestrator { let orchestrator_id = uuid::Uuid::new_v4().to_string(); + // Persistence names the run by an id of its own, which the park + // owner, the checkpoint and `RunParked` carry; the run's `RunContext` + // is named by its `RunId`. #780 makes them one id. let run_id_str = persistence.lock().await.run_id().to_string(); // One guard per park-mode run; `ParkGuard` documents arming and drop. let park_guard = agent_config @@ -1156,11 +1159,16 @@ impl Orchestrator { /// The run the orchestration serves, for its workers and coordinator to /// begin theirs within. That is the run in scope; a test driving the - /// orchestrator outside one gets a run nobody observes, under the - /// request id the config carries. + /// orchestrator outside one gets a run nobody observes, under the run id + /// the config carries. fn orchestration_run(&self) -> Arc { crate::run_context::current_run().unwrap_or_else(|| { - RunContext::channel(self.agent_config.request_id.clone().unwrap_or_default()).0 + RunContext::channel( + self.agent_config + .run_id + .unwrap_or_else(crate::config::RunId::mint), + ) + .0 }) } @@ -7703,8 +7711,8 @@ mod tests { }; use crate::session_store::{InMemoryApprovalStore, InMemoryEventBus}; - let request_id = format!("req_cancel_{}", uuid::Uuid::new_v4().simple()); - let (run, mut events) = crate::run_context::RunContext::channel(request_id.as_str()); + let run_id = aura_events::RunId::mint(); + let (run, mut events) = crate::run_context::RunContext::channel(run_id); let store: Arc = Arc::new(InMemoryApprovalStore::new()); @@ -7719,7 +7727,7 @@ mod tests { }), park_enabled: true, }), - request_id: Some(request_id.clone()), + run_id: Some(run_id), ..AgentRuntimeConfig::default() }; let orchestrator = Orchestrator::new(config).await.unwrap(); @@ -7835,7 +7843,7 @@ mod tests { }), memory_dir: Some(memory_dir.to_string_lossy().into_owned()), session_id: Some("park-sess".to_string()), - request_id: Some(format!("req_park_{}", uuid::Uuid::new_v4().simple())), + run_id: Some(aura_events::RunId::mint()), ..AgentRuntimeConfig::default() }; let orchestrator = Orchestrator::new(config).await.unwrap(); @@ -8083,12 +8091,11 @@ mod tests { let dir = tempfile::tempdir().unwrap(); let (orchestrator, store, registry, run_id) = park_orchestrator(dir.path()).await; let (plan, records, pending) = awaiting_plan_with_parked_calls(®istry, &run_id).await; - let request_id = orchestrator + let this_run = orchestrator .agent_config - .request_id - .clone() - .unwrap_or_default(); - let (run, mut events) = crate::run_context::RunContext::channel(request_id.as_str()); + .run_id + .expect("a park orchestrator is built for a run"); + let (run, mut events) = crate::run_context::RunContext::channel(this_run); // Arming inside the scope is what gives the guard the run its drop // sweep reports to. crate::run_context::with_run(Arc::clone(&run), arm_guard(&orchestrator, &plan)).await; @@ -8178,12 +8185,11 @@ mod tests { let dir = tempfile::tempdir().unwrap(); let (orchestrator, store, registry, run_id) = park_orchestrator(dir.path()).await; let (plan, records, pending) = awaiting_plan_with_parked_calls(®istry, &run_id).await; - let request_id = orchestrator + let this_run = orchestrator .agent_config - .request_id - .clone() - .unwrap_or_default(); - let (run, mut events) = crate::run_context::RunContext::channel(request_id.as_str()); + .run_id + .expect("a park orchestrator is built for a run"); + let (run, mut events) = crate::run_context::RunContext::channel(this_run); // The human decides the first call before the commit is attempted. let decided = pending[0].decision_id; @@ -8270,12 +8276,11 @@ mod tests { let (orchestrator, registry, run_id) = park_orchestrator_over(store.clone(), dir.path()).await; let (plan, records, pending) = awaiting_plan_with_parked_calls(®istry, &run_id).await; - let request_id = orchestrator + let this_run = orchestrator .agent_config - .request_id - .clone() - .unwrap_or_default(); - let (run, mut events) = crate::run_context::RunContext::channel(request_id.as_str()); + .run_id + .expect("a park orchestrator is built for a run"); + let (run, mut events) = crate::run_context::RunContext::channel(this_run); crate::run_context::with_run(Arc::clone(&run), arm_guard(&orchestrator, &plan)).await; let (event_tx, mut event_rx) = tokio::sync::mpsc::channel(32); @@ -8398,7 +8403,8 @@ mod tests { skills: None, }, )]); - let request_id = format!("req_orphan_{}", uuid::Uuid::new_v4().simple()); + let run_id = aura_events::RunId::mint(); + let request_id = run_id.to_string(); let config = AgentRuntimeConfig { hitl: Some(crate::hitl::HitlRuntime { patterns: Arc::from(["echo_tool".into()]), @@ -8410,7 +8416,7 @@ mod tests { }), memory_dir: Some(memory_dir.to_string_lossy().into_owned()), session_id: Some("orphan-sess".to_string()), - request_id: Some(request_id.clone()), + run_id: Some(run_id), orchestration: Some(OrchestrationConfig { enabled: true, workers, @@ -8522,7 +8528,8 @@ mod tests { let dir = tempfile::tempdir().unwrap(); let (orchestrator, store, _registry, request_id) = override_park_orchestrator(dir.path(), 1).await; - let (run, mut events) = crate::run_context::RunContext::channel(request_id.as_str()); + let (run, mut events) = + crate::run_context::RunContext::channel(request_id.parse().expect("a run id")); // Depth 1 gives the loop three turns (the rig's +1 safety net), so // the gated call must land on the third: the first two turns burn @@ -8581,7 +8588,8 @@ mod tests { let dir = tempfile::tempdir().unwrap(); let (orchestrator, store, _registry, request_id) = override_park_orchestrator(dir.path(), 4).await; - let (run, mut events) = crate::run_context::RunContext::channel(request_id.as_str()); + let (run, mut events) = + crate::run_context::RunContext::channel(request_id.parse().expect("a run id")); let (_model, gated_invocations) = gated_worker_override(vec![ScriptedTurn::tool_calls_then_stream_failure(vec![ @@ -8701,7 +8709,7 @@ mod tests { }), memory_dir: Some(memory_dir.to_string_lossy().into_owned()), session_id: Some(session_id.to_string()), - request_id: Some(format!("req_resume_{}", uuid::Uuid::new_v4().simple())), + run_id: Some(aura_events::RunId::mint()), orchestration: Some(OrchestrationConfig { enabled: true, workers, diff --git a/crates/aura/src/orchestration/park/commit.rs b/crates/aura/src/orchestration/park/commit.rs index c385962a6..fbb813efe 100644 --- a/crates/aura/src/orchestration/park/commit.rs +++ b/crates/aura/src/orchestration/park/commit.rs @@ -558,8 +558,8 @@ mod tests { let (registry, store) = conv_registry(); let run_id = "0191e8c0-ffff-7000-8000-000000000006"; let owner = run_owner_id(run_id); - let request_id = format!("req_sweep_{}", uuid::Uuid::new_v4().simple()); - let (run, mut events) = crate::run_context::RunContext::channel(request_id.as_str()); + let this_run = aura_events::RunId::mint(); + let (run, mut events) = crate::run_context::RunContext::channel(this_run); let now = chrono::Utc::now(); let decided = DecisionId::generate(); @@ -626,8 +626,8 @@ mod tests { let (registry, store) = conv_registry(); let run_id = "0191e8c0-aaaa-7000-8000-000000000007"; let owner = run_owner_id(run_id); - let request_id = format!("req_sweep_{}", uuid::Uuid::new_v4().simple()); - let (run, mut events) = crate::run_context::RunContext::channel(request_id.as_str()); + let this_run = aura_events::RunId::mint(); + let (run, mut events) = crate::run_context::RunContext::channel(this_run); let now = chrono::Utc::now(); let first = DecisionId::generate(); diff --git a/crates/aura/src/orchestration/park/guard.rs b/crates/aura/src/orchestration/park/guard.rs index 2297b3c19..ec2f2e05d 100644 --- a/crates/aura/src/orchestration/park/guard.rs +++ b/crates/aura/src/orchestration/park/guard.rs @@ -144,8 +144,8 @@ mod tests { async fn unpublished_guard_drop_cancels_run_approvals() { let (registry, store) = registry_with_store(); let run_id: RunId = "0191e8c0-2222-7000-8000-000000000042".parse().unwrap(); - let request_id = format!("req_guard_{}", uuid::Uuid::new_v4().simple()); - let (run, mut events) = crate::run_context::RunContext::channel(request_id.as_str()); + let this_run = aura_events::RunId::mint(); + let (run, mut events) = crate::run_context::RunContext::channel(this_run); let scope = worker_scope(run_id); let decision_id = DecisionId::generate(); diff --git a/crates/aura/src/orchestration/persistence_wrapper.rs b/crates/aura/src/orchestration/persistence_wrapper.rs index 470250690..2033f1c99 100644 --- a/crates/aura/src/orchestration/persistence_wrapper.rs +++ b/crates/aura/src/orchestration/persistence_wrapper.rs @@ -1449,7 +1449,7 @@ mod tests { session_id: None, }; - let request_id = format!("req_w2_{}", uuid::Uuid::new_v4().simple()); + let this_run = aura_events::RunId::mint(); let gate = Arc::new(HitlApprovalWrapper::new( Arc::from(["kubectl_*".into()]), @@ -1460,7 +1460,7 @@ mod tests { )); // `WrappedTool` runs `pre_call` in its own task, which no scope // crosses, so the gate is bound the way `stream` binds it. - let (run, mut rx) = crate::run_context::RunContext::channel(request_id.as_str()); + let (run, mut rx) = crate::run_context::RunContext::channel(this_run); gate.bind_run(run); let gate: Arc = gate; let persistence: Arc = Arc::new(test_wrapper(Arc::new(Mutex::new( diff --git a/crates/aura/src/orchestration/test_rig.rs b/crates/aura/src/orchestration/test_rig.rs index c95b2128c..1ab821de6 100644 --- a/crates/aura/src/orchestration/test_rig.rs +++ b/crates/aura/src/orchestration/test_rig.rs @@ -742,7 +742,7 @@ pub(crate) async fn park_orchestrator_in( }), memory_dir: Some(memory_dir.to_string_lossy().into_owned()), session_id: Some("park-sess".to_string()), - request_id: Some(format!("req_rig_{}", uuid::Uuid::new_v4().simple())), + run_id: Some(aura_events::RunId::mint()), orchestration: Some(super::OrchestrationConfig { enabled: true, workers, diff --git a/crates/aura/src/orchestration/types.rs b/crates/aura/src/orchestration/types.rs index d10d4a427..2eab5b543 100644 --- a/crates/aura/src/orchestration/types.rs +++ b/crates/aura/src/orchestration/types.rs @@ -4,12 +4,9 @@ //! queries into tasks, track their execution, and manage dependencies. use std::collections::HashMap; -use std::fmt; -use std::str::FromStr; use std::sync::Mutex; use serde::{Deserialize, Serialize}; -use uuid::Uuid; use crate::hitl::DecisionId; @@ -23,28 +20,7 @@ const MAX_STEP_NESTING: usize = 2; // Domain identifiers for orchestration runs and tasks, modeled as simple types // (opaque newtypes reached through canonical conversion traits). -/// Identifier for a single orchestration run. -/// -/// A run is an orchestration concept; single-agent requests have none. Run ids -/// are v4 UUIDs; parse one from its string form via `FromStr`. Serializes as -/// the bare UUID string. -#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)] -#[serde(transparent)] -pub struct RunId(Uuid); - -impl fmt::Display for RunId { - fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - self.0.fmt(f) - } -} - -impl FromStr for RunId { - type Err = uuid::Error; - - fn from_str(s: &str) -> Result { - Uuid::from_str(s).map(Self) - } -} +pub use aura_events::RunId; /// Identity of a worker task within a run. /// diff --git a/crates/aura/src/rig_builder.rs b/crates/aura/src/rig_builder.rs index 223ef6ee3..0a032ab61 100644 --- a/crates/aura/src/rig_builder.rs +++ b/crates/aura/src/rig_builder.rs @@ -12,7 +12,7 @@ use crate::builder::{ Agent, ClientTool, PreparedAgent, RunToolFactory, build_streaming_agent_with_tools, no_run_tools, }; -use crate::config::{AgentRuntimeConfig, WorkerSkills}; +use crate::config::{AgentRuntimeConfig, RunId, WorkerSkills}; use crate::error::BuilderError; use crate::forwarded_headers::ForwardedHeaders; use crate::hitl::PendingApprovals; @@ -177,8 +177,8 @@ impl RigBuilder { .map_err(|e| BuilderError::AgentError(format!("Failed to build agent: {e}"))) } - /// Prepare an agent and begin its run for `request_id`, recording skill - /// invocations with the builder's skill recorder. + /// Prepare an agent and begin its run `run_id` (a fresh one when `None`), + /// recording skill invocations with the builder's skill recorder. /// /// See [`Self::prepare_agent`] for the parameters. Each call prepares a /// fresh agent; callers that want to reuse one across requests call @@ -188,7 +188,7 @@ impl RigBuilder { req_headers: Option<&HashMap>, additional_tools: Vec>, client_tools: Option>, - request_id: Option, + run_id: Option, session_id: Option, ) -> Result { let prepared = self @@ -196,7 +196,7 @@ impl RigBuilder { .await?; prepared .begin_run( - request_id.unwrap_or_default(), + run_id.unwrap_or_else(RunId::mint), req_headers, self.skill_recorder.clone(), ) @@ -218,13 +218,13 @@ impl RigBuilder { req_headers: Option<&HashMap>, session_id: Option, client_tools: Option>, - request_id: Option, + run_id: Option, ) -> Result, BuilderError> { self.build_streaming_agent_with_tools( req_headers, session_id, client_tools, - request_id, + run_id, no_run_tools(), ) .await @@ -239,13 +239,13 @@ impl RigBuilder { req_headers: Option<&HashMap>, session_id: Option, client_tools: Option>, - request_id: Option, + run_id: Option, run_tools: RunToolFactory, ) -> Result, BuilderError> { let mut agent_config = self.discovered_agent_config(req_headers)?; resolve_mcp_headers(&mut agent_config, req_headers); agent_config.session_id = session_id; - agent_config.request_id = request_id; + agent_config.run_id = run_id; agent_config.skill_recorder = self.skill_recorder.clone(); build_streaming_agent_with_tools(&agent_config, client_tools, run_tools) diff --git a/crates/aura/src/run_context.rs b/crates/aura/src/run_context.rs index 5a31ebaac..4677d2446 100644 --- a/crates/aura/src/run_context.rs +++ b/crates/aura/src/run_context.rs @@ -21,7 +21,7 @@ use futures::Stream; use aura_events::agent::AgentEvent; use tokio::sync::mpsc; -use aura_events::ToolCallId; +use aura_events::{RunId, ToolCallId}; use crate::scratchpad::ContextBudget; use crate::skill_tool::SkillInvocationRecorder; @@ -33,7 +33,7 @@ pub const EVENT_CHANNEL_CAPACITY: usize = 1024; /// One run — what its own work needs to correlate, where its events go, and /// the state an agent's tools keep for it. pub struct RunContext { - id: Arc, + id: RunId, tool_calls: Mutex>, events: mpsc::Sender, cancel: CancellationToken, @@ -50,7 +50,7 @@ const MAX_PENDING_TOOL_CALLS: usize = 256; impl RunContext { /// A run and the receiver its observer reads, on a token of its own. - pub fn channel(id: impl Into>) -> (Arc, mpsc::Receiver) { + pub fn channel(id: RunId) -> (Arc, mpsc::Receiver) { Self::channel_on(id, CancellationToken::new()) } @@ -58,7 +58,7 @@ impl RunContext { /// to stop on — a child of its own caller's, so one run ending leaves the /// others alone. pub fn channel_on( - id: impl Into>, + id: RunId, cancel: CancellationToken, ) -> (Arc, mpsc::Receiver) { Self::channel_for_agent(id, cancel, None, None, None) @@ -67,7 +67,7 @@ impl RunContext { /// A run on `cancel` carrying the state a prepared agent's tools keep for /// it, and the receiver its observer reads. pub fn channel_for_agent( - id: impl Into>, + id: RunId, cancel: CancellationToken, scratchpad_budget: Option, turn_nudge: Option>, @@ -75,7 +75,7 @@ impl RunContext { ) -> (Arc, mpsc::Receiver) { let (events, receiver) = mpsc::channel(EVENT_CHANNEL_CAPACITY); let run = Arc::new(Self { - id: id.into(), + id, tool_calls: Mutex::new(VecDeque::new()), events, cancel, @@ -97,7 +97,7 @@ impl RunContext { skill_recorder: Option>, ) -> Arc { Arc::new(Self { - id: Arc::clone(&parent.id), + id: parent.id, tool_calls: Mutex::new(VecDeque::new()), events: parent.events.clone(), cancel: parent.cancel.clone(), @@ -111,14 +111,14 @@ impl RunContext { /// reading what it emits. Production names a run it can reach an observer /// through, or names none. #[cfg(test)] - pub fn detached(id: impl Into>) -> Arc { + pub fn detached(id: RunId) -> Arc { Self::channel(id).0 } /// [`detached`](Self::detached), carrying tool state. #[cfg(test)] pub(crate) fn detached_with( - id: impl Into>, + id: RunId, scratchpad_budget: Option, turn_nudge: Option>, ) -> Arc { @@ -163,8 +163,13 @@ impl RunContext { delivered } - pub fn id(&self) -> &Arc { - &self.id + pub fn id(&self) -> RunId { + self.id + } + + /// Whether `id` names this run. A string that is not a run id names none. + pub fn has_id(&self, id: &str) -> bool { + id.parse::().is_ok_and(|parsed| parsed == self.id) } /// The token that cancels this run. @@ -273,13 +278,17 @@ impl BoundRun { /// A slot holding an unobserved run that carries only `budget`. #[cfg(test)] pub(crate) fn pinned_budget(budget: ContextBudget) -> Self { - Self::holding(RunContext::detached_with("pinned", Some(budget), None)) + Self::holding(RunContext::detached_with(RunId::mint(), Some(budget), None)) } /// A slot holding an unobserved run that carries only `turn_nudge`. #[cfg(test)] pub(crate) fn pinned_nudge(turn_nudge: Arc) -> Self { - Self::holding(RunContext::detached_with("pinned", None, Some(turn_nudge))) + Self::holding(RunContext::detached_with( + RunId::mint(), + None, + Some(turn_nudge), + )) } } @@ -308,10 +317,10 @@ impl RunLease { /// A prepared agent was asked to begin a run while it still serves another. #[derive(Debug, thiserror::Error)] -#[error("prepared agent already serves request `{active}`; it runs one request at a time")] +#[error("prepared agent already serves run `{active}`; it serves one run at a time")] pub struct RunInProgress { /// Id of the run holding the agent. - pub active: String, + pub active: RunId, } tokio::task_local! { @@ -341,8 +350,8 @@ pub fn current_run() -> Option> { RUN.try_with(Arc::clone).ok() } -pub fn current_run_id() -> Option> { - RUN.try_with(|run| Arc::clone(run.id())).ok() +pub fn current_run_id() -> Option { + RUN.try_with(|run| run.id()).ok() } /// Runs `f` with `run` in scope. Task-locals do not cross `tokio::spawn`, so @@ -377,8 +386,8 @@ impl Stream for ScopedStream { /// Runs `f` with a fresh run in scope and returns what it emitted, for a test /// that asserts on a run's events without standing up an observer. #[cfg(test)] -pub(crate) async fn observing(id: &str, f: F) -> (F::Output, Vec) { - let (run, mut events) = RunContext::channel(id); +pub(crate) async fn observing(name: &str, f: F) -> (F::Output, Vec) { + let (run, mut events) = RunContext::channel(named_run_id(name)); let out = with_run(run, f).await; let mut seen = Vec::new(); @@ -388,22 +397,29 @@ pub(crate) async fn observing(id: &str, f: F) -> (F::Output, Vec RunId { + uuid::Uuid::new_v5(&uuid::Uuid::NAMESPACE_OID, name.as_bytes()).into() +} + #[cfg(test)] mod tests { use super::*; use futures::StreamExt; - fn run(id: &str) -> Arc { - RunContext::detached(id) + fn run(name: &str) -> Arc { + RunContext::detached(named_run_id(name)) } #[tokio::test] async fn a_scope_established_inside_a_spawn_holds() { - let (run, _rx) = RunContext::channel("spawned"); + let (run, _rx) = RunContext::channel(named_run_id("spawned")); let seen = tokio::spawn(with_run(run, async { current_run_id() })) .await .unwrap(); - assert_eq!(seen.as_deref(), Some("spawned")); + assert_eq!(seen, Some(named_run_id("spawned"))); } /// A run built on the caller's token stops when the caller does. The work @@ -413,7 +429,7 @@ mod tests { #[tokio::test] async fn a_run_stops_on_the_token_it_was_built_on() { let caller = CancellationToken::new(); - let (run, _events) = RunContext::channel_on("run_on_token", caller.clone()); + let (run, _events) = RunContext::channel_on(named_run_id("run_on_token"), caller.clone()); assert!(!run.cancel_token().is_cancelled()); caller.cancel(); @@ -424,8 +440,8 @@ mod tests { /// cancels something else. #[tokio::test] async fn a_run_given_no_token_has_its_own() { - let (a, _ea) = RunContext::channel("run_a"); - let (b, _eb) = RunContext::channel("run_b"); + let (a, _ea) = RunContext::channel(named_run_id("run_a")); + let (b, _eb) = RunContext::channel(named_run_id("run_b")); a.cancel_token().cancel(); assert!(a.cancel_token().is_cancelled()); @@ -444,7 +460,7 @@ mod tests { #[tokio::test] async fn a_scope_supplies_the_run() { let seen = with_run(run("run_1"), async { current_run_id() }).await; - assert_eq!(seen.as_deref(), Some("run_1")); + assert_eq!(seen, Some(named_run_id("run_1"))); } #[tokio::test] @@ -458,8 +474,8 @@ mod tests { current_run_id() })); - assert_eq!(a.await.unwrap().as_deref(), Some("run_a")); - assert_eq!(b.await.unwrap().as_deref(), Some("run_b")); + assert_eq!(a.await.unwrap(), Some(named_run_id("run_a"))); + assert_eq!(b.await.unwrap(), Some(named_run_id("run_b"))); } #[tokio::test] @@ -467,10 +483,7 @@ mod tests { let inner = futures::stream::iter(0..3).map(|_| current_run_id()); let seen: Vec<_> = scope_stream(run("run_s"), inner).collect().await; - assert_eq!( - seen.iter().map(|id| id.as_deref()).collect::>(), - vec![Some("run_s"); 3] - ); + assert_eq!(seen, vec![Some(named_run_id("run_s")); 3]); } /// Orchestration drives workers with `FuturesUnordered` inside the run's @@ -496,10 +509,7 @@ mod tests { }) .await; - assert_eq!( - seen.iter().map(|id| id.as_deref()).collect::>(), - vec![Some("run_w"); 3] - ); + assert_eq!(seen, vec![Some(named_run_id("run_w")); 3]); } /// A spawned task does not inherit its parent's scope, which is why every @@ -534,12 +544,12 @@ mod tests { let nudge = TurnNudgeState::new(true, None, 2).unwrap(); let slot = BoundRun::default(); slot.bind(RunContext::detached_with( - "req_a", + named_run_id("req_a"), Some(budget.clone()), Some(Arc::clone(&nudge)), )); - assert_eq!(slot.id_or_empty(), "req_a"); + assert_eq!(slot.id_or_empty(), named_run_id("req_a").to_string()); slot.scratchpad_budget().unwrap().record_intercepted(7); assert_eq!( budget.scratchpad_usage().0, @@ -548,8 +558,8 @@ mod tests { ); assert!(Arc::ptr_eq(&slot.turn_nudge().unwrap(), &nudge)); - slot.bind(RunContext::detached("req_b")); - assert_eq!(slot.id_or_empty(), "req_b"); + slot.bind(RunContext::detached(named_run_id("req_b"))); + assert_eq!(slot.id_or_empty(), named_run_id("req_b").to_string()); assert!(slot.scratchpad_budget().is_none()); } @@ -557,11 +567,11 @@ mod tests { /// cancellation see it, with the worker's own tool state. #[tokio::test] async fn a_child_shares_its_parents_identity_and_keeps_its_own_state() { - let (parent, mut events) = RunContext::channel("req_parent"); + let (parent, mut events) = RunContext::channel(named_run_id("req_parent")); let nudge = TurnNudgeState::new(true, None, 2).unwrap(); let child = RunContext::child(&parent, None, Some(Arc::clone(&nudge)), None); - assert_eq!(child.id().as_ref(), "req_parent"); + assert_eq!(child.id(), parent.id()); assert!(Arc::ptr_eq(child.turn_nudge().unwrap(), &nudge)); assert!(parent.turn_nudge().is_none(), "the parent keeps none of it"); diff --git a/crates/aura/src/scratchpad/setup.rs b/crates/aura/src/scratchpad/setup.rs index 6c39d8725..d13161547 100644 --- a/crates/aura/src/scratchpad/setup.rs +++ b/crates/aura/src/scratchpad/setup.rs @@ -306,7 +306,7 @@ mod tests { assert!(build.tools_config.run.scratchpad_budget().is_none()); let run_budget = build.budget.fresh(); run.bind(crate::run_context::RunContext::detached_with( - "req", + crate::run_context::named_run_id("req"), Some(run_budget.clone()), None, )); diff --git a/crates/aura/src/skill_rehydration.rs b/crates/aura/src/skill_rehydration.rs index c6dfcda48..465096bbe 100644 --- a/crates/aura/src/skill_rehydration.rs +++ b/crates/aura/src/skill_rehydration.rs @@ -411,7 +411,7 @@ mod tests { // Turn N: history [user], anchor = 0 + 1. The LLM calls load_skill. let recorder = Arc::new(SkillInvocationRecorder::new(store.clone(), log.clone(), 1)); let (run, _events) = crate::run_context::RunContext::channel_for_agent( - "req-turn-n", + crate::run_context::named_run_id("req-turn-n"), tokio_util::sync::CancellationToken::new(), None, None, diff --git a/crates/aura/src/skill_tool.rs b/crates/aura/src/skill_tool.rs index b308be76a..4203c45d4 100644 --- a/crates/aura/src/skill_tool.rs +++ b/crates/aura/src/skill_tool.rs @@ -877,8 +877,14 @@ mod tests { log.clone(), anchor, )); - RunContext::channel_for_agent(id, CancellationToken::new(), None, None, Some(recorder)) - .0 + RunContext::channel_for_agent( + crate::run_context::named_run_id(id), + CancellationToken::new(), + None, + None, + Some(recorder), + ) + .0 }; let load = |name: &str| LoadSkillArgs { name: name.to_string(), diff --git a/crates/aura/src/streaming.rs b/crates/aura/src/streaming.rs index 794f7dfac..49a1bb7a7 100644 --- a/crates/aura/src/streaming.rs +++ b/crates/aura/src/streaming.rs @@ -283,7 +283,8 @@ pub trait StreamingAgent: Send + Sync { /// holds. The returned handle owns the events, the token that cancels them, /// and the usage they accumulate. /// - /// `request_id` correlates MCP progress and tool events for this run. + /// The run streams under the id its agent was built for, which a caller + /// passes back as `request_id`; a different value is logged, not adopted. async fn stream( &self, query: &str, @@ -481,7 +482,8 @@ mod tests { agent: aura_events::AgentContext, items: Vec>, ) -> (Vec, usize, Vec) { - let (run, mut events) = RunContext::channel("run_tee"); + let (run, mut events) = + RunContext::channel(crate::run_context::named_run_id("run_tee")); let passed = tee_content(run, agent, futures::stream::iter(items)) .collect::>() .await @@ -564,7 +566,8 @@ mod tests { /// the items must still pass. #[tokio::test] async fn an_unobserved_run_still_streams_its_items() { - let (run, events) = RunContext::channel("run_unobserved"); + let (run, events) = + RunContext::channel(crate::run_context::named_run_id("run_unobserved")); drop(events); let passed = tee_content( diff --git a/crates/aura/src/streaming_request_hook.rs b/crates/aura/src/streaming_request_hook.rs index ab872f267..11ac47932 100644 --- a/crates/aura/src/streaming_request_hook.rs +++ b/crates/aura/src/streaming_request_hook.rs @@ -106,7 +106,7 @@ impl Drop for ParkCellRegistration { /// correlates nothing until #732. Its tool events stay off the run for the same /// reason, and orchestration reports the worker's calls itself. fn queue_owner(stream_id: &str) -> Option> { - current_run().filter(|run| run.id().as_ref() == stream_id) + current_run().filter(|run| run.has_id(stream_id)) } /// Sends a tool event raised by this stream to the run [`queue_owner`] gives it. @@ -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(); } @@ -775,10 +776,11 @@ mod tests { /// scope, and its tool events must not reach the run as the run's own. #[tokio::test] async fn only_the_run_s_own_stream_sends_it_tool_events() { - let (run, mut events) = RunContext::channel("req_1"); + let id = crate::run_context::named_run_id("req_1").to_string(); + let (run, mut events) = RunContext::channel(id.parse().unwrap()); with_run(run, async { - emit_from_stream("req_1:task:0:attempt:1", requested("worker")).await; - emit_from_stream("req_1", requested("own")).await; + emit_from_stream(&format!("{id}:task:0:attempt:1"), requested("worker")).await; + emit_from_stream(&id, requested("own")).await; }) .await; @@ -796,17 +798,19 @@ mod tests { #[tokio::test] async fn the_run_s_own_stream_owns_the_queue() { - let run = RunContext::detached("req_1"); - let owned = with_run(run, async { queue_owner("req_1").is_some() }).await; + let id = crate::run_context::named_run_id("req_1"); + let run = RunContext::detached(id); + let owned = with_run(run, async { queue_owner(&id.to_string()).is_some() }).await; assert!(owned); } /// An orchestration worker streams under its task attempt. #[tokio::test] async fn a_worker_s_stream_owns_no_queue() { - let run = RunContext::detached("req_1"); + let id = crate::run_context::named_run_id("req_1"); + let run = RunContext::detached(id); let owned = with_run(run, async { - queue_owner("req_1:task:0:attempt:1").is_some() + queue_owner(&format!("{id}:task:0:attempt:1")).is_some() }) .await; assert!( diff --git a/crates/aura/src/turn_nudge.rs b/crates/aura/src/turn_nudge.rs index 6ab647f96..9eead3502 100644 --- a/crates/aura/src/turn_nudge.rs +++ b/crates/aura/src/turn_nudge.rs @@ -451,7 +451,7 @@ mod tests { let first_nudge = seed.fresh(); slot.bind(RunContext::detached_with( - "req_a", + crate::run_context::named_run_id("req_a"), None, Some(Arc::clone(&first_nudge)), )); @@ -464,7 +464,11 @@ mod tests { "the first run is on its penultimate turn", ); - slot.bind(RunContext::detached_with("req_b", None, Some(seed.fresh()))); + slot.bind(RunContext::detached_with( + crate::run_context::named_run_id("req_b"), + None, + Some(seed.fresh()), + )); assert_eq!( tool.call("hello".to_string()).await.unwrap(), "hello",