Skip to content

fix(agents): refuse a stream under another run's id or a second time - #797

Open
justintime4tea wants to merge 6 commits into
justingross/GH-778-run-id-on-run-contextfrom
justingross/GH-790-refuse-stream-misuse
Open

justintime4tea wants to merge 6 commits into
justingross/GH-778-run-id-on-run-contextfrom
justingross/GH-790-refuse-stream-misuse

Conversation

@justintime4tea

@justintime4tea justintime4tea commented Oct 9, 2026 •

Copy link
Copy Markdown
Collaborator

StreamingAgent::stream on a built Agent streamed under the id it began with whatever request_id it was given, and a second call reused the already-cancelled run. The run-id layer (#796) had given the orchestration factory the same mismatch. Both now claim the run's one stream through StreamClaim: a request_id that isn't the run's id (in the run's own spelling), or a second call, returns a run whose stream yields one StreamRefused and does nothing. The id is checked before the claim is taken. Only the trait method claims, so the coordinator's transient retry (stream_chat_with_depth) still re-streams. An OrchestratorFactory fixes its run id when built, minting one if the config names none. StreamingAgent::run_id() returns the run an agent streams, so a caller holding only the trait object a builder returns can stream an agent built with run_id: None; it replaces the inherent run_id on Agent and the factory. StreamClaim is public so a StreamingAgent outside aura keeps the same contract; the test mock does, and its callers stream under its run_id(). The trait doc states the contract and the doc examples build with the id they stream under; the server paths (chat, A2A, Slack) already did.

Verification

  • cargo +nightly fmt --check, cargo clippy --workspace --all-targets --all-features -- -D warnings, and cargo test --workspace (2,588 passed) on this layer.

Fixes: GH-790
Ref: GH-778

@greptile-apps

greptile-apps Bot commented Oct 9, 2026 •

Copy link
Copy Markdown

RetriggerConfidence Score: 5/5

[Critical impact] The PR appears safe to merge; no actionable new issue was found.

Summary

The PR makes each agent accept one stream under its own run id.

  • StreamClaim rejects a wrong id without taking the stream.
  • StreamingAgent::run_id() exposes the id through the builder’s returned trait object.
  • OrchestratorFactory fixes its id when built.
  • The mock now follows the same contract, addressing the previous unnumbered finding.

Diagram

%%{init: {'theme': 'neutral'}}%%
flowchart TD
  A["StreamingAgent::stream"] --> B{"Id matches?"}
  B -->|No| R["Return one StreamRefused"]
  B -->|Yes| C{"Stream already claimed?"}
  C -->|Yes| R
  C -->|No| D["Claim and start stream"]
  E["Coordinator retry"] --> F["stream_chat_with_depth"]
  F --> G["Stream without another claim"]
Loading

Reviews (6) · Last reviewed commit: "fix(agents): refuse a stream under anoth..." · Reviewed by Greptile

Comment thread crates/aura-test-utils/src/mock_agent.rs
@justintime4tea
justintime4tea added this pull request to stack #798 October 9, 2026 00:32
@justintime4tea
justintime4tea marked this pull request as ready for review October 9, 2026 00:45
@justintime4tea
justintime4tea requested a review from a team October 9, 2026 00:45
@justintime4tea
justintime4tea force-pushed the justingross/GH-790-refuse-stream-misuse branch from 262ab78 to b063bad Compare October 9, 2026 00:46
@justintime4tea
justintime4tea force-pushed the justingross/GH-790-refuse-stream-misuse branch from b063bad to c339c83 Compare October 9, 2026 05:55
@justintime4tea
justintime4tea force-pushed the justingross/GH-790-refuse-stream-misuse branch from c339c83 to a42a20e Compare October 9, 2026 14:05
@justintime4tea
justintime4tea force-pushed the justingross/GH-790-refuse-stream-misuse branch from a42a20e to c020de9 Compare October 9, 2026 14:26
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
A built Agent streamed under the id it began with, whatever request_id
it was handed, logging a mismatch at debug. A library caller streaming
under its own id therefore had approvals, MCP tracking and the hook
keyed under another one, and cancel_and_close_mcp with its own id
matched nothing. A second stream reused the run: the first AgentRun's
drop guard had cancelled its token, so the second stopped at its first
hook, with no observer and the first's leftover tool-call ids. The
orchestration factory had the same mismatch, streaming under the run id
it was built with.

Both now claim the run's one stream through StreamClaim. A request_id
that is not the run's id, spelled as the run spells it, or a second
call, returns a run whose stream yields one StreamRefused and does
nothing. The id is checked before the claim is taken, so a misnamed call
leaves the stream to the caller that names the run. Only the trait
method claims: the orchestrator's transient retry re-streams a
coordinator through the inherent stream_chat_with_depth, which is
untouched. An OrchestratorFactory fixes its run id when built, minting
one when the config names none, and the orchestration it spawns carries
the same id.

StreamingAgent::run_id reads the run an agent streams, so a caller
holding only the trait object, as every builder returns, can stream an
agent built without an id. It replaces the inherent run_id on Agent and
OrchestratorFactory, and a refusal for another run points at it.

StreamClaim is public, so a StreamingAgent outside aura keeps the same
contract. The test mock does: a wrong id or a second stream gets the
refusal a real agent gives, and its callers stream under its run_id().

The StreamingAgent::stream doc states the contract, and the doc
examples build with the id they stream under. The server paths already
do: chat completions, A2A and Slack each build for the run id they
stream under, and stream once.

Fixes: GH-790
Ref: GH-778
@justintime4tea
justintime4tea force-pushed the justingross/GH-790-refuse-stream-misuse branch from c020de9 to a417f91 Compare October 9, 2026 20:52
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant