Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions crates/aura-web-server/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
78 changes: 74 additions & 4 deletions crates/aura-web-server/src/a2a/agent_executor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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",

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

agent.stream exports as an OpenInference LLM span (openinference_exporter.rs), but nothing records the model, input, or output on it, the way chat does through StreamOtelContext. With OTel on, every A2A task adds an empty LLM root span in Phoenix.

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();

Expand All @@ -269,7 +280,9 @@ impl AgentExecutor for AuraAgentExecutor {
metadata: None,
}));

let request_id = format!("a2a_{}", task_id);
// In string form, the run's id is the request id everything
// request-keyed reads: MCP cancellation, approvals, the cancel map.
let request_id = run_id.to_string();

// Registered before the agent build and history fetch, both of which
// await, so a cancelTask during those has a token to cancel. Its
Expand All @@ -295,7 +308,7 @@ impl AgentExecutor for AuraAgentExecutor {
Some(&req_headers),
session_id,
None,
Some(request_id.clone()),
Some(run_id),
run_tools,
)
.await
Expand Down Expand Up @@ -558,7 +571,8 @@ impl AgentExecutor for AuraAgentExecutor {
metadata: None,
}));
}
})
};
in_span(span, Box::pin(execution))
}

fn cancel(&self, ctx: ExecutorContext) -> BoxStream<'static, Result<StreamResponse, A2AError>> {
Expand Down Expand Up @@ -621,6 +635,18 @@ fn extract_text(parts: Vec<Part>) -> Result<String, A2AError> {
Ok(strings.join("\n"))
}

/// `stream`, polled inside `span`, so its log lines carry the span's fields.
fn in_span<T: 'static>(
span: tracing::Span,
mut stream: BoxStream<'static, T>,

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: in_span takes a BoxStream and boxes it again, and the caller already passes Box::pin(execution). Taking S: Stream + Send + 'static instead would avoid the second box.

) -> 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(),
Expand Down Expand Up @@ -1046,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::<u8>::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::<Vec<_>>().await, vec![1]);

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This builds its own span and calls in_span directly, so it never exercises execute(). It would still pass if execute() stopped wrapping the stream, or dropped run.id, a2a.task_id, or a2a.context_id from the span.


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<std::sync::Mutex<Vec<u8>>>);

impl std::io::Write for SharedWriter {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
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]
Expand Down
25 changes: 14 additions & 11 deletions crates/aura-web-server/src/handlers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -176,7 +176,7 @@ async fn build_agent_for_request(
req_headers: &HashMap<String, String>,
additional_tools: Vec<Box<dyn aura::ToolDyn>>,
client_tools: Option<&[ClientToolDefinition]>,
request_id: String,
run_id: aura::RunId,
session_id: String,
) -> Result<Arc<aura::Agent>, PrepareError> {
let client_tool_defs =
Expand All @@ -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
Expand Down Expand Up @@ -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:

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: the same explanation of how the run id's string form keys the registries appears here, in the A2A executor, in the Slack runner, and on RunContext::has_id. CLAUDE.md asks for it in one place, referenced from the others.

// 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).
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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?;
Expand Down Expand Up @@ -768,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);

Expand Down Expand Up @@ -813,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,
Expand All @@ -822,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);

Expand Down
41 changes: 26 additions & 15 deletions crates/aura-web-server/src/slack/runner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -235,6 +235,16 @@ fn queue_answer(ingress: &Arc<SlackIngress>, inbound: Inbound, earlier: Option<V
);
return;
};
// The answer runs inside this span, so its log lines, the agent's
// included, carry the run's id and the message it answers.
let run_id = aura::RunId::mint();
let span = tracing::info_span!(
parent: None,
"agent.stream",
run.id = %run_id,
slack.channel = %inbound.channel,
slack.ts = %inbound.ts,
);
let ingress = Arc::clone(ingress);
let tracker = Arc::clone(&ingress.state.active_requests);
let active = ActiveRequestGuard::new(Arc::clone(&tracker));
Expand All @@ -247,12 +257,10 @@ fn queue_answer(ingress: &Arc<SlackIngress>, inbound: Inbound, earlier: Option<V
};
drop(waiting);
if slot.is_ok() {
ingress.answer(inbound, earlier).await;
ingress.answer(inbound, earlier, run_id).await;
}
};
tracker.track_task(tokio::spawn(
task.instrument(tracing::info_span!(parent: None, "agent.stream")),
));
tracker.track_task(tokio::spawn(task.instrument(span)));
}

impl SlackIngress {
Expand All @@ -263,8 +271,15 @@ impl SlackIngress {
/// participation verdict stands, since the trimmed read may no longer
/// hold the bot turn that earned it. A conversation read here is
/// checked here.
async fn answer(&self, inbound: Inbound, prefetched: Option<Vec<SlackMessage>>) {
let request_id = format!("slack_{}_{}", inbound.channel, inbound.ts);
async fn answer(
&self,
inbound: Inbound,
prefetched: Option<Vec<SlackMessage>>,
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();

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Verbose logs lose message identity

Under --verbose, changing request_id to a UUID removes the Slack channel and message timestamp from log lines. Their replacements, slack.channel and slack.ts, live on the span, but TruncatingFormatter prints only the log event's fields. Operators can no longer connect errors such as a failed history read to the originating message. A2A logs lose the task identity for the same reason.

Update the verbose formatter to include span fields, and test with that formatter rather than only the default formatter.

Note: If this suggestion doesn't match your team's coding style, reply to this and let me know. I'll remember it for next time!

let earlier = match prefetched {
Some(earlier) => earlier,
None => {
Expand All @@ -291,7 +306,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
Expand Down Expand Up @@ -344,8 +359,10 @@ impl SlackIngress {
&self,
inbound: &Inbound,
earlier: &[SlackMessage],
request_id: &str,
run_id: aura::RunId,
) -> Result<String, RunError> {
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),
Expand All @@ -359,13 +376,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()))?;

Expand Down
Loading
Loading