Repository navigation
refactor(agents): name a run by its RunId - #796
justintime4tea wants to merge 6 commits into
Conversation
|
0346b07 to
a8058ca
Compare
a8058ca to
c5911a0
Compare
| ) -> crate::streaming::AgentRun { | ||
| // 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); |
There was a problem hiding this comment.
This takes the run id from the config, not from the caller.
- Built with
run_id: None, everystream()mints a fresh id and ignores therequest_idpassed in, apart from a debug log. The base version streamed underrequest_id. A caller that then sweeps approvals or cancels MCP calls under its own id matches nothing, andArc<dyn StreamingAgent>gives it no way to learn the real id. - Built with
Some(run_id), everystream()reuses it. There's noRunInProgresscheck likebegin_runhas, so two concurrent streams share an id, and one request's approval sweep cancels the other's.
The id minted here also isn't written back, so Orchestrator::new mints a second one (orchestrator.rs:610) and stores that as the run it serves.
The chat, A2A, and Slack paths all pass an id, so today this only affects library callers.
| Uuid::from_str(s).map(Self) | ||
| } | ||
| } | ||
| pub use aura_events::RunId; |
There was a problem hiding this comment.
Re-exporting aura_events::RunId here gives two different values the same type. Persistence still mints its own id per run (persistence.rs:332). That one goes into the park owner, AgentScope::Worker.run_id, and RunParked, while the span, the request, and the RunContext carry the other. The compiler can no longer tell them apart.
Whether a resumed run keeps its id (#780) decides whether these should be one type.
There was a problem hiding this comment.
Agreed. f4e2779 drops the re-export and puts persistence's id back in its own type, PersistenceRunId, so AgentScope::Worker, the orchestrator's scope building, and the stored-approval restore are typed by what they hold, and aura::RunId means only the run's. It refuses the nil UUID as RunId does. RunParked.run_id stays a String until #780 decides whether a resumed run keeps its id.
| crate::run_context::current_run().unwrap_or_else(|| { | ||
| RunContext::channel(self.agent_config.request_id.clone().unwrap_or_default()).0 | ||
| }) | ||
| crate::run_context::current_run().unwrap_or_else(|| RunContext::channel(self.run_id).0) |
There was a problem hiding this comment.
Outside a scope this builds a new RunContext on every call, so the coordinator and each worker share only the id. Each one gets its own cancel token, its own event channel with the receiver dropped, and its own tool-call queue, so cancelling the coordinator doesn't reach the workers.
Holding one Arc<RunContext>, for example in a OnceLock, would share the run itself. It's latent, since the factory always runs this inside with_run.
| let run_id = aura::RunId::mint(); | ||
| let span = tracing::info_span!( | ||
| parent: None, | ||
| "agent.stream", |
There was a problem hiding this comment.
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.
| 1 | ||
| }) | ||
| .boxed(); | ||
| assert_eq!(in_span(span, run).collect::<Vec<_>>().await, vec![1]); |
There was a problem hiding this comment.
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.
| /// `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>, |
There was a problem hiding this comment.
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.
| 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: |
There was a problem hiding this comment.
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.
c5911a0 to
3b49718
Compare
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_<hex>, a2a_<task_id>, and slack_<channel>_<ts>; 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
`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
`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
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
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
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
3b49718 to
f4e2779
Compare
| ) { | ||
| // In string form, the run's id is the request id everything | ||
| // request-keyed reads. | ||
| let request_id = run_id.to_string(); |
There was a problem hiding this comment.
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!
RunContext::idbecomes theRunId, so the envelope's id and the task-local id are one value.RunIdand use its string form as the request id, so approvals, their sweep, MCP cancellation, and the A2A cancel map keep one value per run.begin_run,AgentRuntimeConfig, andRigBuilder's build methods take aRunId; an agent built without one mints its own. The orchestration factory streams under the run id it was built for, asAgentalready did.PersistenceRunId, so the compiler keeps it apart from the run'sRunId; it refuses the nil UUID asRunIddoes. Adopting the run's id there waits for [FEATURE]: A runtime that owns sessions and their runs #780, which owns resume.Visible change: request ids are hyphenated UUIDs rather than
req_<hex>,a2a_<task_id>, andslack_<channel>_<ts>, so the id no longer says where a run came from. Each ingress opens its run'sagent.streamroot span withrun.idand the origin instead:a2a.task_idanda2a.context_idfor A2A,slack.channelandslack.tsfor Slack (chat's carriesrun.id, which its HTTP span already records ashttp.request_id). Log lines emitted inside the run, the agent's included, print those fields, and the trace carries them as attributes. A2A gains theagent.streamroot span chat and Slack already had; the Slack run id is minted when the message is queued, where its span opens.Verification
cargo +nightly fmt --check,cargo clippy --workspace --all-targets --all-features -- -D warnings, andcargo test --workspace(2,575 passed) on this layer.Fixes: GH-778
Ref: GH-578
Ref: GH-780