diff --git a/docs/concepts/subagent-yield-handoff.md b/docs/concepts/subagent-yield-handoff.md index 68a1b7029b5e..4da549b83d9f 100644 --- a/docs/concepts/subagent-yield-handoff.md +++ b/docs/concepts/subagent-yield-handoff.md @@ -34,6 +34,14 @@ The implementation owners are `subagent-registry-requester-yield.ts`, `agent-task-tracking.ts`. `adoptPausedSubagentRunForFollowUp` uses the existing registry replacement operation; it does not create a second delegated task. +Private child results wait for their spawning turn to settle before individual +announcement admission. Normal settlement resumes each finished private child, +even while siblings are still running. Explicit yield assigns the frozen batch +first, then resumes child cleanup under that owner. Late announcement failures +cannot replace the batch's delivery state; already committed delivery evidence +remains valid. Restart activation reconciles retained requester-turn bindings +before resuming child completion. + Settlement dispatch uses `subagent_settle` input provenance. Individual announcements and the older descendant-wake path retain `subagent_announce`: the latter already owns its run replacement after dispatch and must not trigger diff --git a/docs/tools/subagents/announce.md b/docs/tools/subagents/announce.md index dbaa1d5f04ca..5c5666553510 100644 --- a/docs/tools/subagents/announce.md +++ b/docs/tools/subagents/announce.md @@ -48,7 +48,9 @@ This option supports hidden, native, one-shot runs only. It cannot be combined with ACP, `collect: true`, `visible: true`, `thread: true`, `mode: "session"`, or `expectsCompletionMessage: false`. It does not change the default completion mode. -Busy parents receive a separate private turn after their current work. A reset or +Finished private results remain in the registry until the spawning parent turn +settles. A normal parent finish releases each ready result for private review; +`sessions_yield` hands the results to its existing child batch instead. A reset or removed parent does not transfer the result to another session. When a settled batch contains a private result, its combined review stays private; ordinary siblings retain their individual completion delivery. 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 09770bb71a68..f2c4b7d28440 100644 --- a/src/agents/subagents/registry/subagent-registry-lifecycle-announce-cleanup.ts +++ b/src/agents/subagents/registry/subagent-registry-lifecycle-announce-cleanup.ts @@ -396,6 +396,16 @@ export const startSubagentAnnounceCleanupFlow = ( } const cleanup = entry.cleanup; const skipRequesterDelivery = entry.suppressCompletionDelivery === true; + // The spawning turn decides between individual review and a yielded batch. + // Keep private results durable without admitting a competing requester turn. + if ( + entry.completionTarget === "parent" && + entry.requesterTurnRunId && + !skipRequesterDelivery && + entry.delivery?.status !== "delivered" + ) { + return false; + } // A terminal delivery failure closes upward delivery, not live descendants. // Their completion callback re-enters this same cleanup path without a timer. if ( @@ -523,6 +533,10 @@ export const startSubagentAnnounceCleanupFlow = ( } const pendingPayload = loadPendingFinalDeliveryPayload(entry); const requesterOrigin = normalizeDeliveryContext(pendingPayload.requesterOrigin); + const requesterSettleGeneration = entry.requesterSettleWake?.rearmGeneration; + const requesterTookCompletion = () => + entry.requesterTurnYielded === true || + entry.requesterSettleWake?.rearmGeneration !== requesterSettleGeneration; let latestDeliveryError = getDeliveryLastError(entry); let committedDelivery: SubagentRunRecord["delivery"]; const finalizeAnnounceCleanup = async (announceOutcome: SubagentAnnounceFlowOutcome) => { @@ -540,7 +554,8 @@ export const startSubagentAnnounceCleanupFlow = ( } // Requester-settle can commit delivery while the mirror lookup is pending. const shouldCreditPriorDelivery = entry.delivery?.status === "delivered" || hasDeliveryMirror; - if (shouldCreditPriorDelivery) { + const handedOff = requesterTookCompletion(); + if (shouldCreditPriorDelivery || handedOff) { latestDeliveryError = undefined; } if (announceOutcome !== "delivered" && latestDeliveryError) { @@ -550,7 +565,11 @@ export const startSubagentAnnounceCleanupFlow = ( context, runId, cleanup, - shouldCreditPriorDelivery ? "delivered" : announceOutcome, + shouldCreditPriorDelivery + ? "delivered" + : handedOff + ? "intentional_non_delivery" + : announceOutcome, cleanupGeneration, ); }; @@ -631,6 +650,11 @@ export const startSubagentAnnounceCleanupFlow = ( if (entry.delivery?.status === "delivered") { return; } + // A late failure cannot rearm an announcement transferred to the batch. + // Committed sends still retain their delivery evidence. + if (!delivery.delivered && requesterTookCompletion()) { + return; + } recordAnnounceDeliveryResult(entry, delivery, params.runs); if (delivery.delivered) { const deliveryState = ensureDeliveryState(entry); diff --git a/src/agents/subagents/registry/subagent-registry-lifecycle.test.ts b/src/agents/subagents/registry/subagent-registry-lifecycle.test.ts index e6db3738e3e0..ea8bc3b0f95f 100644 --- a/src/agents/subagents/registry/subagent-registry-lifecycle.test.ts +++ b/src/agents/subagents/registry/subagent-registry-lifecycle.test.ts @@ -4672,6 +4672,136 @@ describe("subagent registry lifecycle hardening", () => { expect(persist).toHaveBeenCalled(); }); + it.each([false, true])( + "holds private completion until requester settlement (yielded: %s)", + async (requesterYielded) => { + const entry = createRunEntry({ + endedAt: Date.now(), + outcome: { status: "ok" }, + requesterTurnRunId: "run-requester", + completionTarget: "parent", + expectsCompletionMessage: true, + retainAttachmentsOnKeep: true, + completion: { required: true, resultText: "private child result" }, + delivery: { status: "pending" }, + }); + const sibling = createRunEntry({ + runId: "slow-sibling", + childSessionKey: "agent:main:subagent:slow-sibling", + requesterSessionKey: entry.requesterSessionKey, + requesterTurnRunId: "run-requester", + expectsCompletionMessage: true, + }); + const runSubagentAnnounceFlow = vi.fn( + async (params) => + params.isCompletionOwnedByRequesterYield?.() ? "intentional_non_delivery" : "delivered", + ); + const runs = new Map([ + [entry.runId, entry], + [sibling.runId, sibling], + ]); + const controller = createLifecycleController({ + entry, + runs, + runSubagentAnnounceFlow, + resumeSubagentRun: (runId) => { + controller.startSubagentAnnounceCleanupFlow(runId, runs.get(runId)!); + }, + maybeWakeRequesterAfterAllChildrenSettled: async () => false, + }); + try { + expect(controller.startSubagentAnnounceCleanupFlow(entry.runId, entry)).toBe(false); + expect(runSubagentAnnounceFlow).not.toHaveBeenCalled(); + expect(entry.cleanupHandled).not.toBe(true); + expect(entry.completion?.resultText).toBe("private child result"); + if (requesterYielded) { + markRequesterTurnYieldedInRuns({ + requesterSessionKey: entry.requesterSessionKey, + requesterTurnRunId: "run-requester", + runs, + persistOrThrow: () => undefined, + }); + } + expect( + controller.settleRequesterTurnAfterSessionSpawns({ + requesterSessionKey: entry.requesterSessionKey, + requesterTurnRunId: "run-requester", + requesterYielded, + acceptedSessionSpawns: [entry, sibling].map((child) => ({ + runId: child.runId, + childSessionKey: child.childSessionKey, + expectsCompletionMessage: true, + })), + }), + ).toBe(true); + await waitForLifecycleState(() => expect(entry.cleanupCompletedAt).toBeTypeOf("number")); + expect(entry.requesterTurnRunId).toBeUndefined(); + expect(entry.delivery?.status).toBe(requesterYielded ? "pending" : "delivered"); + expect(entry.requesterSettleWake?.requesterYieldBatch).toBe( + requesterYielded ? true : undefined, + ); + expect(sibling.execution.endedAt).toBeUndefined(); + expect(runSubagentAnnounceFlow).toHaveBeenCalledOnce(); + } finally { + controller.clearScheduledResumeTimers(); + } + }, + ); + + it("does not let a late announce failure reclaim a pending yielded batch", async () => { + const entry = createRunEntry({ + endedAt: Date.now(), + outcome: { status: "ok" }, + requesterTurnRunId: "run-requester", + expectsCompletionMessage: true, + retainAttachmentsOnKeep: true, + completion: { required: true, resultText: "child result" }, + delivery: { status: "pending" }, + }); + const announce = createDeferredCore(); + const runSubagentAnnounceFlow = vi.fn( + async (params) => { + await announce.promise; + params.onDeliveryResult?.({ + delivered: false, + path: "direct", + error: "obsolete announce failure", + disposition: "retryable", + }); + return "retryable"; + }, + ); + const controller = createLifecycleController({ + entry, + runSubagentAnnounceFlow, + maybeWakeRequesterAfterAllChildrenSettled: async () => false, + }); + try { + controller.startSubagentAnnounceCleanupFlow(entry.runId, entry); + await waitForLifecycleState(() => expect(runSubagentAnnounceFlow).toHaveBeenCalledOnce()); + entry.requesterTurnYielded = true; + controller.settleRequesterTurnAfterSessionSpawns({ + requesterSessionKey: entry.requesterSessionKey, + requesterTurnRunId: "run-requester", + requesterYielded: true, + acceptedSessionSpawns: [{ runId: entry.runId, childSessionKey: entry.childSessionKey }], + }); + const batch = structuredClone(entry.requesterSettleWake); + announce.resolve(); + await waitForLifecycleState(() => expect(entry.cleanupCompletedAt).toBeTypeOf("number")); + expect(entry.requesterSettleWake).toEqual(batch); + expect(entry.delivery).toMatchObject({ + status: "pending", + disposition: "intentional_non_delivery", + }); + expect(entry.delivery?.lastError).toBeUndefined(); + expect(entry.delivery?.nextAttemptAt).toBeUndefined(); + } finally { + announce.resolve(); + controller.clearScheduledResumeTimers(); + } + }); + it.each([ { waitingOn: "announce", outcome: "intentional_non_delivery" }, { waitingOn: "announce", outcome: "retryable" }, diff --git a/src/agents/subagents/registry/subagent-registry-lifecycle.ts b/src/agents/subagents/registry/subagent-registry-lifecycle.ts index 151a05b74a0c..736fc69a0787 100644 --- a/src/agents/subagents/registry/subagent-registry-lifecycle.ts +++ b/src/agents/subagents/registry/subagent-registry-lifecycle.ts @@ -321,7 +321,14 @@ export class SubagentLifecycleController { ...args, runs: this.options.runs, persistOrThrow: (...runIds) => this.options.persistOrThrow(...runIds), - schedule: (runId, entry) => { + schedule: (runId, entry, kind) => { + if (kind === "completion") { + if (!this.hasCleanupFailure(entry)) { + this.options.resumedRuns.delete(runId); + this.options.resumeSubagentRun(runId); + } + return; + } if (this.hasScheduledRequesterSettleWakeRun(entry)) { this.markRequesterSettleWakeRearm(entry); return; 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 16451f5b962e..219d2c2015a7 100644 --- a/src/agents/subagents/registry/subagent-registry-requester-yield.test.ts +++ b/src/agents/subagents/registry/subagent-registry-requester-yield.test.ts @@ -324,7 +324,7 @@ describe("settleRequesterTurnAfterSessionSpawns", () => { expect(persistOrThrow).toHaveBeenCalledTimes(expected ? 2 : 1); if (expected) { expect(entry.requesterSettleWake?.batchRunIds).toEqual([entry.runId]); - expect(schedule).toHaveBeenCalledExactlyOnceWith(entry.runId, entry); + expect(schedule).toHaveBeenCalledExactlyOnceWith(entry.runId, entry, "settle"); } else { expect(entry.requesterSettleWake).toBeUndefined(); expect(entry.requesterTurnRunId).toBe(REQUESTER_TURN); @@ -378,7 +378,7 @@ describe("settleRequesterTurnAfterSessionSpawns", () => { afterRequesterYield: true, }); expect(entry.delivery?.disposition).toBe("intentional_non_delivery"); - expect(schedule).toHaveBeenCalledExactlyOnceWith(entry.runId, entry); + expect(schedule).toHaveBeenCalledExactlyOnceWith(entry.runId, entry, "settle"); }); it("persists a mixed delivered and in-progress yielded batch before scheduling", () => { @@ -418,7 +418,7 @@ describe("settleRequesterTurnAfterSessionSpawns", () => { expect(beta.requesterTurnRunId).toBeUndefined(); expect(beta.delivery?.disposition).toBe("intentional_non_delivery"); expect(calls).toEqual(["persist", "schedule"]); - expect(schedule).toHaveBeenCalledExactlyOnceWith(alpha.runId, alpha); + expect(schedule).toHaveBeenCalledExactlyOnceWith(alpha.runId, alpha, "settle"); }); it.each([true, false])( @@ -465,10 +465,10 @@ describe("settleRequesterTurnAfterSessionSpawns", () => { batchRunIds: [completion.runId], afterRequesterYield: true, }); - expect(schedule).toHaveBeenCalledExactlyOnceWith(completion.runId, completion); + expect(schedule).toHaveBeenCalledExactlyOnceWith(completion.runId, completion, "settle"); } else { expect(completion.requesterSettleWake).toBeUndefined(); - expect(schedule).toHaveBeenCalledExactlyOnceWith(completion.runId, completion); + expect(schedule).toHaveBeenCalledExactlyOnceWith(completion.runId, completion, "settle"); } expect(inline.requesterTurnRunId).toBe(REQUESTER_TURN); expect(inline.requesterTurnYielded).toBeUndefined(); diff --git a/src/agents/subagents/registry/subagent-registry-requester-yield.ts b/src/agents/subagents/registry/subagent-registry-requester-yield.ts index 9f35ace627a2..87c5fd9a9177 100644 --- a/src/agents/subagents/registry/subagent-registry-requester-yield.ts +++ b/src/agents/subagents/registry/subagent-registry-requester-yield.ts @@ -63,7 +63,7 @@ export function settleRequesterTurnAfterSessionSpawns(params: { acceptedSessionSpawns: readonly AcceptedSessionSpawn[]; runs: Map; persistOrThrow(...runIds: string[]): void; - schedule(runId: string, entry: SubagentRunRecord): void; + schedule(runId: string, entry: SubagentRunRecord, kind: "completion" | "settle"): void; }): boolean { const requesterSessionKey = params.requesterSessionKey.trim(); const requesterTurnRunId = params.requesterTurnRunId.trim(); @@ -220,12 +220,21 @@ export function settleRequesterTurnAfterSessionSpawns(params: { scheduleYieldedSubagentRunProgress(entry); } } + for (const entry of entries) { + if ( + entry.completionTarget === "parent" && + typeof entry.execution.endedAt === "number" && + params.runs.has(entry.runId) + ) { + params.schedule(entry.runId, entry, "completion"); + } + } if ( rearmGeneration !== undefined && entries.every((entry) => typeof entry.execution.endedAt === "number") ) { // Active children keep the frozen batch; their normal completion owner schedules it. - params.schedule(firstEntry.runId, firstEntry); + params.schedule(firstEntry.runId, firstEntry, "settle"); } else if ( !params.requesterYielded && entries.every((entry) => typeof entry.execution.endedAt === "number") @@ -234,7 +243,7 @@ export function settleRequesterTurnAfterSessionSpawns(params: { // Once a normal parent response settles, resume its original per-child delivery. for (const entry of entries) { if (params.runs.has(entry.runId)) { - params.schedule(entry.runId, entry); + params.schedule(entry.runId, entry, "settle"); } } } diff --git a/test/subagent-requester-owner.e2e.test.ts b/test/subagent-requester-owner.e2e.test.ts index cee6c90083fb..63b0490a9764 100644 --- a/test/subagent-requester-owner.e2e.test.ts +++ b/test/subagent-requester-owner.e2e.test.ts @@ -7,8 +7,11 @@ import { saveSubagentRegistryToSqlite, } from "../src/agents/subagents/registry/subagent-registry.store.sqlite.js"; import type { SubagentRunRecord } from "../src/agents/subagents/registry/subagent-registry.types.js"; +import { getSessionKysely } from "../src/config/sessions/session-accessor.sqlite-scope.js"; import type { OpenClawConfig } from "../src/config/types.openclaw.js"; import { connectGatewayClient, disconnectGatewayClient } from "../src/gateway/test-helpers.e2e.js"; +import { executeSqliteQuerySync } from "../src/infra/kysely-sync.js"; +import { withOpenClawAgentDatabaseReadOnly } from "../src/state/openclaw-agent-db-readonly.js"; import { closeOpenClawStateDatabaseForTest } from "../src/state/openclaw-state-db.js"; import { writeOpenAiResponsesSse, @@ -18,6 +21,7 @@ import { createOpenClawTestInstance, type OpenClawTestInstance, } from "./helpers/openclaw-test-instance.js"; +import { createDeferred } from "./helpers/promise.js"; const TEST_TIMEOUT_MS = 180_000; const MODEL_REF = "requester-owner/synthetic"; @@ -38,6 +42,7 @@ type ProofModelServer = { bodies: () => readonly string[]; close: () => Promise; countRequestsContaining: (marker: string) => number; + completionResponseCount: () => number; requestCount: () => number; url: string; }; @@ -57,6 +62,132 @@ afterEach(async () => { }); describe("REQUESTER-OWNER requester agent id survives completion dispatch", () => { + it( + "delivers a private result once when the child finishes before the parent yields", + { timeout: TEST_TIMEOUT_MS }, + async () => { + const yieldGate = createDeferred(); + const modelServer = await startProofModelServer({ yieldAfterSpawn: yieldGate.promise }); + modelServers.push(modelServer); + const instance = await createOpenClawTestInstance({ + name: "private-completion-before-yield", + config: createTestConfig(modelServer.url), + env: { OPENCLAW_SKIP_PROVIDERS: undefined, OPENCLAW_TEST_MINIMAL_GATEWAY: undefined }, + }); + instances.push(instance); + instance.state.applyEnv(); + const sessionId = "private-yield-requester-session"; + const sessionKey = `agent:${REQUESTER_AGENT_ID}:${REQUESTER_KEY}`; + await writeSubagentSessionEntry({ + stateDir: instance.stateDir, + agentId: REQUESTER_AGENT_ID, + sessionKey, + sessionId, + defaultSessionId: sessionId, + }); + closeOpenClawStateDatabaseForTest(); + await instance.startGateway(); + const chatErrors: unknown[] = []; + const client = await connectGatewayClient({ + url: instance.url, + token: instance.gatewayToken, + onEvent: (event) => { + if (event.event === "chat" && (event.payload as { state?: string })?.state === "error") { + chatErrors.push(event.payload); + } + }, + }); + const readInputs = () => + withOpenClawAgentDatabaseReadOnly( + ({ db }) => + executeSqliteQuerySync( + db, + getSessionKysely(db) + .selectFrom("session_pending_inputs") + .selectAll() + .where("session_id", "=", sessionId), + ).rows, + { agentId: REQUESTER_AGENT_ID }, + ); + try { + const parent = client.request( + "agent", + { + sessionKey, + agentId: REQUESTER_AGENT_ID, + idempotencyKey: "private-yield-parent-turn", + message: PARENT_PROMPT, + deliver: false, + }, + { expectFinal: true }, + ); + void parent.catch(() => {}); + instance.state.applyEnv(); + await vi.waitFor( + () => { + const runs = [...loadSubagentRegistryFromSqlite().values()]; + expect(runs, instance.logs()).toHaveLength(1); + expect(runs[0]).toMatchObject({ + completionTarget: "parent", + requesterTurnRunId: "private-yield-parent-turn", + execution: { status: "terminal", outcome: { status: "ok" } }, + completion: { resultText: CHILD_MARKER }, + }); + expect(modelServer.countRequestsContaining(CHILD_MARKER)).toBe(0); + }, + { interval: 50, timeout: 60_000 }, + ); + // The parent checks the finished child through its real tool before yielding. + yieldGate.resolve(); + expect(await parent, instance.logs()).toMatchObject({ status: "ok" }); + await vi.waitFor( + () => { + const runs = [...loadSubagentRegistryFromSqlite().values()]; + expect(runs, instance.logs()).toHaveLength(1); + expect(runs[0]?.delivery?.status, instance.logs()).toBe("delivered"); + expect(runs[0]?.requesterSettleWake).toBeUndefined(); + }, + { interval: 50, timeout: 30_000 }, + ); + const receipts = withOpenClawAgentDatabaseReadOnly( + ({ db }) => + executeSqliteQuerySync( + db, + getSessionKysely(db) + .selectFrom("session_input_completions") + .selectAll() + .where("session_id", "=", sessionId), + ).rows, + { agentId: REQUESTER_AGENT_ID }, + ); + expect(receipts.found).toBe(true); + if (!receipts.found) { + throw new Error("Expected durable private completion receipts"); + } + expect(receipts.value.filter((receipt) => receipt.succeeded === 1)).toHaveLength(1); + expect(modelServer.completionResponseCount()).toBe(1); + expect(chatErrors).toEqual([]); + const history = await client.request<{ messages: unknown[] }>("chat.history", { + sessionKey, + agentId: REQUESTER_AGENT_ID, + limit: 30, + }); + expect(JSON.stringify(history.messages)).not.toContain("This turn ended before a reply"); + const inputs = readInputs(); + expect(inputs.found && inputs.value.filter((input) => input.state !== "cancelled")).toEqual( + [], + ); + } finally { + yieldGate.resolve(); + await disconnectGatewayClient(client); + await instance.stopGateway(); + } + expect(instance.logs()).not.toContain( + "subagent source lifecycle changed before completion delivery", + ); + }, + ); + it( "preserves the requester owner through a fresh normalized spawn", { timeout: TEST_TIMEOUT_MS }, @@ -361,8 +492,13 @@ function buildToolCallEvents(name: string, args: Record): SseEv ]; } -async function startProofModelServer(): Promise { +async function startProofModelServer(options?: { + yieldAfterSpawn: Promise; +}): Promise { const requestBodies: string[] = []; + let parentCheckedChildren = false; + let parentYielded = false; + let completionResponses = 0; const server = createServer((request, response) => { void handleModelRequest(request, response).catch((error: unknown) => { if (!response.headersSent) { @@ -391,10 +527,16 @@ async function startProofModelServer(): Promise { body += typeof chunk === "string" ? chunk : Buffer.from(chunk).toString("utf8"); } requestBodies.push(body); + if (options?.yieldAfterSpawn && parentCheckedChildren && !parentYielded) { + parentYielded = true; + writeOpenAiResponsesSse(response, buildToolCallEvents("sessions_yield", {})); + return; + } const completion = [RESTORED_CHILD_RESULT, CHILD_MARKER].find((marker) => body.includes(marker), ); if (completion) { + completionResponses += 1; writeOpenAiResponsesText(response, { text: completion, responseId: `response-${++responseSequence}`, @@ -419,10 +561,17 @@ async function startProofModelServer(): Promise { label: "requester-owner-child", thread: false, mode: "run", + ...(options?.yieldAfterSpawn ? { completionTarget: "parent" } : {}), }), ); return; } + if (options?.yieldAfterSpawn) { + await options.yieldAfterSpawn; + parentCheckedChildren = true; + writeOpenAiResponsesSse(response, buildToolCallEvents("subagents", { action: "list" })); + return; + } writeOpenAiResponsesText(response, { text: "REQUESTER-OWNER-PARENT-OK", responseId: `response-${++responseSequence}`, @@ -439,6 +588,7 @@ async function startProofModelServer(): Promise { bodies: () => requestBodies, countRequestsContaining: (marker) => requestBodies.filter((entry) => entry.includes(marker)).length, + completionResponseCount: () => completionResponses, requestCount: () => requestBodies.length, url: `http://127.0.0.1:${address.port}`, close: async () => {