From 97fe34f132ca0903e8d3adeeebeb12dc1e62edc0 Mon Sep 17 00:00:00 2001 From: Brian Love Date: Wed, 2 Sep 2026 17:01:23 -0700 Subject: [PATCH] feat(runtimes): stand-down guard so the mastra subagent injector retires MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The injector exists only because @ag-ui/mastra 1.1.2 drops the child stream (`case "tool-output": break`). Upstream is expected to grow its own sub-agent surface — @ag-ui/core 0.0.59 already ships the SUBAGENT_* schemas and the LangGraph integration already emits them behind `subagent_visibility`. Once the installed bridge emits SUBAGENT_* too, injecting alongside it puts duplicates on the wire: a duplicate SUBAGENT_STARTED for one subagentRunId is a hard AG-UI verifier error, and two distinct ids paint two cards for one delegation. The first bridge SUBAGENT_* now latches `bridgeEmitsSubagentEvents`, after which `chunk()` injects nothing and `eventsFor()` is a passthrough. Detection is sound because `chunk()` output is written straight to the socket and never re-enters `eventsFor()`, so a SUBAGENT_* seen there is always the bridge's. Two details that carry the design: - The latch is process-wide, not per-injector. `createSubagentInjector()` is called per request, so a per-run flag would re-learn every run and duplicate the first delegation of each one. - Standing down mid-delegation closes what we already announced, or RUN_FINISHED trips the verifier's "subagents are still active" rule. The close is neutral (SUBAGENT_FINISHED with no `outcome`, which the schema allows) because the delegation neither succeeded nor failed. `synthesized` is kept so the bridge's buffered TOOL_CALL_* copies still dedupe against the eager ones already on the wire. Known limit, documented in the module header: the tee sees a chunk before the bridge does, so the delegation in flight at detection already has our SUBAGENT_STARTED on the wire and cannot be retracted. Only the first delegation of the first run after an upgrade is affected. The clean exit remains deleting this module and the tee once upstream lands. Co-Authored-By: Claude Opus 5 --- deployments/ag-ui-mastra/server.mjs | 9 +- deployments/ag-ui-mastra/subagent-emitter.mjs | 71 +++++++++++- .../test/subagent-emitter.test.mjs | 107 +++++++++++++++++- 3 files changed, 183 insertions(+), 4 deletions(-) diff --git a/deployments/ag-ui-mastra/server.mjs b/deployments/ag-ui-mastra/server.mjs index 6641f90f5..fdfb6a919 100644 --- a/deployments/ag-ui-mastra/server.mjs +++ b/deployments/ag-ui-mastra/server.mjs @@ -22,9 +22,14 @@ import { pathToFileURL } from 'node:url'; import { dirname, resolve } from 'node:path'; import { MastraAgent } from '@ag-ui/mastra'; import { createMastra } from './agents.mjs'; -import { createSubagentInjector } from './subagent-emitter.mjs'; +import { createSubagentInjector, createBridgeCapability } from './subagent-emitter.mjs'; import { withDelegationTee } from './streaming-tee.mjs'; +// Process-wide, so the stand-down latch survives across runs: once the +// installed @ag-ui/mastra is seen emitting SUBAGENT_* itself, every later run +// starts already retired instead of re-learning (and double-emitting) each time. +const bridgeCapability = createBridgeCapability(); + const AG_UI_INTERNAL_TOKEN = process.env.AG_UI_INTERNAL_TOKEN; if (!AG_UI_INTERNAL_TOKEN) { // Same boot contract as ag-ui-dev's `os.environ["AG_UI_INTERNAL_TOKEN"]`: @@ -110,7 +115,7 @@ export function createAgUiServer() { // SUBAGENT_FINISHED/ERROR on `tool-result`. // - `eventsFor()`: the bridge's own AG-UI events, with its later buffered // TOOL_CALL_START/ARGS/END copies for a synthesized id dropped. - const injector = createSubagentInjector(); + const injector = createSubagentInjector(bridgeCapability); const write = (event) => res.write(sseFrame(event)); const observe = (chunk) => { for (const e of injector.chunk(chunk)) write(e); diff --git a/deployments/ag-ui-mastra/subagent-emitter.mjs b/deployments/ag-ui-mastra/subagent-emitter.mjs index 17c13993a..4dd14b8ca 100644 --- a/deployments/ag-ui-mastra/subagent-emitter.mjs +++ b/deployments/ag-ui-mastra/subagent-emitter.mjs @@ -42,9 +42,41 @@ // Terminal cleanup: RUN_ERROR / RUN_FINISHED with delegations still pending // closes open child messages and emits SUBAGENT_ERROR per pending id, // exactly once, before the terminal frame. +// +// STAND-DOWN (forward compatibility). This whole module exists only because +// @ag-ui/mastra 1.1.2 drops the child stream. Upstream is expected to grow its +// own sub-agent surface (@ag-ui/core 0.0.59 already ships the SUBAGENT_* +// schemas, and the LangGraph integration already emits them behind +// `subagent_visibility`). The day the installed bridge emits SUBAGENT_* itself, +// injecting too would put duplicates on the wire — and a duplicate +// SUBAGENT_STARTED for one subagentRunId is a HARD AG-UI verifier error, while +// two distinct ids paint two cards for one delegation. So the first bridge +// SUBAGENT_* latches `bridgeEmitsSubagentEvents` (see createBridgeCapability), +// after which `chunk()` injects nothing and `eventsFor()` is a passthrough. +// +// KNOWN LIMIT: the tee sees a chunk BEFORE the bridge does, so the delegation +// already in flight at the moment of detection has our SUBAGENT_STARTED on the +// wire already and cannot be retracted — it is closed neutrally instead. Only +// that first delegation of the first run after an upgrade is affected; the +// latch is process-wide, so every run after it starts retired. The clean exit +// is still to DELETE this module and the tee once upstream lands. const AGENT_TOOL_PREFIX = 'agent-'; +/** + * Shared, sticky record of whether the installed bridge emits SUBAGENT_* itself. + * + * Create ONE per process and pass it to every run's injector. The latch must + * outlive a single run: `createSubagentInjector` is called per request, so a + * per-injector flag would re-learn on every run and duplicate the first + * delegation of each one. + * + * @returns {{ bridgeEmitsSubagentEvents: boolean }} + */ +export function createBridgeCapability() { + return { bridgeEmitsSubagentEvents: false }; +} + /** * @typedef {object} Entry * @property {string} subagentRunId @@ -83,12 +115,15 @@ function failureMessage(parsed) { /** * Create a per-run injector. * + * @param {{ bridgeEmitsSubagentEvents: boolean }} [capability] the shared + * stand-down latch (see {@link createBridgeCapability}). Defaults to a + * private one, which is only right for a single-run test. * @returns {{ chunk(chunk: object): object[], eventsFor(event: object): object[] }} * `chunk` returns the AG-UI events to write for a raw Mastra chunk (usually * none); `eventsFor` returns, for each outbound bridge event, the ordered * list of frames to write (injections plus, unless deduped, the event). */ -export function createSubagentInjector() { +export function createSubagentInjector(capability = createBridgeCapability()) { /** @type {Map} pending delegations by toolCallId */ const pending = new Map(); /** Ids whose TOOL_CALL_START/ARGS/END were synthesized — bridge copies drop. */ @@ -153,8 +188,29 @@ export function createSubagentInjector() { return [...closeMessage(entry), { type: 'SUBAGENT_ERROR', subagentRunId: entry.subagentRunId, message }]; } + /** + * The bridge emitted SUBAGENT_* itself, so it owns the surface from here on + * and this injector retires (for this run and, via the shared latch, every + * later one). Anything we already announced has to be closed first: an open + * subagent at RUN_FINISHED is a hard verifier error. The close is NEUTRAL — + * SUBAGENT_FINISHED with no `outcome`, since the delegation neither succeeded + * nor failed, we simply stopped owning it. `synthesized` is deliberately kept + * so the bridge's buffered TOOL_CALL_* copies still dedupe against the eager + * ones already on the wire. + */ + function standDown() { + capability.bridgeEmitsSubagentEvents = true; + const out = [...pending.values()].flatMap((entry) => [ + ...closeMessage(entry), + { type: 'SUBAGENT_FINISHED', subagentRunId: entry.subagentRunId }, + ]); + pending.clear(); + return out; + } + return { chunk(chunk) { + if (capability.bridgeEmitsSubagentEvents) return []; const payload = chunk?.payload ?? {}; switch (chunk?.type) { case 'start': @@ -243,6 +299,19 @@ export function createSubagentInjector() { eventsFor(event) { switch (event.type) { + // Forward compatibility: the day the bridge emits these itself, retire + // rather than double-emit (a duplicate SUBAGENT_STARTED for one + // subagentRunId is a hard verifier error; two distinct ids paint two + // cards). Our own injections are written straight to the socket and + // never re-enter here, so a SUBAGENT_* at this point is always the + // bridge's. + case 'SUBAGENT_STARTED': + case 'SUBAGENT_FINISHED': + case 'SUBAGENT_ERROR': { + if (capability.bridgeEmitsSubagentEvents) return [event]; + return [...standDown(), event]; + } + case 'TOOL_CALL_START': { if (synthesized.has(event.toolCallId)) return []; // eager copy already on the wire const name = event.toolCallName ?? ''; diff --git a/deployments/ag-ui-mastra/test/subagent-emitter.test.mjs b/deployments/ag-ui-mastra/test/subagent-emitter.test.mjs index b883655aa..8ad1bb1f7 100644 --- a/deployments/ag-ui-mastra/test/subagent-emitter.test.mjs +++ b/deployments/ag-ui-mastra/test/subagent-emitter.test.mjs @@ -4,7 +4,7 @@ // cockpit/runtimes/mastra/angular/docs/wire-capture-subagents.md. import { test } from 'node:test'; import assert from 'node:assert/strict'; -import { createSubagentInjector } from '../subagent-emitter.mjs'; +import { createSubagentInjector, createBridgeCapability } from '../subagent-emitter.mjs'; const TID = 'call_aUfV9K0RCDZdZK3NWt9dRDKx'; @@ -464,3 +464,108 @@ test('chunk: agent-* tool-call with NO prior start/step-start omits parentMessag assert.equal(start.type, 'TOOL_CALL_START'); assert.equal('parentMessageId' in start, false, 'key must be absent, not present with an undefined value'); }); + +// --- stand-down guard ------------------------------------------------------ +// Forward compatibility: the day @ag-ui/mastra emits its own SUBAGENT_* events, +// this injector must retire rather than double-emit. A duplicate +// SUBAGENT_STARTED for one subagentRunId is a hard AG-UI verifier error, and +// two distinct ids paint two cards for one delegation. + +const bridgeStarted = (tid = TID) => ({ + type: 'SUBAGENT_STARTED', + subagentRunId: `bridge-${tid}`, + name: 'weather_forecaster', + parentToolCallId: tid, +}); + +test('stand-down: a bridge SUBAGENT_STARTED stops all further chunk injection', () => { + const injector = createSubagentInjector(); + injector.eventsFor(bridgeStarted()); + + assert.deepEqual(injector.chunk(stepStartChunk), []); + assert.deepEqual(injector.chunk(toolCallChunk()), [], 'must not synthesize TOOL_CALL_*/SUBAGENT_STARTED'); + assert.deepEqual(injector.chunk(childTextStart), []); + assert.deepEqual(injector.chunk(childDelta('hi')), []); + assert.deepEqual(injector.chunk(toolResultChunk(SUCCESS_RESULT)), []); +}); + +test('stand-down: the bridge SUBAGENT_* event itself passes through untouched', () => { + const injector = createSubagentInjector(); + const ev = bridgeStarted(); + assert.deepEqual(injector.eventsFor(ev), [ev]); +}); + +test('stand-down: SUBAGENT_FINISHED and SUBAGENT_ERROR also latch it', () => { + for (const ev of [ + { type: 'SUBAGENT_FINISHED', subagentRunId: 'b-1', outcome: { type: 'success' } }, + { type: 'SUBAGENT_ERROR', subagentRunId: 'b-1', message: 'nope' }, + ]) { + const injector = createSubagentInjector(); + assert.deepEqual(injector.eventsFor(ev), [ev]); + assert.deepEqual(injector.chunk(toolCallChunk()), [], `${ev.type} must latch stand-down`); + } +}); + +test('stand-down mid-delegation closes our in-flight child message and subagent first', () => { + const injector = createSubagentInjector(); + injector.chunk(stepStartChunk); + injector.chunk(toolCallChunk()); + injector.chunk(childTextStart); + injector.chunk(childDelta('partial')); + + const ev = bridgeStarted(); + // Our already-announced subagent must be closed, or RUN_FINISHED trips the + // verifier's "subagents are still active" rule. Neutral close: no outcome — + // the delegation did not actually succeed or fail, we simply stopped owning it. + assert.deepEqual(injector.eventsFor(ev), [ + { type: 'TEXT_MESSAGE_END', messageId: M1, subagentRunId: SUB }, + { type: 'SUBAGENT_FINISHED', subagentRunId: SUB }, + ev, + ]); +}); + +test('stand-down still drops bridge copies of TOOL_CALL_* we already synthesized', () => { + const injector = createSubagentInjector(); + injector.chunk(stepStartChunk); + injector.chunk(toolCallChunk()); + injector.eventsFor(bridgeStarted()); + + // Those three are already on the wire from the eager synthesis — letting the + // bridge's buffered copies through now would duplicate the tool call. + assert.deepEqual(injector.eventsFor(delegationStart()), []); + assert.deepEqual(injector.eventsFor({ type: 'TOOL_CALL_ARGS', toolCallId: TID, delta: '{}' }), []); + assert.deepEqual(injector.eventsFor({ type: 'TOOL_CALL_END', toolCallId: TID }), []); +}); + +test('stand-down: TOOL_CALL_RESULT and RUN_FINISHED pass through with no injection', () => { + const injector = createSubagentInjector(); + injector.chunk(stepStartChunk); + injector.chunk(toolCallChunk()); + injector.eventsFor(bridgeStarted()); + + const result = delegationResult(JSON.stringify(SUCCESS_RESULT)); + assert.deepEqual(injector.eventsFor(result), [result]); + const finished = { type: 'RUN_FINISHED', threadId: 't-1', runId: 'r-1' }; + assert.deepEqual(injector.eventsFor(finished), [finished], 'no cleanup — we no longer own any subagent'); +}); + +test('stand-down latches across runs through a shared capability', () => { + const capability = createBridgeCapability(); + const firstRun = createSubagentInjector(capability); + firstRun.eventsFor(bridgeStarted()); + + // A per-injector flag would re-learn every run and duplicate the first + // delegation of each one; the latch is shared so later runs start retired. + const secondRun = createSubagentInjector(capability); + secondRun.chunk(stepStartChunk); + assert.deepEqual(secondRun.chunk(toolCallChunk()), []); +}); + +test('injectors with independent capabilities do not affect each other', () => { + const stoodDown = createSubagentInjector(createBridgeCapability()); + stoodDown.eventsFor(bridgeStarted()); + + const fresh = createSubagentInjector(createBridgeCapability()); + fresh.chunk(stepStartChunk); + assert.deepEqual(fresh.chunk(toolCallChunk()), eagerToolCall); +});