From 094873902baf86f3a99d661a7bf2cc33650cfb8f Mon Sep 17 00:00:00 2001 From: Vito Cappello Date: Thu, 20 Aug 2026 12:02:08 -0400 Subject: [PATCH] fix(agents): retire delivered requester finals (#123285) * fix(agents): retire delivered requester finals * fix(agents): bind requester final receipts before yielded settlement --------- Co-authored-by: VACInc <3279061+VACInc@users.noreply.github.com> Co-authored-by: Peter Steinberger --- .../subagent-announce-delivery.test.ts | 1 + .../subagent-announce-direct-delivery.ts | 67 ++++++----- .../announce/subagent-announce-dispatch.ts | 2 + ...ent-registry-lifecycle-announce-cleanup.ts | 2 +- .../subagent-registry-lifecycle-delivery.ts | 30 +++++ .../subagent-registry-requester-yield.test.ts | 105 ++++++++++++++++++ .../subagent-registry-requester-yield.ts | 22 +++- ...registry.lifecycle-retry-grace.e2e.test.ts | 40 +++++++ .../subagent-registry.store.sqlite.test.ts | 18 +++ .../registry/subagent-registry.types.ts | 2 + 10 files changed, 260 insertions(+), 29 deletions(-) diff --git a/src/agents/subagents/announce/subagent-announce-delivery.test.ts b/src/agents/subagents/announce/subagent-announce-delivery.test.ts index 1352aef53350..a1238dd077f1 100644 --- a/src/agents/subagents/announce/subagent-announce-delivery.test.ts +++ b/src/agents/subagents/announce/subagent-announce-delivery.test.ts @@ -1918,6 +1918,7 @@ describe("deliverSubagentAnnouncement completion delivery", () => { }); expectDeliveryPath(result, "direct"); + expect(result).toMatchObject({ requesterVisibleFinalDelivered: true }); expect(callGateway).not.toHaveBeenCalled(); expectInProcessAgentParams(dispatchGatewayMethodInProcess, { deliver: true, diff --git a/src/agents/subagents/announce/subagent-announce-direct-delivery.ts b/src/agents/subagents/announce/subagent-announce-direct-delivery.ts index 2f9e15b20a03..c3b4adfbf0ef 100644 --- a/src/agents/subagents/announce/subagent-announce-direct-delivery.ts +++ b/src/agents/subagents/announce/subagent-announce-direct-delivery.ts @@ -519,35 +519,43 @@ export async function sendSubagentAnnounceDirectly(params: { error: "completion agent did not use the message tool for message-tool-only delivery", }; } - const hasVisibleCompletionReply = Boolean( + const requesterVisibleFinalDelivered = Boolean( directAnnounceResult && - ((params.requireVisibleReply - ? hasMessagingToolDeliveryToSource(directAnnounceResult, deliveryTarget, { - requireFinalReply: true, - }) - : hasMessagingToolDelivery) || - (hasVisibleAgentPayload( - params.requireVisibleReply - ? { - payloads: Array.isArray(directAnnounceResult.payloads) - ? directAnnounceResult.payloads.filter((payload) => { - const flags = payload as Record; - return ( - flags?.isCommentary !== true && - flags?.isCompactionNotice !== true && - flags?.isFallbackNotice !== true && - flags?.isStatusNotice !== true && - flags?.visible !== false - ); - }) - : [], - } - : directAnnounceResult, - { ...completionPayloadVisibility, includeSilentReplyPayloads: false }, - ) && - (!params.requireVisibleReply || - directAnnounceResult.deliveryStatus?.status !== "suppressed"))), + (hasMessagingToolDeliveryToSource(directAnnounceResult, deliveryTarget, { + requireFinalReply: true, + }) || + (shouldDeliverAgentFinal && + !requiresMessageToolDelivery && + hasVisibleAgentPayload( + { + payloads: Array.isArray(directAnnounceResult.payloads) + ? directAnnounceResult.payloads.filter((payload) => { + const flags = payload as Record; + return ( + flags?.isCommentary !== true && + flags?.isCompactionNotice !== true && + flags?.isFallbackNotice !== true && + flags?.isStatusNotice !== true && + flags?.visible !== false + ); + }) + : [], + }, + { ...completionPayloadVisibility, includeSilentReplyPayloads: false }, + ) && + directAnnounceResult.deliveryStatus?.status !== "suppressed")), ); + const hasVisibleCompletionReply = + requesterVisibleFinalDelivered || + (!params.requireVisibleReply && + Boolean( + directAnnounceResult && + (hasMessagingToolDelivery || + hasVisibleAgentPayload(directAnnounceResult, { + ...completionPayloadVisibility, + includeSilentReplyPayloads: false, + })), + )); const acceptsIntentionalSilentCompletion = hasIntentionalSilentCompletionReply && !isSubagentCompletion; if ( @@ -583,6 +591,11 @@ export async function sendSubagentAnnounceDirectly(params: { return { delivered: true, path: "direct", + ...(params.expectsCompletionMessage && + !params.requesterIsSubagent && + requesterVisibleFinalDelivered + ? { requesterVisibleFinalDelivered: true } + : {}), }; } catch (err) { const permanent = isPermanentAnnounceDeliveryError(err); diff --git a/src/agents/subagents/announce/subagent-announce-dispatch.ts b/src/agents/subagents/announce/subagent-announce-dispatch.ts index 351f45d8722c..20137f4f867f 100644 --- a/src/agents/subagents/announce/subagent-announce-dispatch.ts +++ b/src/agents/subagents/announce/subagent-announce-dispatch.ts @@ -32,6 +32,8 @@ export type SubagentAnnounceDeliveryResult = { path: SubagentDeliveryPath; deliveredAt?: number; enqueuedAt?: number; + /** Direct completion that already sent the yielded requester's visible final. */ + requesterVisibleFinalDelivered?: true; reason?: SubagentAnnounceDeliveryFailureReason; error?: string; // Stops fallback delivery when ownership changed or another terminal result diff --git a/src/agents/subagents/registry/subagent-registry-lifecycle-announce-cleanup.ts b/src/agents/subagents/registry/subagent-registry-lifecycle-announce-cleanup.ts index dff2b6363d4d..d2212903e670 100644 --- a/src/agents/subagents/registry/subagent-registry-lifecycle-announce-cleanup.ts +++ b/src/agents/subagents/registry/subagent-registry-lifecycle-announce-cleanup.ts @@ -580,7 +580,7 @@ export const startSubagentAnnounceCleanupFlow = ( retireSupersededCleanupInBackground(context, runId, entry, cleanupGeneration); return; } - recordAnnounceDeliveryResult(entry, delivery); + recordAnnounceDeliveryResult(entry, delivery, params.runs); if (delivery.delivered) { const deliveryState = ensureDeliveryState(entry); deliveryState.status = "delivered"; diff --git a/src/agents/subagents/registry/subagent-registry-lifecycle-delivery.ts b/src/agents/subagents/registry/subagent-registry-lifecycle-delivery.ts index a22afd3c1f9a..14d895911de0 100644 --- a/src/agents/subagents/registry/subagent-registry-lifecycle-delivery.ts +++ b/src/agents/subagents/registry/subagent-registry-lifecycle-delivery.ts @@ -39,6 +39,7 @@ import type { } from "./subagent-registry-lifecycle-context.js"; import type { PendingFinalDeliveryPayload, SubagentRunRecord } from "./subagent-registry.types.js"; import { compareSubagentRunGeneration } from "./subagent-run-generation.js"; +import { hasSubagentRunEnded } from "./subagent-run-liveness.js"; const DELIVERY_MIRROR_HISTORY_MAX_CHARS = 128 * 1024; @@ -78,6 +79,7 @@ export const formatAnnounceDeliveryError = (delivery: SubagentAnnounceDeliveryRe export const recordAnnounceDeliveryResult = ( entry: SubagentRunRecord, delivery: SubagentAnnounceDeliveryResult, + runs?: ReadonlyMap, ) => { const deliveryState = ensureDeliveryState(entry); if (typeof delivery.enqueuedAt === "number") { @@ -88,6 +90,34 @@ export const recordAnnounceDeliveryResult = ( typeof delivery.deliveredAt === "number" ? delivery.deliveredAt : Date.now(); deliveryState.deliveredAt = deliveredAt; deliveryState.lastDropReason = undefined; + const requesterTurnRunId = entry.requesterTurnRunId?.trim(); + if ( + delivery.path === "direct" && + delivery.requesterVisibleFinalDelivered && + requesterTurnRunId + ) { + const siblings = [...(runs?.values() ?? [])].filter( + (sibling) => + sibling.requesterSessionKey === entry.requesterSessionKey && + sibling.requesterTurnRunId === requesterTurnRunId && + sibling.expectsCompletionMessage === true, + ); + if ( + siblings.some((sibling) => sibling === entry) && + siblings.every( + (sibling) => + sibling.execution.status === "terminal" && + hasSubagentRunEnded(sibling) && + (sibling === entry || sibling.delivery?.status === "delivered"), + ) + ) { + // Bind final evidence before yielding; direct delivery is fenced once a yield is frozen. + deliveryState.requesterVisibleFinal = { + requesterTurnRunId, + batchRunIds: siblings.map((sibling) => sibling.runId).toSorted(), + }; + } + } } deliveryState.disposition = delivery.disposition ?? (delivery.delivered ? "delivered" : "retryable"); diff --git a/src/agents/subagents/registry/subagent-registry-requester-yield.test.ts b/src/agents/subagents/registry/subagent-registry-requester-yield.test.ts index 838d245e631c..859b260763e8 100644 --- a/src/agents/subagents/registry/subagent-registry-requester-yield.test.ts +++ b/src/agents/subagents/registry/subagent-registry-requester-yield.test.ts @@ -79,6 +79,86 @@ describe("settleRequesterTurnAfterSessionSpawns", () => { expect(schedule).toHaveBeenCalledOnce(); }); + it("retires a completed yielded batch whose requester already produced its final", () => { + const entry = makeRun("run-child"); + entry.cleanupCompletedAt = 2_100; + entry.delivery = { + status: "delivered", + requesterVisibleFinal: { requesterTurnRunId: REQUESTER_TURN, batchRunIds: [entry.runId] }, + }; + const schedule = vi.fn(); + + expect( + settleRequesterTurnAfterSessionSpawns({ + requesterSessionKey: REQUESTER, + requesterTurnRunId: REQUESTER_TURN, + requesterYielded: true, + acceptedSessionSpawns: [accepted(entry)], + runs: new Map([[entry.runId, entry]]), + persistOrThrow: vi.fn(), + schedule, + }), + ).toBe(true); + expect(entry.requesterSettleWake).toBeUndefined(); + expect(entry.requesterTurnRunId).toBeUndefined(); + expect(entry.delivery?.requesterVisibleFinal).toBeUndefined(); + expect(schedule).not.toHaveBeenCalled(); + }); + + it.each([ + [ + "another requester turn", + (entry: SubagentRunRecord) => { + entry.delivery!.requesterVisibleFinal!.requesterTurnRunId = "run-other"; + }, + ], + [ + "changed child membership", + (entry: SubagentRunRecord) => { + entry.delivery!.requesterVisibleFinal!.batchRunIds.push("run-later"); + }, + ], + [ + "unfinished cleanup", + (entry: SubagentRunRecord) => { + entry.cleanupCompletedAt = undefined; + }, + ], + [ + "unfinished delivery", + (entry: SubagentRunRecord) => { + entry.delivery!.status = "in_progress"; + }, + ], + [ + "a replayed running child", + (entry: SubagentRunRecord) => { + entry.execution.status = "running"; + }, + ], + ] as const)("keeps requester settlement when the final receipt has %s", (_, invalidate) => { + const entry = makeRun("run-child"); + entry.cleanupCompletedAt = 2_100; + entry.delivery = { + status: "delivered", + requesterVisibleFinal: { requesterTurnRunId: REQUESTER_TURN, batchRunIds: [entry.runId] }, + }; + invalidate(entry); + + expect( + settleRequesterTurnAfterSessionSpawns({ + requesterSessionKey: REQUESTER, + requesterTurnRunId: REQUESTER_TURN, + requesterYielded: true, + acceptedSessionSpawns: [accepted(entry)], + runs: new Map([[entry.runId, entry]]), + persistOrThrow: vi.fn(), + schedule: vi.fn(), + }), + ).toBe(true); + expect(entry.requesterSettleWake?.requesterYieldBatch).toBe(true); + }); + it.each([ ["matches", "agent:main:subagent:worker", true], ["rejects", "agent:main:subagent:other", false], @@ -291,6 +371,31 @@ describe("settleRequesterTurnAfterSessionSpawns", () => { expect(entry.retireAfterRequesterTurn).toBeUndefined(); }); + it("retires a delete-mode row after its requester-owned final is already delivered", () => { + const entry = makeRun("run-delete"); + entry.cleanup = "delete"; + entry.cleanupCompletedAt = 2_100; + entry.retireAfterRequesterTurn = true; + entry.delivery = { + status: "delivered", + requesterVisibleFinal: { requesterTurnRunId: REQUESTER_TURN, batchRunIds: [entry.runId] }, + }; + const runs = new Map([[entry.runId, entry]]); + + expect( + settleRequesterTurnAfterSessionSpawns({ + requesterSessionKey: REQUESTER, + requesterTurnRunId: REQUESTER_TURN, + requesterYielded: true, + acceptedSessionSpawns: [accepted(entry)], + runs, + persistOrThrow: vi.fn(), + schedule: vi.fn(), + }), + ).toBe(true); + expect(runs.has(entry.runId)).toBe(false); + }); + it("retires a completed delete-mode row after a normal requester answer", () => { const entry = makeRun("run-delete", false); entry.retireAfterRequesterTurn = true; diff --git a/src/agents/subagents/registry/subagent-registry-requester-yield.ts b/src/agents/subagents/registry/subagent-registry-requester-yield.ts index 97d4a9b7c8c6..468f1df20c95 100644 --- a/src/agents/subagents/registry/subagent-registry-requester-yield.ts +++ b/src/agents/subagents/registry/subagent-registry-requester-yield.ts @@ -91,8 +91,25 @@ export function settleRequesterTurnAfterSessionSpawns(params: { requesterTurnYielded: entry.requesterTurnYielded, retireAfterRequesterTurn: entry.retireAfterRequesterTurn, })); + const requesterAlreadyDeliveredFinal = + params.requesterYielded && + entries.every( + (entry) => + entry.execution.status === "terminal" && + typeof entry.execution.endedAt === "number" && + entry.delivery?.status === "delivered" && + typeof entry.cleanupCompletedAt === "number", + ) && + entries.some((entry) => { + const receipt = entry.delivery?.requesterVisibleFinal; + return ( + receipt?.requesterTurnRunId === requesterTurnRunId && + receipt.batchRunIds.length === batchRunIds.length && + receipt.batchRunIds.every((runId, index) => runId === batchRunIds[index]) + ); + }); let rearmGeneration: number | undefined; - if (params.requesterYielded) { + if (params.requesterYielded && !requesterAlreadyDeliveredFinal) { rearmGeneration = Math.max(0, ...entries.map((entry) => entry.requesterSettleWake?.rearmGeneration ?? 0)) + 1; for (const entry of entries) { @@ -126,6 +143,9 @@ export function settleRequesterTurnAfterSessionSpawns(params: { } } else { for (const entry of entries) { + if (entry.delivery) { + delete entry.delivery.requesterVisibleFinal; + } entry.requesterTurnRunId = undefined; entry.requesterTurnYielded = undefined; if (entry.retireAfterRequesterTurn === true) { diff --git a/src/agents/subagents/registry/subagent-registry.lifecycle-retry-grace.e2e.test.ts b/src/agents/subagents/registry/subagent-registry.lifecycle-retry-grace.e2e.test.ts index 03fb1b863db8..55bfb02d94bb 100644 --- a/src/agents/subagents/registry/subagent-registry.lifecycle-retry-grace.e2e.test.ts +++ b/src/agents/subagents/registry/subagent-registry.lifecycle-retry-grace.e2e.test.ts @@ -387,6 +387,46 @@ describe("subagent registry lifecycle error grace", () => { }); } + it("does not replay a requester-owned final already delivered before its turn yields", async () => { + const requesterTurnRunId = "run-requester-already-delivered"; + const runId = "run-completed-before-yield"; + const childSessionKey = "agent:main:subagent:completed-before-yield"; + registerCompletionRun(runId, "completed-before-yield", "finish once", requesterTurnRunId); + setAssistantOutput(childSessionKey, "child complete"); + + emitLifecycleEvent(runId, { phase: "end", endedAt: Date.now() }); + await waitForDeliveredCleanup(runId); + + const completed = mod + .listSubagentRunsForRequester(MAIN_REQUESTER_SESSION_KEY) + .find((run) => run.runId === runId); + expect(completed?.delivery?.requesterVisibleFinal).toEqual({ + requesterTurnRunId, + batchRunIds: [runId], + }); + expect(getAgentCalls()).toHaveLength(1); + expect( + mod.markRequesterTurnYielded({ + requesterSessionKey: MAIN_REQUESTER_SESSION_KEY, + requesterTurnRunId, + }), + ).toBe(1); + expect( + mod.settleRequesterAfterSessionSpawns({ + requesterSessionKey: MAIN_REQUESTER_SESSION_KEY, + requesterTurnRunId, + requesterYielded: true, + acceptedSessionSpawns: [{ runId, childSessionKey }], + }), + ).toBe(true); + + await vi.advanceTimersByTimeAsync(30_000); + await flushAsync(); + expect(getAgentCalls()).toHaveLength(1); + expect(getRequesterWakeCalls()).toHaveLength(0); + expect(completed?.delivery?.requesterVisibleFinal).toBeUndefined(); + }); + it("lets requester settlement own a yielded batch after sibling deliveries race", async () => { const requesterTurnRunId = "run-requester-yield-race"; const alphaSessionKey = "agent:main:subagent:yield-alpha"; diff --git a/src/agents/subagents/registry/subagent-registry.store.sqlite.test.ts b/src/agents/subagents/registry/subagent-registry.store.sqlite.test.ts index 88fe352d2aab..727b31c43d8b 100644 --- a/src/agents/subagents/registry/subagent-registry.store.sqlite.test.ts +++ b/src/agents/subagents/registry/subagent-registry.store.sqlite.test.ts @@ -139,6 +139,24 @@ describe("subagent registry sqlite store", () => { }); }); + it("preserves requester-owned final receipts in the existing SQLite payload", async () => { + await withTempStateEnv(async () => { + const requesterVisibleFinal = { + requesterTurnRunId: "run-requester", + batchRunIds: ["run-one"], + }; + const run = createRun({ delivery: { status: "delivered", requesterVisibleFinal } }); + + saveSubagentRegistryToSqlite(new Map([[run.runId, run]])); + closeOpenClawStateDatabaseForTest(); + + expect(loadSubagentRegistryFromSqlite().get(run.runId)?.delivery).toMatchObject({ + status: "delivered", + requesterVisibleFinal, + }); + }); + }); + it.each([ { name: "visible", diff --git a/src/agents/subagents/registry/subagent-registry.types.ts b/src/agents/subagents/registry/subagent-registry.types.ts index df10f8e879cb..90a1954557eb 100644 --- a/src/agents/subagents/registry/subagent-registry.types.ts +++ b/src/agents/subagents/registry/subagent-registry.types.ts @@ -139,6 +139,8 @@ export type SubagentCompletionDeliveryState = { enqueuedAt?: number; deliveredAt?: number; announcedAt?: number; + /** Exact requester turn and completed child batch that already produced its visible final. */ + requesterVisibleFinal?: { requesterTurnRunId: string; batchRunIds: string[] }; lastAttemptAt?: number; attemptCount?: number; lastError?: string | null;