From 179eab1369bcd58d57df8ad3c9df40d61c718537 Mon Sep 17 00:00:00 2001 From: Justin Gross Date: Mon, 5 Oct 2026 15:48:18 -0400 Subject: [PATCH 1/6] refactor(agents): name a run by its RunId The envelope names a run by RunId, a UUID, while the task-local run context named it by the HTTP request id string, so the two could never be the same value. RunContext::id is now a RunId. The chat, A2A, and Slack handlers mint one before building the agent and use its string form as the request id, so every request-keyed registry (HITL approvals and their sweep, MCP cancellation, the A2A cancel map) keeps one value per run. begin_run, AgentRuntimeConfig, and RigBuilder's build_agent, build_streaming_agent_with_headers and build_streaming_agent_with_tools take a RunId, and an agent built without one mints its own. The orchestration factory streams under the run id it was built for, as Agent already did. RunContext::has_id compares a request id string against the run without allocating. Request ids become hyphenated UUIDs rather than req_, a2a_, and slack__; the A2A executor and the Slack runner log each run id beside the task or message it serves, at debug, so the two stay joinable. The orchestration RunId becomes a re-export of aura_events::RunId, so there is one run id type. Orchestration persistence still mints its own value for checkpoints and the park owner key; adopting the run's id there waits for the runtime, which owns resume. Fixes: GH-778 Ref: GH-578 Ref: GH-780 --- .../aura-web-server/src/a2a/agent_executor.rs | 8 +- crates/aura-web-server/src/handlers.rs | 19 +-- crates/aura-web-server/src/slack/runner.rs | 25 ++-- crates/aura/src/builder.rs | 118 ++++++++++-------- crates/aura/src/config.rs | 10 +- crates/aura/src/hitl/gate.rs | 11 +- crates/aura/src/hitl/route.rs | 10 +- crates/aura/src/hitl/tool.rs | 8 +- crates/aura/src/mcp/client.rs | 60 +++++---- crates/aura/src/mcp/progress.rs | 48 +++---- crates/aura/src/orchestration/factory.rs | 20 ++- crates/aura/src/orchestration/orchestrator.rs | 62 +++++---- crates/aura/src/orchestration/park/commit.rs | 8 +- crates/aura/src/orchestration/park/guard.rs | 4 +- .../src/orchestration/persistence_wrapper.rs | 4 +- crates/aura/src/orchestration/test_rig.rs | 2 +- crates/aura/src/orchestration/types.rs | 26 +--- crates/aura/src/rig_builder.rs | 18 +-- crates/aura/src/run_context.rs | 100 ++++++++------- crates/aura/src/scratchpad/setup.rs | 2 +- crates/aura/src/skill_rehydration.rs | 2 +- crates/aura/src/skill_tool.rs | 10 +- crates/aura/src/streaming.rs | 9 +- crates/aura/src/streaming_request_hook.rs | 19 +-- crates/aura/src/turn_nudge.rs | 8 +- 25 files changed, 333 insertions(+), 278 deletions(-) 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 1eae4aef2..018dc9eef 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 7e19c726d..e51d49400 100644 --- a/crates/aura/src/config.rs +++ b/crates/aura/src/config.rs @@ -39,7 +39,7 @@ pub enum WorkerSkills { Override(Vec), } -pub use aura_events::SessionId; +pub use aura_events::{RunId, SessionId}; /// Runtime build context for constructing agents. /// @@ -111,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, @@ -153,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(), @@ -198,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/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 21c6d204c..df606795b 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..2b9bd379b 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,33 @@ pub(crate) async fn observing(id: &str, f: F) -> (F::Output, Vec RunId { + RunId::try_from(uuid::Uuid::new_v5( + &uuid::Uuid::NAMESPACE_OID, + name.as_bytes(), + )) + .expect("a v5 UUID is never nil") +} + #[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 +433,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 +444,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 +464,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 +478,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 +487,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 +513,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 +548,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 +562,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 +571,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 44c9909a3..51b80b396 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. @@ -778,10 +778,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; @@ -799,17 +800,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", From e4abf8d836e8d0157f1f26ffbcbcfc1b8eeaa837 Mon Sep 17 00:00:00 2001 From: Justin Gross Date: Fri, 9 Oct 2026 00:29:07 -0400 Subject: [PATCH 2/6] fix(orchestration): name an unscoped orchestration by one run id `orchestration_run()` fell back to a run under the config's id, or a fresh one when the config carried none, and it minted that fresh id on every call. An orchestrator built from a config without a run id and driven outside a run's scope therefore began its coordinator and each of its workers within a different run, and the doc's "the run the orchestration serves" named several. `Orchestrator::new` now mints the id once, as `OrchestratorFactory::new` does, writes it back to the config and keeps it on the orchestrator; `orchestration_run()` reads that. The note beside the persistence run id states the two ids as they are rather than as a change to come. Ref: GH-778 --- crates/aura/src/orchestration/orchestrator.rs | 56 ++++++++++++++----- 1 file changed, 43 insertions(+), 13 deletions(-) diff --git a/crates/aura/src/orchestration/orchestrator.rs b/crates/aura/src/orchestration/orchestrator.rs index df606795b..61d46299e 100644 --- a/crates/aura/src/orchestration/orchestrator.rs +++ b/crates/aura/src/orchestration/orchestrator.rs @@ -450,6 +450,9 @@ pub struct Orchestrator { /// The underlying agent configuration (for creating workers) agent_config: AgentRuntimeConfig, + /// The run the orchestration serves. + run_id: crate::config::RunId, + /// Tool call observer for coordinator visibility into worker tool execution. /// Wired to emit AgentEventPayload for real-time SSE streaming via spawn_tool_event_forwarder. pub(super) tool_call_observer: ToolCallObserver, @@ -595,10 +598,16 @@ enum LoopStep { } impl Orchestrator { - /// Create a new orchestrator from configuration. + /// Create a new orchestrator from configuration, for the run + /// `agent_config.run_id` names, or for a fresh run when it names none. pub async fn new( - agent_config: AgentRuntimeConfig, + mut agent_config: AgentRuntimeConfig, ) -> Result> { + // Minted once, here, so every agent of the orchestration begins its + // run within the same one. + let run_id = *agent_config + .run_id + .get_or_insert_with(crate::config::RunId::mint); let orchestration_config = agent_config.orchestration.clone().unwrap_or_default(); // Initialize MCP manager (shared across coordinator and all workers via Arc) @@ -637,7 +646,7 @@ impl Orchestrator { // 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. + // is named by its `RunId`. The two are distinct. 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 @@ -665,6 +674,7 @@ impl Orchestrator { Ok(Self { orchestrator_id, + run_id, config: orchestration_config, agent_config, tool_call_observer, @@ -1159,17 +1169,10 @@ 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 run id - /// the config carries. + /// orchestrator outside one gets a run nobody observes, under the id the + /// orchestrator was created with. fn orchestration_run(&self) -> Arc { - crate::run_context::current_run().unwrap_or_else(|| { - RunContext::channel( - self.agent_config - .run_id - .unwrap_or_else(crate::config::RunId::mint), - ) - .0 - }) + crate::run_context::current_run().unwrap_or_else(|| RunContext::channel(self.run_id).0) } /// Park-mode wiring for one worker attempt: `Some` only when this run @@ -7507,6 +7510,33 @@ mod tests { assert_eq!(other["content"], "kept"); } + /// Every agent of an orchestration begins its run within the + /// orchestration's, so outside a scope they all get the same unobserved + /// run, under the id the orchestrator was created with. + #[tokio::test] + async fn an_unscoped_orchestration_names_one_run() { + let orchestrator = Orchestrator::new(crate::config::AgentRuntimeConfig::default()) + .await + .unwrap(); + let first = orchestrator.orchestration_run(); + let second = orchestrator.orchestration_run(); + assert_eq!(first.id(), second.id()); + assert_eq!( + orchestrator.agent_config.run_id, + Some(first.id()), + "the config carries the id the orchestrator minted" + ); + + let named = crate::run_context::named_run_id("req_named"); + let orchestrator = Orchestrator::new(crate::config::AgentRuntimeConfig { + run_id: Some(named), + ..Default::default() + }) + .await + .unwrap(); + assert_eq!(orchestrator.orchestration_run().id(), named); + } + /// Attached artifacts are listed in the worker's context after any /// dependency results, and alone when the task has no dependencies. #[tokio::test] From 7481e118ab1742cf17017b9d75a966bc33f77fba Mon Sep 17 00:00:00 2001 From: Justin Gross Date: Fri, 9 Oct 2026 00:29:07 -0400 Subject: [PATCH 3/6] fix(agents): match a run's id by its spelling `RunContext::has_id` parsed the string it was given and compared UUIDs, so the same id in upper case, in simple form or as a URN named the run; `StreamClaim::claim` demands the run's own `Display` form, because that string is the key every request-keyed registry holds. The two answers to "does this string name this run" disagreed, and the lenient one is what the MCP registries use to route cancellation and tool events. `has_id` now compares the spelling, encoded into a stack buffer rather than allocated, so a run answers to one string everywhere. Ref: GH-778 --- crates/aura/src/run_context.rs | 24 ++++++++++++++++++++++-- 1 file changed, 22 insertions(+), 2 deletions(-) diff --git a/crates/aura/src/run_context.rs b/crates/aura/src/run_context.rs index 2b9bd379b..822536419 100644 --- a/crates/aura/src/run_context.rs +++ b/crates/aura/src/run_context.rs @@ -167,9 +167,13 @@ impl RunContext { self.id } - /// Whether `id` names this run. A string that is not a run id names none. + /// Whether `id` is this run's id as the run spells it — the form its + /// `Display` gives, which is the key every request-keyed registry holds. + /// The same UUID spelled another way names no run, and neither does a + /// string that is not one. pub fn has_id(&self, id: &str) -> bool { - id.parse::().is_ok_and(|parsed| parsed == self.id) + let mut spelled = uuid::Uuid::encode_buffer(); + *self.id.as_uuid().hyphenated().encode_lower(&mut spelled) == *id } /// The token that cancels this run. @@ -461,6 +465,22 @@ mod tests { assert_eq!(current_run_id(), None); } + /// Registries key a run by the string its id displays as, so that is the + /// one spelling a run answers to: not the same UUID in another case or + /// form, and not a string that is no run id at all. + #[test] + fn a_run_answers_only_to_its_id_as_it_spells_it() { + let run = run("run_spelled"); + let spelled = run.id().to_string(); + + assert!(run.has_id(&spelled)); + assert!(!run.has_id(&spelled.to_uppercase())); + assert!(!run.has_id(&run.id().as_uuid().simple().to_string())); + assert!(!run.has_id(&format!("urn:uuid:{spelled}"))); + assert!(!run.has_id("req_1")); + assert!(!run.has_id(&named_run_id("run_other").to_string())); + } + #[tokio::test] async fn a_scope_supplies_the_run() { let seen = with_run(run("run_1"), async { current_run_id() }).await; From 94efb944355a1111eff2817ac13120f4fe4c770f Mon Sep 17 00:00:00 2001 From: Justin Gross Date: Fri, 9 Oct 2026 10:18:46 -0400 Subject: [PATCH 4/6] fix(server): tie a run's id to its origin on the agent span A run's id is a bare UUID, so a log line that names it no longer says which A2A task or Slack message the run serves, and the line that tied them together was debug-only. Each ingress now opens its run's agent.stream root span with the run's id and its origin: the A2A task and context ids, or the Slack channel and message ts. Chat's span carries the run's id, which its HTTP span already records. A2A had no agent.stream span; its execution stream is now polled inside one, as chat's and Slack's runs already are. The Slack run id is minted when the message is queued, where its span is opened. Ref: GH-778 --- Cargo.lock | 1 + crates/aura-web-server/Cargo.toml | 1 + .../aura-web-server/src/a2a/agent_executor.rs | 76 +++++++++++++++++-- crates/aura-web-server/src/handlers.rs | 6 +- crates/aura-web-server/src/slack/runner.rs | 34 +++++---- 5 files changed, 97 insertions(+), 21 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 19753da9f..025d106fe 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -414,6 +414,7 @@ dependencies = [ "tower-http", "tracing", "tracing-opentelemetry", + "tracing-subscriber", "uuid", "wiremock", ] diff --git a/crates/aura-web-server/Cargo.toml b/crates/aura-web-server/Cargo.toml index 40ecaf885..dff8a4013 100644 --- a/crates/aura-web-server/Cargo.toml +++ b/crates/aura-web-server/Cargo.toml @@ -104,6 +104,7 @@ http-body-util = "0.1" tempfile = "3" aura-config = { path = "../aura-config", features = ["test_util"] } wiremock = "0.6" +tracing-subscriber = { workspace = true } # Exhaustively explores interleavings of the task-cancel map, which a stress # test cannot reach: the window is one mutex release and reacquire wide. diff --git a/crates/aura-web-server/src/a2a/agent_executor.rs b/crates/aura-web-server/src/a2a/agent_executor.rs index a75dcc428..55272d6a6 100644 --- a/crates/aura-web-server/src/a2a/agent_executor.rs +++ b/crates/aura-web-server/src/a2a/agent_executor.rs @@ -243,7 +243,18 @@ impl AgentExecutor for AuraAgentExecutor { let run_tools = self.app_state.run_tools(&config); let mut append_tracker: HashMap<(String, String, String), bool> = HashMap::new(); - Box::pin(async_stream::stream! { + // The run is polled inside this span, so its log lines, the agent's + // included, carry the run's id and the task it executes. + let run_id = aura::RunId::mint(); + let span = tracing::info_span!( + parent: None, + "agent.stream", + run.id = %run_id, + a2a.task_id = %ctx.task_id, + a2a.context_id = %ctx.context_id, + ); + + let execution = async_stream::stream! { let task_id = ctx.task_id.clone(); let context_id = ctx.context_id.clone(); @@ -269,11 +280,9 @@ impl AgentExecutor for AuraAgentExecutor { metadata: None, })); - // The run's id, and in string form the request id everything + // In string form, the run's id is 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 @@ -562,7 +571,8 @@ impl AgentExecutor for AuraAgentExecutor { metadata: None, })); } - }) + }; + in_span(span, Box::pin(execution)) } fn cancel(&self, ctx: ExecutorContext) -> BoxStream<'static, Result> { @@ -625,6 +635,18 @@ fn extract_text(parts: Vec) -> Result { Ok(strings.join("\n")) } +/// `stream`, polled inside `span`, so its log lines carry the span's fields. +fn in_span( + span: tracing::Span, + mut stream: BoxStream<'static, T>, +) -> BoxStream<'static, T> { + futures_util::stream::poll_fn(move |cx| { + let _entered = span.enter(); + stream.poll_next_unpin(cx) + }) + .boxed() +} + pub(super) fn fail_status(task_id: &str, context_id: &str, error_msg: &str) -> StreamResponse { StreamResponse::StatusUpdate(TaskStatusUpdateEvent { task_id: task_id.to_string(), @@ -1050,6 +1072,50 @@ mod tests { assert!(claim_agent(&state, &task_id, &agent) == CancelRaced::Yes); } + /// A line the run logs carries its span's fields, which is what ties a + /// run's id to the task it executes. + #[tokio::test] + async fn a_line_logged_inside_the_span_carries_its_fields() { + let written = Arc::new(std::sync::Mutex::new(Vec::::new())); + let writer = { + let written = Arc::clone(&written); + move || SharedWriter(Arc::clone(&written)) + }; + let subscriber = tracing_subscriber::fmt() + .with_writer(writer) + .with_ansi(false) + .finish(); + let _default = tracing::subscriber::set_default(subscriber); + + let span = tracing::info_span!("agent.stream", a2a.task_id = %"t_1"); + let run = futures_util::stream::once(async { + tracing::warn!("inside the run"); + 1 + }) + .boxed(); + assert_eq!(in_span(span, run).collect::>().await, vec![1]); + + let written = String::from_utf8(written.lock().unwrap().clone()).unwrap(); + let line = written + .lines() + .find(|line| line.contains("inside the run")) + .expect("the line is written"); + assert!(line.contains("agent.stream{a2a.task_id=t_1}"), "{line}"); + } + + struct SharedWriter(Arc>>); + + impl std::io::Write for SharedWriter { + fn write(&mut self, buf: &[u8]) -> std::io::Result { + self.0.lock().unwrap().extend_from_slice(buf); + Ok(buf.len()) + } + + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } + } + /// The guard is created before the entry, because the agent build and the /// history fetch can both return early and would otherwise leave it behind. #[test] diff --git a/crates/aura-web-server/src/handlers.rs b/crates/aura-web-server/src/handlers.rs index 018dc9eef..80c3848df 100644 --- a/crates/aura-web-server/src/handlers.rs +++ b/crates/aura-web-server/src/handlers.rs @@ -769,9 +769,10 @@ async fn handle_non_streaming_completion( let (result_tx, result_rx) = oneshot::channel(); + let agent_span = tracing::info_span!(parent: None, "agent.stream", run.id = %config.request_id); let handle = tokio::spawn( execute_completion(setup, config, DeliveryMode::Collect { result_tx }) - .instrument(tracing::info_span!(parent: None, "agent.stream")), + .instrument(agent_span), ); data.active_requests.track_task(handle); @@ -814,6 +815,7 @@ async fn handle_streaming_completion( let heartbeat_interval = std::time::Duration::from_secs(15); + let agent_span = tracing::info_span!(parent: None, "agent.stream", run.id = %config.request_id); let handle = tokio::spawn( execute_completion( setup, @@ -823,7 +825,7 @@ async fn handle_streaming_completion( heartbeat_interval, }, ) - .instrument(tracing::info_span!(parent: None, "agent.stream")), + .instrument(agent_span), ); data.active_requests.track_task(handle); diff --git a/crates/aura-web-server/src/slack/runner.rs b/crates/aura-web-server/src/slack/runner.rs index bd0b3bc36..57aba8980 100644 --- a/crates/aura-web-server/src/slack/runner.rs +++ b/crates/aura-web-server/src/slack/runner.rs @@ -235,6 +235,16 @@ fn queue_answer(ingress: &Arc, inbound: Inbound, earlier: Option, inbound: Inbound, earlier: Option>) { - // 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(); + async fn answer( + &self, + inbound: Inbound, + prefetched: Option>, + run_id: aura::RunId, + ) { + // In string form, the run's id is the request id everything + // request-keyed reads. 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 => { From 1b232f489874d25ee4fccfefc45307d2088917f2 Mon Sep 17 00:00:00 2001 From: Justin Gross Date: Fri, 9 Oct 2026 10:25:08 -0400 Subject: [PATCH 5/6] doc(orchestration): describe only the factory's own run choice The comment in OrchestratorFactory::stream restated how an Agent picks its run, which StreamingAgent::stream already states as the contract both implementors keep. It now says only what is particular to the factory: the run it was built for, or a fresh one. Ref: GH-778 --- crates/aura/src/orchestration/factory.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/crates/aura/src/orchestration/factory.rs b/crates/aura/src/orchestration/factory.rs index 17bb50c25..56f1d51fa 100644 --- a/crates/aura/src/orchestration/factory.rs +++ b/crates/aura/src/orchestration/factory.rs @@ -218,8 +218,8 @@ 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. + // The run this factory was built for, or a fresh one when it was + // built for none. let run_id = self.agent_config.run_id.unwrap_or_else(RunId::mint); if run_id.to_string() != request_id { tracing::debug!( From f4e27792c75a53184df188a9afda028e36a173b9 Mon Sep 17 00:00:00 2001 From: Justin Gross Date: Fri, 9 Oct 2026 18:42:36 -0400 Subject: [PATCH 6/6] refactor(orchestration): type persistence's run id apart from the run's Orchestration persistence names a run by an id of its own, which the park owner key, the worker approval scope, and a stored approval record carry. Re-exporting aura_events::RunId as orchestration's RunId gave that value the run's own type, so nothing told the two apart. Those sites now hold a PersistenceRunId, and aura::RunId names only the run's id. The persistence id refuses the nil UUID as the run's does; its string and wire forms are unchanged. RunParked.run_id stays a string until GH-780 decides whether a resumed run keeps its id. Ref: GH-778 --- crates/aura/src/hitl/decision.rs | 6 +- crates/aura/src/hitl/events.rs | 4 +- crates/aura/src/hitl/route.rs | 2 +- crates/aura/src/lib.rs | 4 +- crates/aura/src/orchestration/mod.rs | 4 +- crates/aura/src/orchestration/orchestrator.rs | 4 +- crates/aura/src/orchestration/park/guard.rs | 10 +-- .../src/orchestration/persistence_wrapper.rs | 4 +- crates/aura/src/orchestration/types.rs | 73 ++++++++++++++++++- crates/aura/src/session_store/record.rs | 4 +- 10 files changed, 93 insertions(+), 22 deletions(-) diff --git a/crates/aura/src/hitl/decision.rs b/crates/aura/src/hitl/decision.rs index e975ebf0f..881365710 100644 --- a/crates/aura/src/hitl/decision.rs +++ b/crates/aura/src/hitl/decision.rs @@ -17,7 +17,7 @@ use uuid::Uuid; use crate::request_cancellation::RequestCancelToken; use crate::config::SessionId; -use crate::orchestration::{RunId, TaskIdentity}; +use crate::orchestration::{PersistenceRunId, TaskIdentity}; /// Wall-clock timestamp. /// @@ -92,14 +92,14 @@ pub enum AgentScope { session_id: Option, }, Worker { - run_id: RunId, + run_id: PersistenceRunId, task: TaskIdentity, session_id: Option, }, /// Future coordinator-mediated surface, declared now, constructed by no /// current code path. Coordinator { - run_id: RunId, + run_id: PersistenceRunId, }, } diff --git a/crates/aura/src/hitl/events.rs b/crates/aura/src/hitl/events.rs index ef426e12d..fd6521546 100644 --- a/crates/aura/src/hitl/events.rs +++ b/crates/aura/src/hitl/events.rs @@ -296,12 +296,12 @@ fn outcome_to_wire(outcome: &ApprovalOutcome) -> ApprovalOutcomeWire { mod tests { use super::*; use crate::hitl::decision::{AgentScope, CancelReason, DecisionId}; - use crate::orchestration::{RunId, TaskIdentity}; + use crate::orchestration::{PersistenceRunId, TaskIdentity}; #[test] fn sender_dropped_outcome_serializes_as_sender_dropped() { let id = DecisionId::generate(); - let run_id: RunId = "0191e8c0-1111-7000-8000-000000000000".parse().unwrap(); + let run_id: PersistenceRunId = "0191e8c0-1111-7000-8000-000000000000".parse().unwrap(); let scope = AgentScope::Worker { run_id, task: TaskIdentity::new(0, Some("ops".to_string())), diff --git a/crates/aura/src/hitl/route.rs b/crates/aura/src/hitl/route.rs index e30815c90..4ec3f5404 100644 --- a/crates/aura/src/hitl/route.rs +++ b/crates/aura/src/hitl/route.rs @@ -857,7 +857,7 @@ mod tests { #[test] fn worker_request_wire_shape_flattens_task_and_keeps_session() { - let run_id: crate::orchestration::RunId = + let run_id: crate::orchestration::PersistenceRunId = "0191e8c0-1111-7000-8000-000000000000".parse().unwrap(); let request = ApprovalRequest { version: PROTOCOL_VERSION, diff --git a/crates/aura/src/lib.rs b/crates/aura/src/lib.rs index c579380be..6b35e7359 100644 --- a/crates/aura/src/lib.rs +++ b/crates/aura/src/lib.rs @@ -55,7 +55,7 @@ pub use builder::{ Agent, AgentBuilder, BeginRunError, FilesystemTools, PreparedAgent, RunToolFactory, build_streaming_agent, build_streaming_agent_with_tools, no_run_tools, }; -pub use config::{AgentRuntimeConfig, SessionId, ToolContextFactory}; +pub use config::{AgentRuntimeConfig, RunId, SessionId, ToolContextFactory}; // Pure config types are owned by `aura-config` and re-exported here for // ergonomic consumption (`aura::LlmConfig`, etc.). pub use aura_config::{ @@ -70,7 +70,7 @@ pub use orchestration::tools::{ }; pub use orchestration::{ ArtifactsConfig, EventContext, OrchestrationConfig, OrchestrationStreamEvent, Orchestrator, - OrchestratorFactory, Plan, PlanningResponse, RoutingMode, RunId, Task, TaskIdentity, TaskJson, + OrchestratorFactory, Plan, PlanningResponse, RoutingMode, Task, TaskIdentity, TaskJson, TaskState, TaskStatus, TimeoutsConfig, agent_info, agent_info_with_tools, summarize_tools, worker_overview, }; diff --git a/crates/aura/src/orchestration/mod.rs b/crates/aura/src/orchestration/mod.rs index 6c147b85f..4e293a8fa 100644 --- a/crates/aura/src/orchestration/mod.rs +++ b/crates/aura/src/orchestration/mod.rs @@ -90,6 +90,6 @@ pub use prompt_constants::{context, fields, sections}; #[cfg(test)] pub(crate) use test_rig::{ScriptedAgent, ScriptedCompletionModel, ScriptedTurn}; pub use types::{ - BlockedCell, CellOutcome, ParkSnapshot, PendingCall, Plan, PlanningResponse, RunId, StepInput, - StructuredTaskOutput, Task, TaskIdentity, TaskJson, TaskState, TaskStatus, + BlockedCell, CellOutcome, ParkSnapshot, PendingCall, PersistenceRunId, Plan, PlanningResponse, + StepInput, StructuredTaskOutput, Task, TaskIdentity, TaskJson, TaskState, TaskStatus, }; diff --git a/crates/aura/src/orchestration/orchestrator.rs b/crates/aura/src/orchestration/orchestrator.rs index 61d46299e..d458c8b61 100644 --- a/crates/aura/src/orchestration/orchestrator.rs +++ b/crates/aura/src/orchestration/orchestrator.rs @@ -966,7 +966,7 @@ impl Orchestrator { let p = self.persistence.lock().await; (p.run_id().to_string(), p.session_id().map(String::from)) }; - let run_id = run_id_str.parse::().map_err( + let run_id = run_id_str.parse::().map_err( |e| -> Box { format!("HITL: orchestration run id '{run_id_str}' is not a valid UUID: {e}") .into() @@ -1213,7 +1213,7 @@ impl Orchestrator { (p.run_id().to_string(), p.session_id().map(String::from)) }; run_id - .parse::() + .parse::() .ok() .map(|run_id| crate::hitl::AgentScope::Worker { run_id, diff --git a/crates/aura/src/orchestration/park/guard.rs b/crates/aura/src/orchestration/park/guard.rs index ec2f2e05d..a6c74b8ea 100644 --- a/crates/aura/src/orchestration/park/guard.rs +++ b/crates/aura/src/orchestration/park/guard.rs @@ -83,7 +83,7 @@ mod tests { AgentScope, ApprovalItem, ApprovalOrigin, ApprovalRequest, DecisionId, PROTOCOL_VERSION, ParkedApproval, PendingApprovals, }; - use crate::orchestration::{RunId, TaskIdentity, run_owner_id}; + use crate::orchestration::{PersistenceRunId, TaskIdentity, run_owner_id}; use crate::session_store::{ApprovalStore, InMemoryApprovalStore, InMemoryEventBus}; fn registry_with_store() -> (PendingApprovals, Arc) { @@ -95,7 +95,7 @@ mod tests { (registry, store) } - fn worker_scope(run_id: RunId) -> AgentScope { + fn worker_scope(run_id: PersistenceRunId) -> AgentScope { AgentScope::Worker { run_id, task: TaskIdentity::new(0, Some("operations".to_string())), @@ -143,7 +143,7 @@ mod tests { #[tokio::test] 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 run_id: PersistenceRunId = "0191e8c0-2222-7000-8000-000000000042".parse().unwrap(); let this_run = aura_events::RunId::mint(); let (run, mut events) = crate::run_context::RunContext::channel(this_run); @@ -193,7 +193,7 @@ mod tests { #[tokio::test] async fn published_guard_drop_leaves_approvals_parked() { let (registry, store) = registry_with_store(); - let run_id: RunId = "0191e8c0-3333-7000-8000-000000000042".parse().unwrap(); + let run_id: PersistenceRunId = "0191e8c0-3333-7000-8000-000000000042".parse().unwrap(); let scope = worker_scope(run_id); let decision_id = DecisionId::generate(); registry @@ -220,7 +220,7 @@ mod tests { #[tokio::test] async fn unrecorded_guard_drop_is_inert() { let (registry, store) = registry_with_store(); - let run_id: RunId = "0191e8c0-4444-7000-8000-000000000042".parse().unwrap(); + let run_id: PersistenceRunId = "0191e8c0-4444-7000-8000-000000000042".parse().unwrap(); let other = DecisionId::generate(); // The ticket belongs to this run's owner id, so the arming condition // is load-bearing: an armed guard's drop would sweep and clear it. diff --git a/crates/aura/src/orchestration/persistence_wrapper.rs b/crates/aura/src/orchestration/persistence_wrapper.rs index 2033f1c99..e0742cfb7 100644 --- a/crates/aura/src/orchestration/persistence_wrapper.rs +++ b/crates/aura/src/orchestration/persistence_wrapper.rs @@ -1426,7 +1426,7 @@ mod tests { AgentScope, ApprovalDecision, DecisionId, DecisionRoute, HitlApprovalWrapper, PendingApprovals, }; - use crate::orchestration::{RunId, TaskIdentity}; + use crate::orchestration::{PersistenceRunId, TaskIdentity}; use crate::session_store::{ ApprovalStore, EventBus, InMemoryApprovalStore, InMemoryEventBus, }; @@ -1442,7 +1442,7 @@ mod tests { timeout: std::time::Duration::from_secs(60), }); - let run_id: RunId = "0191e8c0-1111-7000-8000-000000000000".parse().unwrap(); + let run_id: PersistenceRunId = "0191e8c0-1111-7000-8000-000000000000".parse().unwrap(); let scope = AgentScope::Worker { run_id, task: TaskIdentity::new(2, Some("k8s-agent".to_string())), diff --git a/crates/aura/src/orchestration/types.rs b/crates/aura/src/orchestration/types.rs index 2eab5b543..6cc34a1fd 100644 --- a/crates/aura/src/orchestration/types.rs +++ b/crates/aura/src/orchestration/types.rs @@ -4,9 +4,13 @@ //! 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 aura_events::InvalidRunId; use serde::{Deserialize, Serialize}; +use uuid::Uuid; use crate::hitl::DecisionId; @@ -20,7 +24,52 @@ const MAX_STEP_NESTING: usize = 2; // Domain identifiers for orchestration runs and tasks, modeled as simple types // (opaque newtypes reached through canonical conversion traits). -pub use aura_events::RunId; +/// Identifier orchestration persistence names a run by. It is not the run's +/// [`RunId`](aura_events::RunId). +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)] +#[serde(try_from = "Uuid", into = "Uuid")] +pub struct PersistenceRunId(Uuid); + +impl PersistenceRunId { + pub fn as_uuid(&self) -> &Uuid { + &self.0 + } +} + +impl TryFrom for PersistenceRunId { + type Error = InvalidRunId; + + /// The nil UUID names nothing, so it is refused. + fn try_from(uuid: Uuid) -> Result { + if uuid.is_nil() { + Err(InvalidRunId::Nil) + } else { + Ok(Self(uuid)) + } + } +} + +impl From for Uuid { + fn from(id: PersistenceRunId) -> Self { + id.0 + } +} + +impl fmt::Display for PersistenceRunId { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + self.0.fmt(f) + } +} + +impl FromStr for PersistenceRunId { + type Err = InvalidRunId; + + fn from_str(s: &str) -> Result { + Uuid::from_str(s) + .map_err(InvalidRunId::Malformed) + .and_then(Self::try_from) + } +} /// Identity of a worker task within a run. /// @@ -3118,4 +3167,26 @@ mod tests { chain_line ); } + + /// The nil UUID names nothing, so persistence's id refuses it as the + /// run's own id does; a well-formed one round-trips as the bare string. + #[test] + fn a_persistence_run_id_refuses_the_nil_uuid() { + const NIL: &str = "00000000-0000-0000-0000-000000000000"; + const ID: &str = "0191e8c0-1111-7000-8000-000000000000"; + assert!(matches!( + NIL.parse::(), + Err(aura_events::InvalidRunId::Nil) + )); + assert!("not-a-uuid".parse::().is_err()); + assert!(serde_json::from_value::(serde_json::json!(NIL)).is_err()); + + let id: PersistenceRunId = ID.parse().unwrap(); + assert_eq!(id.to_string(), ID); + assert_eq!(serde_json::to_value(id).unwrap(), serde_json::json!(ID)); + assert_eq!( + serde_json::from_value::(serde_json::json!(ID)).unwrap(), + id + ); + } } diff --git a/crates/aura/src/session_store/record.rs b/crates/aura/src/session_store/record.rs index 48d020650..cf25c1ae4 100644 --- a/crates/aura/src/session_store/record.rs +++ b/crates/aura/src/session_store/record.rs @@ -18,7 +18,7 @@ use crate::hitl::{ AgentScope, ApprovalDecision, ApprovalItem, ApprovalOrigin, ApprovalRequest, DecisionId, ParkedApproval, Timestamp, }; -use crate::orchestration::{RunId, TaskIdentity}; +use crate::orchestration::{PersistenceRunId, TaskIdentity}; /// Round-trippable storage form of a [`ParkedApproval`]. Field and tag names /// are a persisted contract shared by every instance reading the store — rename @@ -233,7 +233,7 @@ impl From for ApprovalOrigin { } } -fn parse_run_id(raw: &str) -> Result { +fn parse_run_id(raw: &str) -> Result { raw.parse().map_err(|e| InvalidRecord { reason: format!("run_id '{raw}': {e}"), })