From de512088dbf110f9527f3eb464f35efcff2ac774 Mon Sep 17 00:00:00 2001 From: Jacob Hull Date: Tue, 6 Oct 2026 14:14:06 -0700 Subject: [PATCH 1/5] refactor: brand the request id RequestId is a newtype that only RequestId::generate, RequestId::for_a2a_task, and RequestId::for_slack_message construct. It replaces the String alias and the bare request id strings in HITL, MCP cancellation, the session stores, scratchpad storage, StreamingAgent, and the web server. A run's id is its request's id. A parked approval's ApprovalOwner is the request that raised it, the orchestration run that parked it, or no one for an agent built without a request. Stored records and the webhook payload keep the bare-string request_id field. Signed-off-by: Jacob Hull --- crates/aura-cli/src/repl/mcp/wizard.rs | 4 +- crates/aura-test-utils/src/mock_agent.rs | 63 +++++-- .../aura-web-server/src/a2a/agent_executor.rs | 60 +++---- crates/aura-web-server/src/handlers.rs | 30 ++-- .../src/session_store/redis/approval_store.rs | 10 +- crates/aura-web-server/src/slack/runner.rs | 20 +-- .../aura-web-server/src/streaming/handlers.rs | 8 +- crates/aura-web-server/src/streaming/otel.rs | 4 +- crates/aura-web-server/tests/common/mod.rs | 17 +- .../tests/file_session_store_test.rs | 13 +- .../tests/redis_session_store_test.rs | 42 +++-- crates/aura/src/builder.rs | 76 +++++---- crates/aura/src/config.rs | 2 +- crates/aura/src/domain/mod.rs | 5 + crates/aura/src/domain/request_id.rs | 95 +++++++++++ crates/aura/src/hitl/decision.rs | 38 +++++ crates/aura/src/hitl/events.rs | 2 +- crates/aura/src/hitl/gate.rs | 87 ++++++---- crates/aura/src/hitl/mod.rs | 15 +- crates/aura/src/hitl/protocol.rs | 7 +- crates/aura/src/hitl/registry.rs | 41 ++--- crates/aura/src/hitl/route.rs | 47 +++--- crates/aura/src/hitl/tool.rs | 26 +-- crates/aura/src/lib.rs | 4 +- crates/aura/src/mcp/client.rs | 158 +++++++++--------- crates/aura/src/mcp/manager.rs | 58 ++++--- crates/aura/src/mcp/progress.rs | 23 ++- crates/aura/src/orchestration/factory.rs | 16 +- crates/aura/src/orchestration/mod.rs | 8 +- crates/aura/src/orchestration/orchestrator.rs | 57 ++++--- crates/aura/src/orchestration/overview.rs | 4 +- crates/aura/src/orchestration/park/commit.rs | 48 +++--- .../src/orchestration/park/continuation.rs | 4 +- crates/aura/src/orchestration/park/guard.rs | 35 ++-- crates/aura/src/orchestration/park/mod.rs | 4 +- .../src/orchestration/persistence_wrapper.rs | 12 +- crates/aura/src/orchestration/test_rig.rs | 2 +- crates/aura/src/request_cancellation.rs | 2 - crates/aura/src/rig_builder.rs | 6 +- crates/aura/src/run_context.rs | 83 ++++----- crates/aura/src/scratchpad/storage.rs | 78 ++------- crates/aura/src/scratchpad/tools.rs | 6 +- crates/aura/src/scratchpad/wrapper.rs | 72 ++------ crates/aura/src/session_store/fault_store.rs | 6 +- crates/aura/src/session_store/file.rs | 12 +- crates/aura/src/session_store/memory.rs | 34 ++-- crates/aura/src/session_store/mod.rs | 4 +- crates/aura/src/session_store/record.rs | 73 +++++++- crates/aura/src/streaming.rs | 18 +- crates/aura/src/streaming_request_hook.rs | 28 ++-- 50 files changed, 916 insertions(+), 651 deletions(-) create mode 100644 crates/aura/src/domain/mod.rs create mode 100644 crates/aura/src/domain/request_id.rs diff --git a/crates/aura-cli/src/repl/mcp/wizard.rs b/crates/aura-cli/src/repl/mcp/wizard.rs index 393c95464..350c251e2 100644 --- a/crates/aura-cli/src/repl/mcp/wizard.rs +++ b/crates/aura-cli/src/repl/mcp/wizard.rs @@ -400,9 +400,7 @@ fn verify_server( (info.status.clone(), tools) }) .ok_or_else(|| format!("`{name}` missing from the connection status snapshot")); - manager - .cancel_and_close_all("mcp-add-verify", "verification complete") - .await; + manager.close_all().await; result }) } diff --git a/crates/aura-test-utils/src/mock_agent.rs b/crates/aura-test-utils/src/mock_agent.rs index 9796ced3e..1d956068d 100644 --- a/crates/aura-test-utils/src/mock_agent.rs +++ b/crates/aura-test-utils/src/mock_agent.rs @@ -173,9 +173,9 @@ impl StreamingAgent for MockAgent { _query: &str, _chat_history: Vec, options: aura::streaming::RunOptions, - request_id: &str, + request_id: &aura::RequestId, ) -> AgentRun { - let stream = self.start(request_id).await; + let stream = self.start(request_id.as_str()).await; // Carries a caller-supplied token so `cancel_token()` returns the one the // caller named. The scripts do not race it, so cancelling does not end a // mock stream. @@ -186,7 +186,7 @@ impl StreamingAgent for MockAgent { ) } - async fn cancel_and_close_mcp(&self, _request_id: &str, _reason: &str) -> usize { + async fn cancel_and_close_mcp(&self, _request_id: &aura::RequestId, _reason: &str) -> usize { 0 } } @@ -200,7 +200,12 @@ mod tests { async fn a_pending_agent_never_yields() { let agent = MockAgent::pending(); let mut stream = agent - .stream("q", vec![], aura::streaming::RunOptions::default(), "req_1") + .stream( + "q", + vec![], + aura::streaming::RunOptions::default(), + &aura::RequestId::generate(), + ) .await .into_events(); assert!( @@ -216,7 +221,12 @@ mod tests { async fn a_yielding_agent_produces_its_items_then_ends() { let agent = MockAgent::yielding(vec![items::text("hello "), items::text("world")]); let mut stream = agent - .stream("q", vec![], aura::streaming::RunOptions::default(), "req_1") + .stream( + "q", + vec![], + aura::streaming::RunOptions::default(), + &aura::RequestId::generate(), + ) .await .into_events(); @@ -235,16 +245,24 @@ mod tests { async fn the_start_hook_runs_before_the_stream() { let ran = Arc::new(AtomicBool::new(false)); let flag = Arc::clone(&ran); - let agent = MockAgent::pending().on_stream_start(move |request_id| { + let request_id = aura::RequestId::generate(); + let expected = request_id.to_string(); + let agent = MockAgent::pending().on_stream_start(move |seen| { let flag = Arc::clone(&flag); + let expected = expected.clone(); async move { - assert_eq!(request_id, "req_1"); + assert_eq!(seen, expected); flag.store(true, Ordering::SeqCst); } }); let _ = agent - .stream("q", vec![], aura::streaming::RunOptions::default(), "req_1") + .stream( + "q", + vec![], + aura::streaming::RunOptions::default(), + &request_id, + ) .await .into_events(); @@ -253,6 +271,7 @@ mod tests { #[tokio::test(start_paused = true)] async fn effects_run_in_script_order_and_see_the_request_id() { + let request_id = aura::RequestId::generate(); let order = Arc::new(Mutex::new(Vec::new())); let effect_order = Arc::clone(&order); let agent = MockAgent::scripted(vec![ @@ -270,13 +289,16 @@ mod tests { "q", vec![], aura::streaming::RunOptions::default(), - "req_42", + &request_id, ) .await .into_events(); let items: Vec<_> = stream.collect().await; - assert_eq!(order.lock().expect("order lock").as_slice(), ["req_42"]); + assert_eq!( + order.lock().expect("order lock").as_slice(), + [request_id.to_string()] + ); assert_eq!(items.len(), 1, "effects do not yield stream items"); } @@ -285,13 +307,23 @@ mod tests { let agent = MockAgent::yielding([items::text("once")]); let first: Vec<_> = agent - .stream("q", vec![], aura::streaming::RunOptions::default(), "req_1") + .stream( + "q", + vec![], + aura::streaming::RunOptions::default(), + &aura::RequestId::generate(), + ) .await .into_events() .collect() .await; let second: Vec<_> = agent - .stream("q", vec![], aura::streaming::RunOptions::default(), "req_1") + .stream( + "q", + vec![], + aura::streaming::RunOptions::default(), + &aura::RequestId::generate(), + ) .await .into_events() .collect() @@ -318,7 +350,12 @@ mod tests { ]); let mut stream = agent - .stream("q", vec![], aura::streaming::RunOptions::default(), "req_1") + .stream( + "q", + vec![], + aura::streaming::RunOptions::default(), + &aura::RequestId::generate(), + ) .await .into_events(); while stream.next().await.is_some() { diff --git a/crates/aura-web-server/src/a2a/agent_executor.rs b/crates/aura-web-server/src/a2a/agent_executor.rs index 2fed3d215..e4e1b6be9 100644 --- a/crates/aura-web-server/src/a2a/agent_executor.rs +++ b/crates/aura-web-server/src/a2a/agent_executor.rs @@ -16,7 +16,7 @@ use a2a::{ }; use a2a_server::{AgentExecutor, ExecutorContext, TaskStore}; use aura::RigBuilder; -use aura::{StreamItem, StreamedAssistantContent, StreamingAgent}; +use aura::{RequestId, StreamItem, StreamedAssistantContent, StreamingAgent}; use futures_util::StreamExt; use futures_util::stream::BoxStream; use serde_json::Value; @@ -46,7 +46,7 @@ pub struct AuraAgentExecutor { struct TaskCancelEntry { token: CancellationToken, agent: Option>, - request_id: String, + request_id: RequestId, } /// Whether `cancel()` got to a task before the code asking. @@ -269,7 +269,7 @@ impl AgentExecutor for AuraAgentExecutor { metadata: None, })); - let request_id = format!("a2a_{}", task_id); + let request_id = RequestId::for_a2a_task(&task_id); // Registered before the agent build and history fetch, both of which // await, so a cancelTask during those has a token to cancel. Its @@ -358,13 +358,13 @@ impl AgentExecutor for AuraAgentExecutor { let Some(item) = next else { break }; match item { Ok(StreamItem::StreamAssistantItem(StreamedAssistantContent::Text(t))) => { - event!(Level::DEBUG, request_id, t, "stream content received"); + event!(Level::DEBUG, %request_id, t, "stream content received"); let append = append_tracker.entry((task_id.clone(), context_id.clone(), RESPONSE_ARTIFACT_ID.to_owned())) .and_modify(|e| *e = true) .or_insert(false); - event!(Level::DEBUG, request_id, "response returned and should be appended: {}", *append); + event!(Level::DEBUG, %request_id, "response returned and should be appended: {}", *append); let artifact = Artifact { artifact_id: RESPONSE_ARTIFACT_ID.to_owned(), @@ -384,7 +384,7 @@ impl AgentExecutor for AuraAgentExecutor { })); } Ok(StreamItem::StreamAssistantItem(StreamedAssistantContent::ToolCall(tc))) => { - event!(Level::DEBUG, request_id, tool_name = tc.name.as_str(), "tool call received"); + event!(Level::DEBUG, %request_id, tool_name = tc.name.as_str(), "tool call received"); let artifact_id: String = format!("tool_call_{}", tc.id); let append = append_tracker.entry((task_id.clone(), context_id.clone(), artifact_id.to_owned())) @@ -414,7 +414,7 @@ impl AgentExecutor for AuraAgentExecutor { })); } Ok(StreamItem::StreamAssistantItem(StreamedAssistantContent::Reasoning(r))) => { - event!(Level::DEBUG, request_id, reasoning = r, "reasoning received"); + event!(Level::DEBUG, %request_id, reasoning = r, "reasoning received"); reasoning_num += 1; let artifact_id: String = format!("reasoning_{}", reasoning_num); @@ -440,7 +440,7 @@ impl AgentExecutor for AuraAgentExecutor { })); } Ok(StreamItem::ScratchpadUsage { agent_id, tokens_intercepted, tokens_extracted }) => { - event!(Level::DEBUG, request_id, "scratchpad usage"); + event!(Level::DEBUG, %request_id, "scratchpad usage"); let artifact_id: String = format!("scratchpad_{}", agent_id); let append = append_tracker.entry((task_id.clone(), context_id.clone(), artifact_id.to_owned())) @@ -468,25 +468,25 @@ impl AgentExecutor for AuraAgentExecutor { })); } Ok(StreamItem::TurnUsage(..)) => { - event!(Level::DEBUG, request_id, "turn usage"); + event!(Level::DEBUG, %request_id, "turn usage"); } Ok(StreamItem::ContextUsage { .. }) => { - event!(Level::DEBUG, request_id, "context usage"); + event!(Level::DEBUG, %request_id, "context usage"); } Ok(StreamItem::AgentEvent(_)) => { - event!(Level::DEBUG, request_id, "orchestration event"); + event!(Level::DEBUG, %request_id, "orchestration event"); } Ok(StreamItem::McpStatus(_)) => { - event!(Level::DEBUG, request_id, "mcp status"); + event!(Level::DEBUG, %request_id, "mcp status"); } Ok(StreamItem::StreamUserItem(_)) => { - event!(Level::DEBUG, request_id, "stream user item"); + event!(Level::DEBUG, %request_id, "stream user item"); } Ok(StreamItem::StreamAssistantItem(StreamedAssistantContent::ToolCallDelta { .. })) => { - event!(Level::DEBUG, request_id, "stream assistant item"); + event!(Level::DEBUG, %request_id, "stream assistant item"); } Ok(StreamItem::StreamAssistantItem(StreamedAssistantContent::ReasoningDelta { .. })) => { - event!(Level::DEBUG, request_id, "reasoning delta"); + event!(Level::DEBUG, %request_id, "reasoning delta"); } Ok(StreamItem::Final(final_info)) => { let append = append_tracker.entry((task_id.clone(), context_id.clone(), FINAL_ARTIFACT_ID.to_owned())) @@ -516,11 +516,11 @@ impl AgentExecutor for AuraAgentExecutor { break; // done processing } Ok(StreamItem::FinalMarker) => { - event!(Level::DEBUG, request_id, "stream final marker"); + event!(Level::DEBUG, %request_id, "stream final marker"); break; // done processing } Err(e) => { - event!(Level::ERROR, request_id, task_id, error = e.to_string(), "stream error"); + event!(Level::ERROR, %request_id, task_id, error = e.to_string(), "stream error"); yield Ok(fail_status(&task_id, &context_id, &e.to_string())); success = false; break; // done processing @@ -642,7 +642,7 @@ pub(super) fn fail_status(task_id: &str, context_id: &str, error_msg: &str) -> S async fn get_history_for_context( task_store: SharedTaskStore, - request_id: &str, + request_id: &RequestId, context_id: &str, task_id: &str, ) -> Result, A2AError> { @@ -652,7 +652,7 @@ async fn get_history_for_context( loop { event!( Level::DEBUG, - request_id, + %request_id, context_id, "processing history for context" ); @@ -673,7 +673,7 @@ async fn get_history_for_context( event!( Level::DEBUG, - request_id, + %request_id, context_id, "found {} tasks, continue token '{}'", page.tasks.len(), @@ -695,7 +695,7 @@ async fn get_history_for_context( let chat_history: Vec = tasks.iter().flat_map(task_turns).collect(); - event!(Level::DEBUG, request_id, context_id, chat_history = ?chat_history, "determined this following history to use"); + event!(Level::DEBUG, %request_id, context_id, chat_history = ?chat_history, "determined this following history to use"); Ok(chat_history) } @@ -909,7 +909,7 @@ mod tests { #[test] fn dropping_the_cancel_guard_releases_the_entry() { let task_id = format!("t_{}", uuid::Uuid::new_v4()); - let request_id = format!("a2a_{task_id}"); + let request_id = RequestId::for_a2a_task(&task_id); let state: Arc = Arc::new(Mutex::new(HashMap::new())); let guard = TaskCancelGuard { @@ -943,7 +943,7 @@ mod tests { TaskCancelEntry { token: CancellationToken::new(), agent: None, - request_id: format!("a2a_{task_id}"), + request_id: RequestId::for_a2a_task(&task_id), }, ); (state, task_id) @@ -1004,7 +1004,7 @@ mod tests { TaskCancelEntry { token: CancellationToken::new(), agent: None, - request_id: "a2a_t".to_string(), + request_id: RequestId::for_a2a_task("t"), }, ); lock_cancel_state(&state).remove("t"); @@ -1020,7 +1020,7 @@ mod tests { #[test] fn claiming_the_agent_reports_a_cancel_that_already_ran() { let task_id = format!("t_{}", uuid::Uuid::new_v4()); - let request_id = format!("a2a_{task_id}"); + let request_id = RequestId::for_a2a_task(&task_id); let state: Arc = Arc::new(Mutex::new(HashMap::new())); let agent: Arc = Arc::new(MockAgent::pending()); @@ -1051,7 +1051,7 @@ mod tests { #[test] fn an_early_return_before_the_run_releases_the_entry() { let task_id = format!("t_{}", uuid::Uuid::new_v4()); - let request_id = format!("a2a_{task_id}"); + let request_id = RequestId::for_a2a_task(&task_id); let state: Arc = Arc::new(Mutex::new(HashMap::new())); { @@ -1161,7 +1161,7 @@ mod tests { for task in tasks { store.create(task).await.expect("task created"); } - get_history_for_context(store, "req", "ctx", executing) + get_history_for_context(store, &RequestId::for_a2a_task(executing), "ctx", executing) .await .expect("history built") } @@ -1319,7 +1319,7 @@ mod loom_tests { TaskCancelEntry { token: CancellationToken::new(), agent: None, - request_id: "a2a_t1".to_string(), + request_id: RequestId::for_a2a_task("t1"), }, ); @@ -1372,7 +1372,7 @@ mod loom_tests { TaskCancelEntry { token: CancellationToken::new(), agent: None, - request_id: "a2a_t1".to_string(), + request_id: RequestId::for_a2a_task("t1"), }, ); @@ -1423,7 +1423,7 @@ mod loom_tests { TaskCancelEntry { token: CancellationToken::new(), agent: None, - request_id: "a2a_t1".to_string(), + request_id: RequestId::for_a2a_task("t1"), }, ); diff --git a/crates/aura-web-server/src/handlers.rs b/crates/aura-web-server/src/handlers.rs index a5e1ccffd..7e79ecd70 100644 --- a/crates/aura-web-server/src/handlers.rs +++ b/crates/aura-web-server/src/handlers.rs @@ -1,6 +1,7 @@ use a2a::VERSION; use aura::RigBuilder; -use aura::{ResponseContent, StreamingAgent, UsageState}; +use aura::hitl::ApprovalOwner; +use aura::{RequestId, ResponseContent, StreamingAgent, UsageState}; use aura_events::{AgentInfo, NativeToolOverview, ServerInfo}; use axum::Json; use axum::body::Body; @@ -24,12 +25,12 @@ use crate::types::*; /// Guard over a request's pending approvals. struct RequestResourceGuard { - request_id: String, + request_id: RequestId, pending_approvals: aura::hitl::PendingApprovals, } impl RequestResourceGuard { - fn new(request_id: String, pending_approvals: aura::hitl::PendingApprovals) -> Self { + fn new(request_id: RequestId, pending_approvals: aura::hitl::PendingApprovals) -> Self { Self { request_id, pending_approvals, @@ -42,18 +43,17 @@ impl Drop for RequestResourceGuard { fn drop(&mut self) { // Synchronous so parked awaits cancel even when the runtime is // shutting down and the spawn below never polls. - self.pending_approvals - .cancel_request_local(&self.request_id); + let owner = ApprovalOwner::Request(self.request_id.clone()); + self.pending_approvals.cancel_request_local(&owner); // Use try_current to avoid panic during runtime shutdown if let Ok(handle) = tokio::runtime::Handle::try_current() { - let id = self.request_id.clone(); let pending_approvals = self.pending_approvals.clone(); // Instrument with the span current at drop so cleanup events // stay parented to the request's trace. let cleanup = tracing::Instrument::instrument( async move { - pending_approvals.cancel_request(&id).await; + pending_approvals.cancel_request(&owner).await; }, tracing::Span::current(), ); @@ -111,7 +111,7 @@ impl std::fmt::Display for PrepareError { /// Both streaming and non-streaming handlers delegate to the same core logic with different /// delivery modes but the same observability instrumentation and stream processing. pub struct CompletionConfig { - pub request_id: String, + pub request_id: RequestId, pub timeout_duration: Option, pub first_chunk_timeout: Option, pub inactivity_timeout: Option, @@ -176,7 +176,7 @@ async fn build_agent_for_request( req_headers: &HashMap, additional_tools: Vec>, client_tools: Option<&[ClientToolDefinition]>, - request_id: String, + request_id: RequestId, session_id: String, ) -> Result, PrepareError> { let client_tool_defs = @@ -214,7 +214,7 @@ pub struct RequestSetup { /// invokes a passthrough tool. pub has_client_tools: bool, /// Request id (`req_…`) shared by the agent build and the completion stream. - pub request_id: String, + pub request_id: RequestId, /// OpenAI-compatible `user` field, for the `user.id` span attribute. pub user_id: Option, /// Request `metadata` serialized as a JSON object string, for the @@ -247,7 +247,7 @@ pub async fn prepare_request( // 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()); + let request_id = RequestId::generate(); // 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). @@ -770,7 +770,7 @@ async fn handle_non_streaming_completion( { let span = tracing::Span::current(); - aura::logging::set_span_attribute(&span, "http.request_id", config.request_id.clone()); + aura::logging::set_span_attribute(&span, "http.request_id", config.request_id.to_string()); } let response_ctx = ResponseContext { @@ -818,7 +818,7 @@ async fn handle_streaming_completion( { let span = tracing::Span::current(); - aura::logging::set_span_attribute(&span, "http.request_id", config.request_id.clone()); + aura::logging::set_span_attribute(&span, "http.request_id", config.request_id.to_string()); } let chat_session_id = setup.chat_session_id.clone(); @@ -2621,7 +2621,7 @@ url = "http://127.0.0.1:9" version: aura::hitl::PROTOCOL_VERSION, instance_id: "test-instance".to_string(), decision_id: aura::hitl::DecisionId::generate(), - request_id: "req-smoke".into(), + owner: aura::hitl::ApprovalOwner::Request(aura::RequestId::generate()), scope: aura::hitl::AgentScope::Single { session_id: None }, origin: aura::hitl::ApprovalOrigin::ConfigGate { matched_pattern: "test_*".into(), @@ -2830,7 +2830,7 @@ url = "http://127.0.0.1:9" version: aura::hitl::PROTOCOL_VERSION, instance_id: "test-instance".to_string(), decision_id: aura::hitl::DecisionId::generate(), - request_id: "req-hmac".into(), + owner: aura::hitl::ApprovalOwner::Request(aura::RequestId::generate()), scope: aura::hitl::AgentScope::Single { session_id: None }, origin: aura::hitl::ApprovalOrigin::ConfigGate { matched_pattern: "test_*".into(), diff --git a/crates/aura-web-server/src/session_store/redis/approval_store.rs b/crates/aura-web-server/src/session_store/redis/approval_store.rs index 0d59ca6ca..64ad9108c 100644 --- a/crates/aura-web-server/src/session_store/redis/approval_store.rs +++ b/crates/aura-web-server/src/session_store/redis/approval_store.rs @@ -22,7 +22,7 @@ use std::sync::LazyLock; use async_trait::async_trait; -use aura::hitl::{ApprovalDecision, DecisionId, ParkedApproval, ResolveError}; +use aura::hitl::{ApprovalDecision, ApprovalOwner, DecisionId, ParkedApproval, ResolveError}; use aura::session_store::{ApprovalStore, DecisionRecord, ParkedApprovalRecord, SessionStoreError}; use redis::AsyncCommands; use redis::aio::ConnectionManager; @@ -94,8 +94,8 @@ impl RedisApprovalStore { format!("{}:approval:decision:{decision_id}", self.key_prefix) } - fn req_key(&self, request_id: &str) -> String { - format!("{}:approval:req:{request_id}", self.key_prefix) + fn req_key(&self, owner: impl std::fmt::Display) -> String { + format!("{}:approval:req:{owner}", self.key_prefix) } /// `GETDEL` a record and prune its request index best-effort, returning @@ -204,9 +204,9 @@ impl ApprovalStore for RedisApprovalStore { async fn cancel_request( &self, - request_id: &str, + owner: &ApprovalOwner, ) -> Result, SessionStoreError> { - let req_key = self.req_key(request_id); + let req_key = self.req_key(owner); let mut conn = self.conn.clone(); let ids: Vec = conn.smembers(&req_key).await.map_err(request_err)?; if ids.is_empty() { diff --git a/crates/aura-web-server/src/slack/runner.rs b/crates/aura-web-server/src/slack/runner.rs index abbf0d9c7..4e1cbf0b1 100644 --- a/crates/aura-web-server/src/slack/runner.rs +++ b/crates/aura-web-server/src/slack/runner.rs @@ -3,7 +3,7 @@ use std::sync::Arc; use std::time::Duration; -use aura::{Message, RigBuilder, RunToolFactory, StreamItem, ToolDyn, no_run_tools}; +use aura::{Message, RequestId, RigBuilder, RunToolFactory, StreamItem, ToolDyn, no_run_tools}; use futures_util::StreamExt; use tokio::sync::{Semaphore, mpsc}; use tracing::{Instrument, debug, error, info, warn}; @@ -264,14 +264,14 @@ 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); + let request_id = RequestId::for_slack_message(&inbound.channel, &inbound.ts); let earlier = match prefetched { Some(earlier) => earlier, None => { let earlier = match self.earlier_messages(&inbound).await { Ok(earlier) => earlier, Err(e) => { - error!(request_id, error = %e, "could not read slack history"); + error!(%request_id, error = %e, "could not read slack history"); return; } }; @@ -288,7 +288,7 @@ impl SlackIngress { .add_reaction(&inbound.channel, &inbound.ts, ACK_REACTION) .await { - warn!(request_id, error = %e, "could not react to slack message"); + warn!(%request_id, error = %e, "could not react to slack message"); } let reply = match self.run_agent(&inbound, &earlier, &request_id).await { @@ -298,13 +298,13 @@ impl SlackIngress { // tell the user something broke when nothing did. Err(RunError::Cancelled) => { warn!( - request_id, + %request_id, "slack-triggered agent run cancelled by shutdown" ); return; } Err(e) => { - error!(request_id, error = %e, "slack-triggered agent run failed"); + error!(%request_id, error = %e, "slack-triggered agent run failed"); FAILED_REPLY.to_owned() } }; @@ -313,7 +313,7 @@ impl SlackIngress { .post_message(&inbound.channel, inbound.reply_thread(), &reply) .await { - error!(request_id, error = %e, "could not post slack reply"); + error!(%request_id, error = %e, "could not post slack reply"); } } @@ -344,7 +344,7 @@ impl SlackIngress { &self, inbound: &Inbound, earlier: &[SlackMessage], - request_id: &str, + request_id: &RequestId, ) -> Result { let history = thread_history(earlier, &self.identity, &inbound.ts); let session_id = match inbound.reply_thread() { @@ -353,7 +353,7 @@ impl SlackIngress { }; let (config, tools) = run_setup(&self.config, &self.api, inbound); debug!( - request_id, + %request_id, search = inbound.action_token.is_some(), "slack run tools" ); @@ -363,7 +363,7 @@ impl SlackIngress { None, Some(session_id), None, - Some(request_id.to_owned()), + Some(request_id.clone()), tools, ) .await diff --git a/crates/aura-web-server/src/streaming/handlers.rs b/crates/aura-web-server/src/streaming/handlers.rs index 28d8c352b..8447a2a17 100644 --- a/crates/aura-web-server/src/streaming/handlers.rs +++ b/crates/aura-web-server/src/streaming/handlers.rs @@ -42,7 +42,7 @@ use tokio::sync::mpsc; /// Context for cancellation and cleanup callbacks. pub struct StreamingCallbacks { /// The request's id. - pub request_id: String, + pub request_id: aura::RequestId, /// Agent reference for MCP cleanup (cancel_and_close_mcp) pub agent: Arc, /// The run's own events. @@ -2405,7 +2405,7 @@ mod tests { let (run_events_tx, run_events) = mpsc::channel(8); ( StreamingCallbacks { - request_id: "req_inactivity_test".to_string(), + request_id: aura::RequestId::generate(), agent: Arc::new(MockAgent::pending()), agent_events: Some(run_events), usage_state: aura::UsageState::new(), @@ -2873,7 +2873,7 @@ mod tests { ( Senders { run_events_tx }, StreamingCallbacks { - request_id: "req_tool_events".to_string(), + request_id: aura::RequestId::generate(), agent: Arc::new(MockAgent::pending()), agent_events: Some(run_events), usage_state: UsageState::new(), @@ -2928,7 +2928,7 @@ mod tests { "q", vec![], aura::streaming::RunOptions::default(), - "req_tool_events", + &aura::RequestId::generate(), ) .await .into_events(); diff --git a/crates/aura-web-server/src/streaming/otel.rs b/crates/aura-web-server/src/streaming/otel.rs index 7e5a482e6..9ac4d0d56 100644 --- a/crates/aura-web-server/src/streaming/otel.rs +++ b/crates/aura-web-server/src/streaming/otel.rs @@ -22,7 +22,7 @@ use super::StreamTermination; pub struct StreamOtelContext { pub provider: String, pub model: String, - pub request_id: String, + pub request_id: aura::RequestId, pub session_id: String, pub query: String, /// OpenAI-compatible `user` field from the request. @@ -50,7 +50,7 @@ impl StreamOtelContext { if let Some(system_prompt) = &self.system_prompt { aura::logging::set_system_prompt_attribute(&span, system_prompt); } - aura::logging::set_span_attribute(&span, "http.request_id", self.request_id.clone()); + aura::logging::set_span_attribute(&span, "http.request_id", self.request_id.to_string()); aura::logging::set_span_attribute(&span, "session.id", self.session_id.clone()); aura::logging::set_span_attribute( &span, diff --git a/crates/aura-web-server/tests/common/mod.rs b/crates/aura-web-server/tests/common/mod.rs index 2f4fc12b6..25628ace0 100644 --- a/crates/aura-web-server/tests/common/mod.rs +++ b/crates/aura-web-server/tests/common/mod.rs @@ -12,20 +12,25 @@ use std::sync::Arc; use std::time::Duration; use aura::hitl::{ - AgentScope, ApprovalDecision, ApprovalItem, ApprovalOrigin, ApprovalRequest, DecisionId, - PROTOCOL_VERSION, ParkedApproval, ResolveError, + AgentScope, ApprovalDecision, ApprovalItem, ApprovalOrigin, ApprovalOwner, ApprovalRequest, + DecisionId, PROTOCOL_VERSION, ParkedApproval, ResolveError, }; use aura::session_store::{ApprovalStore, ParkedApprovalRecord}; -/// A representative parked approval, expiring in `ttl`. -pub fn make_parked(request_id: &str, ttl: Duration) -> ParkedApproval { +/// The same label always names the same owner; it stores as `a2a_