From 7d129279a44e8a2ef3e2189688fa8224a4f6d34f Mon Sep 17 00:00:00 2001 From: Vito Cappello Date: Sun, 9 Aug 2026 21:19:22 -0400 Subject: [PATCH] fix(agents): prevent requester settle while child is still running (#120601) * fix(agents): keep requester settle attached to live children * fix(agents): gate requester settle on terminal children --------- Co-authored-by: VACInc <3279061+VACInc@users.noreply.github.com> --- ...ent-announce.requester-settle-wake.test.ts | 57 ++++++++++++ ...subagent-announce.requester-settle-wake.ts | 16 +++- ...agent-registry-lifecycle-requester-wake.ts | 5 + .../subagent-registry-lifecycle.test.ts | 20 ++++ ...registry.lifecycle-retry-grace.e2e.test.ts | 92 ++++++++++++++++--- 5 files changed, 175 insertions(+), 15 deletions(-) diff --git a/src/agents/subagent-announce.requester-settle-wake.test.ts b/src/agents/subagent-announce.requester-settle-wake.test.ts index b65c50d2a190..da008d52bd06 100644 --- a/src/agents/subagent-announce.requester-settle-wake.test.ts +++ b/src/agents/subagent-announce.requester-settle-wake.test.ts @@ -437,6 +437,63 @@ describe("maybeWakeRequesterAfterAllChildrenSettled", () => { expect(completeBatchSpy).toHaveBeenCalledWith(["run-b"], 1); }); + it.each([ + ["is running without an end timestamp", { status: "running", startedAt: 2_000 }], + [ + "is still marked running with an end timestamp", + { status: "running", startedAt: 2_000, endedAt: 3_000 }, + ], + ["has no end timestamp", { status: "terminal", startedAt: 2_000 }], + ] as const)( + "does not wake a yielded requester while its only frozen child %s", + async (_description, execution) => { + const activeChild = makeSettledChild({ + runId: "run-b", + execution, + delivery: { status: "pending" }, + requesterSettleWake: { + status: "pending", + attemptCount: 0, + batchRunIds: ["run-b"], + requesterYieldBatch: true, + rearmGeneration: 1, + }, + }); + registryRuntimeMock.listSubagentRunsForRequester.mockReturnValue([activeChild]); + + const woke = await maybeWakeRequesterAfterAllChildrenSettled( + wakeParams({ settledEntry: activeChild }), + ); + + expect(woke).toBe(false); + expect(deliverSpy).not.toHaveBeenCalled(); + expect(completeBatchSpy).not.toHaveBeenCalled(); + }, + ); + + it("wakes after a retired frozen member disappears from the registry", async () => { + const remainingChild = makeSettledChild({ + runId: "run-a", + delivery: { status: "delivered" }, + requesterSettleWake: { + status: "pending", + attemptCount: 0, + batchRunIds: ["run-a", "run-b"], + requesterYieldBatch: true, + rearmGeneration: 1, + }, + }); + registryRuntimeMock.listSubagentRunsForRequester.mockReturnValue([remainingChild]); + + const woke = await maybeWakeRequesterAfterAllChildrenSettled( + wakeParams({ settledEntry: remainingChild }), + ); + + expect(woke).toBe(true); + expect(deliverSpy).toHaveBeenCalledOnce(); + expect(completeBatchSpy).toHaveBeenCalledWith(["run-a"], 1); + }); + it("wakes after a requester yields with one already-delivered completion", async () => { const child = makeSettledChild({ runId: "run-b", diff --git a/src/agents/subagent-announce.requester-settle-wake.ts b/src/agents/subagent-announce.requester-settle-wake.ts index 7ed044fd3ea2..531a82303c4c 100644 --- a/src/agents/subagent-announce.requester-settle-wake.ts +++ b/src/agents/subagent-announce.requester-settle-wake.ts @@ -257,6 +257,8 @@ export async function maybeWakeRequesterAfterAllChildrenSettled(params: { let settledBatch: SubagentRunRecord[]; if (frozenBatchRunIds && frozenBatchRunIds.length > 0) { const runsById = new Map(requesterRuns.map((entry) => [entry.runId, entry])); + // Retired rows no longer own completion, but every surviving frozen member + // must be terminal before this batch can wake its requester. settledBatch = frozenBatchRunIds .map((runId) => runsById.get(runId)) .filter( @@ -264,9 +266,21 @@ export async function maybeWakeRequesterAfterAllChildrenSettled(params: { Boolean(entry?.requesterSettleWake) && entry?.requesterSettleWake?.rearmGeneration === currentRearmGeneration, ); + if ( + settledBatch.some( + (entry) => entry.execution.status === "running" || !hasSubagentRunEnded(entry), + ) + ) { + return false; + } } else { settledBatch = buildConnectedSettledWave( - requesterRuns.filter((entry) => entry.requesterSettleWake && hasSubagentRunEnded(entry)), + requesterRuns.filter( + (entry) => + entry.requesterSettleWake && + entry.execution.status !== "running" && + hasSubagentRunEnded(entry), + ), currentSettledEntry, ); } diff --git a/src/agents/subagent-registry-lifecycle-requester-wake.ts b/src/agents/subagent-registry-lifecycle-requester-wake.ts index baa007b7865b..f8dd88a4272d 100644 --- a/src/agents/subagent-registry-lifecycle-requester-wake.ts +++ b/src/agents/subagent-registry-lifecycle-requester-wake.ts @@ -6,6 +6,7 @@ import type { SubagentRegistryLifecycleState, } from "./subagent-registry-lifecycle-contracts.js"; import type { RequesterSettleWakeState, SubagentRunRecord } from "./subagent-registry.types.js"; +import { hasSubagentRunEnded } from "./subagent-run-liveness.js"; type RequesterSettleWakeBatchState = import("./subagent-announce.requester-settle-wake.js").RequesterSettleWakeBatchState; @@ -207,8 +208,12 @@ export function createSubagentRegistryLifecycleRequesterWake( function scheduleRequesterSettleWake(runId: string, entry: SubagentRunRecord): void { const requesterSessionKey = entry.requesterSessionKey?.trim(); + // A replayed lifecycle start can retain an older endedAt; require both + // terminal status and end evidence so a live child never wakes its requester. if ( entry.collect || + entry.execution.status === "running" || + !hasSubagentRunEnded(entry) || !requesterSessionKey || (entry.requesterTurnRunId && entry.requesterTurnYielded === true) || scheduledRequesterSettleWakeRuns.has(runId) || diff --git a/src/agents/subagent-registry-lifecycle.test.ts b/src/agents/subagent-registry-lifecycle.test.ts index b88da588b98e..7167a3af77b6 100644 --- a/src/agents/subagent-registry-lifecycle.test.ts +++ b/src/agents/subagent-registry-lifecycle.test.ts @@ -4121,6 +4121,26 @@ describe("requester settle wake trigger", () => { expect(later.requesterSettleWake).toEqual({ status: "pending", attemptCount: 0 }); }); + it("does not resume a persisted settle wake until its registry row is terminal", async () => { + const entry = createRunEntry({ + requesterSettleWake: { status: "pending", attemptCount: 0 }, + }); + const settleWake = vi.fn(async () => false); + const controller = createLifecycleController({ + entry, + maybeWakeRequesterAfterAllChildrenSettled: settleWake, + }); + + controller.resumeRequesterSettleWake(entry.runId, entry); + await Promise.resolve(); + expect(settleWake).not.toHaveBeenCalled(); + + entry.execution = { ...entry.execution, status: "terminal", endedAt: 4_000 }; + controller.resumeRequesterSettleWake(entry.runId, entry); + + await waitForLifecycleState(() => expect(settleWake).toHaveBeenCalledOnce()); + }); + it("keeps a yielded completion parked until its requester turn settles", async () => { const entry = createRunEntry({ requesterTurnRunId: "run-requester", diff --git a/src/agents/subagent-registry.lifecycle-retry-grace.e2e.test.ts b/src/agents/subagent-registry.lifecycle-retry-grace.e2e.test.ts index 16990a3c8c25..b4f36a2d9aff 100644 --- a/src/agents/subagent-registry.lifecycle-retry-grace.e2e.test.ts +++ b/src/agents/subagent-registry.lifecycle-retry-grace.e2e.test.ts @@ -4,7 +4,10 @@ import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { testing as subagentAnnounceDeliveryTesting } from "./subagent-announce-delivery.test-support.js"; import { testing as subagentAnnounceOutputTesting } from "./subagent-announce-output.test-support.js"; import { testing as subagentAnnounceTesting } from "./subagent-announce.js"; -import { testing as settleWakeTesting } from "./subagent-announce.requester-settle-wake.js"; +import { + maybeWakeRequesterAfterAllChildrenSettled, + testing as settleWakeTesting, +} from "./subagent-announce.requester-settle-wake.js"; import * as announceRead from "./subagent-registry-announce-read.js"; import * as mod from "./subagent-registry.test-helpers.js"; @@ -363,6 +366,17 @@ describe("subagent registry lifecycle error grace", () => { .filter((request): request is GatewayRequest => request.method === "agent"); } + function getRequesterWakeCalls() { + return getAgentCalls().filter((request) => { + const idempotencyKey = (request.params as Record | undefined) + ?.idempotencyKey; + return ( + typeof idempotencyKey === "string" && + idempotencyKey.startsWith("announce:requester-settle:") + ); + }); + } + function getAgentResultsForChildSession(childSessionKey: string): string[] { return getAgentCalls() .filter((request) => { @@ -479,25 +493,16 @@ describe("subagent registry lifecycle error grace", () => { rearmGeneration: undefined, }, ]); - const requesterWakeCalls = () => - getAgentCalls().filter((request) => { - const idempotencyKey = (request.params as Record | undefined) - ?.idempotencyKey; - return ( - typeof idempotencyKey === "string" && - idempotencyKey.startsWith("announce:requester-settle:") - ); - }); await waitForAgentCallCount(3); - expect(requesterWakeCalls()).toHaveLength(1); + expect(getRequesterWakeCalls()).toHaveLength(1); agentCallGates.delete(betaSessionKey); releaseBetaDelivery?.(); await waitForDeliveredCleanup("run-yield-alpha"); await waitForDeliveredCleanup("run-yield-beta"); - expect(requesterWakeCalls()).toHaveLength(1); - const requesterWakeParams = requesterWakeCalls()[0]?.params as + expect(getRequesterWakeCalls()).toHaveLength(1); + const requesterWakeParams = getRequesterWakeCalls()[0]?.params as | Record | undefined; expect(requesterWakeParams?.idempotencyKey).toContain(":yield-1"); @@ -511,7 +516,66 @@ describe("subagent registry lifecycle error grace", () => { await vi.advanceTimersByTimeAsync(30_000); await flushAsync(); - expect(requesterWakeCalls()).toHaveLength(1); + expect(getRequesterWakeCalls()).toHaveLength(1); + }); + + it("keeps a frozen live child asleep until its real registry row becomes terminal", async () => { + const requesterTurnRunId = "run-requester-live-child"; + const liveChildSessionKey = "agent:main:subagent:frozen-live-child"; + registerCompletionRun( + "run-frozen-live-child", + "frozen-live-child", + "live child", + requesterTurnRunId, + ); + setAssistantOutput(liveChildSessionKey, "live child complete"); + + expect( + mod.markRequesterTurnYielded({ + requesterSessionKey: MAIN_REQUESTER_SESSION_KEY, + requesterTurnRunId, + }), + ).toBe(1); + expect( + mod.settleRequesterAfterSessionSpawns({ + requesterSessionKey: MAIN_REQUESTER_SESSION_KEY, + requesterTurnRunId, + requesterYielded: true, + acceptedSessionSpawns: [ + { runId: "run-frozen-live-child", childSessionKey: liveChildSessionKey }, + ], + }), + ).toBe(true); + + const liveChild = mod + .listSubagentRunsForRequester(MAIN_REQUESTER_SESSION_KEY) + .find((run) => run.runId === "run-frozen-live-child"); + if (!liveChild) { + throw new Error("expected frozen live child"); + } + expect( + await maybeWakeRequesterAfterAllChildrenSettled({ + requesterSessionKey: MAIN_REQUESTER_SESSION_KEY, + settledEntry: liveChild, + transitionBatch: noop, + completeBatch: noop, + }), + ).toBe(false); + await flushAsync(); + + expect(getRequesterWakeCalls()).toHaveLength(0); + expect(liveChild.execution.status).toBe("running"); + + emitLifecycleEvent("run-frozen-live-child", { phase: "end", endedAt: Date.now() + 1 }); + await waitForAgentCallCount(1); + await waitForDeliveredCleanup("run-frozen-live-child"); + + expect( + mod + .listSubagentRunsForRequester(MAIN_REQUESTER_SESSION_KEY) + .find((run) => run.runId === "run-frozen-live-child")?.execution, + ).toMatchObject({ status: "terminal", endedAt: expect.any(Number) }); + expect(getRequesterWakeCalls()).toHaveLength(1); }); it("ignores transient lifecycle errors when run retries and then ends successfully", async () => {