Repository navigation
Give a run one handle, its own identity, and pluggable hooks - #710
Conversation
|
48360c7 to
1f15c0a
Compare
1f15c0a to
f93cec8
Compare
f93cec8 to
89414eb
Compare
89414eb to
7515211
Compare
7515211 to
8880383
Compare
| /// Releases the token when the call ends, including a call whose future is | ||
| /// dropped mid-await. | ||
| /// | ||
| /// A server that trails a notification past its call's result then finds no |
There was a problem hiding this comment.
nit - non-blocking? (see impact below)
This comment says a notification that arrives after its call's result "finds no owner and is counted unowned" but I don't think that is the case if happening mid-run.
Breakdown/step-through:
- builder.rs:1689 — stream() binds every client: bound_call = Some(req_1).
- call_tool_tracked (client.rs:910-918) claims token T and creates this guard.
- The call returns and the guard drops, so token_owners.remove(T). bound_call is still Some(req_1); it's only cleared in cancel_and_close (client.rs:1059).
- The server sends one more notifications/progress for T. owner_of (progress.rs:77-88) misses in token_owners and falls back to bound_call, which returns Some(req_1).
- on_progress routes it as owned. The client gets aura.progress for T after tool_complete, and the orphan counter doesn't move.
It stays inside the same run, so there's no cross-request leak, and it matches what nightly did before this change (binding-only routing). But the comment and the orphan diagnostics suggest
released tokens are dropped, and a_released_token_stops_routing only passes because it builds the handler with bound_call = None.
The fallback is there for a server that answers before the claim lands, so it can't just be removed.
Potential Options:
- Keep a small set of released tokens (or a watermark) and have owner_of return None for those instead of falling back to the binding.
- Or accept the behavior and change this comment (and the orphan branch's comment) to say late notifications are dropped only after the run unbinds. Then add a test with bound_call =
Some(..) that pins it down either way.
The impact here is pretty low so if we would rather create an issue and do a follow-up, that seems reasonable to me. I wouldn't block this PR just because of this.
Here's the impact:
- Stale progress in clients. A client using AURA_CUSTOM_EVENTS=true can get an aura.progress event after that tool's aura.tool_complete.
- In the aura CLI (repl/loop.rs:2901-2911), the late event looks for an active tool with the same progress_token. The tool has finished, so it doesn't find one and falls back to
set_agent_reasoning(message). For example, a stale "processed 40/50 files" briefly replaces the spinner's "Thinking" line while the model generates its next turn. It's transient and
never enters scrollback, but it's wrong. - A third-party client that tracks tools by progress_token might reopen a finished tool, or ignore the event. That depends on how the client is written.
- In the aura CLI (repl/loop.rs:2901-2911), the late event looks for an active tool with the same progress_token. The tool has finished, so it doesn't find one and falls back to
- Blind orphan diagnostics mid-run. The unowned counter and the ORPHANED_PROGRESS_ALARM warning ("server may be ignoring notifications/cancelled") only fire after the binding is cleared. A
server that keeps sending progress for finished calls during a long run won't be flagged until the run ends. - A misleading comment and test. The guard comment says late notifications count as unowned, and a_released_token_stops_routing looks like it proves that because it runs with bound_call =
None. Someone relying on either would misjudge how routing actually works.
justintime4tea
left a comment
There was a problem hiding this comment.
Greptile flagged a cancel race on #730, which stacks on this branch. I traced it and it's real, in this branch, at the executor boundary, and it's new here rather than inherited from nightly. Details are inline on agent_executor.rs at L588, L371, L96, L1349, and L628. The fix is small and keeps your shutdown-after-finish fix intact.
8880383 to
20beeff
Compare
Comments Outside DiffThese findings could not be posted inline.
|
Starting an agent went through two methods whose implementations disagreed about cancellation. One remains, returning a handle that owns a run's events, cancellation and usage. RunOptions carries the bound and, for a caller that must cancel before the run exists, a token the run takes a child of. An unbounded run has no bound rather than the Duration::MAX the orchestrator passed for one. Dropping the handle, or the stream it hands out, cancels the run. An orchestration run spawns before its caller sees it, so a consumer that goes away would otherwise leave it spending provider turns. A2A registers a task's cancel entry before the agent build awaits, so a cancelTask arriving during it has something to cancel. make test-loom models that map's interleavings. Ref: #625 Signed-off-by: Jacob Hull <jacob@planethull.com>
Run identity lived on MCP clients as mutable state, set before every stream. It now lives in a scope the agent establishes, which anything running inside the run reads directly. Two things sit outside that scope. Rig executes tools on its own server task, so a client binds the call it serves and a tool call finds it there. Progress notifications arrive on the transport task, correlated by their call's progress token. Both hold the request id with the agent that call acts for, so events keep their attribution. rmcp mints a progress token inside the send, so a server can answer before the call has claimed it. The call the client is bound to covers that window, and a notification with no live call at all stays unowned. Warning waits for a count of them, because firing on the first named a cause it could not know. Ref: #625 Signed-off-by: Jacob Hull <jacob@planethull.com>
Everything watching a run had to live inside the single hook rig allows per streaming request. An AgentHook trait now carries those concerns independently, and the rig-facing hook fans out to them. Cancellation and client-tool passthrough move behind it. A hook that ends a run says why, so the log can tell a blown deadline from a disconnect without the hook holding the deadline itself. Tool events and usage stay together, sharing the pending tool ids they both need. with_hook registers another, so adding a concern means writing an implementation instead of extending one function. Ref: #625 Signed-off-by: Jacob Hull <jacob@planethull.com>
20beeff to
473360d
Compare
feat(agents): give StreamingAgent one entry point returning a run handle
Starting an agent went through two methods whose implementations disagreed about cancellation. One remains, returning a handle that owns a run's events, cancellation and usage.
RunOptions carries the bound and, for a caller that must cancel before the run exists, a token the run takes a child of. An unbounded run has no bound rather than the Duration::MAX the orchestrator passed for one.
Dropping the handle, or the stream it hands out, cancels the run. An orchestration run spawns before its caller sees it, so a consumer that goes away would otherwise leave it spending provider turns.
A2A registers a task's cancel entry before the agent build awaits, so a cancelTask arriving during it has something to cancel. make test-loom models that map's interleavings.
Ref: #625
feat(agents): carry run identity in task-local scope
Run identity lived on MCP clients as mutable state, set before every stream. It now lives in a scope the agent establishes, which anything running inside the run reads directly.
Two things sit outside that scope. Rig executes tools on its own server task, so a client binds the call it serves and a tool call finds it there. Progress notifications arrive on the transport task, correlated
Run identity lived on MCP clients as mutable state, set before every stream. It now lives in a scope the agent establishes, which anything running inside the run reads directly.
Two things sit outside that scope. Rig executes tools on its own server task, so a client binds the call it serves and a tool call finds it there. Progress notifications arrive on the transport task, correlated by their call's progress token. Both hold the request id with the agent that call acts for, so events keep their attribution.
rmcp mints a progress token inside the send, so a server can answer before the call has claimed it. The call the client is bound to covers that window, and a notification with no live call at all stays unowned. Warning waits for a count of them, because firing on the first named a cause it could not know.
Ref: #625
feat(agents): make hooks an extension point
Everything watching a run had to live inside the single hook rig allows per streaming request. An AgentHook trait now carries those concerns independently, and the rig-facing hook fans out to them.
Cancellation and client-tool passthrough move behind it. A hook that ends a run says why, so the log can tell a blown deadline from a disconnect without the hook holding the deadline itself. Tool events and usage stay together, sharing the pending tool ids they both need. with_hook registers another, so adding a concern means writing an implementation instead of extending one function.
Ref: #625