diff --git a/src/agents/subagents/completion/subagent-completion-admission.store.test.ts b/src/agents/subagents/completion/subagent-completion-admission.store.test.ts index a5886d16319e..47650cdb5a4c 100644 --- a/src/agents/subagents/completion/subagent-completion-admission.store.test.ts +++ b/src/agents/subagents/completion/subagent-completion-admission.store.test.ts @@ -1,7 +1,11 @@ import path from "node:path"; import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { useAutoCleanupTempDirTracker } from "../../../../test/helpers/temp-dir.js"; -import { prepareClaimedSessionDelivery } from "../../../infra/session-delivery-queue-storage.js"; +import { + prepareClaimedSessionDelivery, + SessionDeliveryDeadLetteredError, + SessionDeliveryDeferredError, +} from "../../../infra/session-delivery-queue-storage.js"; import { resolvePreferredOpenClawTmpDir } from "../../../infra/tmp-openclaw-dir.js"; import { closeOpenClawStateDatabaseForTest, @@ -210,6 +214,28 @@ describe("atomic subagent completion admission store", () => { expect(rowCount("task_runs")).toBe(0); }); + it("dead-letters expired orphan generations before resolving their logical owner", () => { + const { queueEntry } = records(); + if (queueEntry.kind !== "agentTurn" || queueEntry.owner?.kind !== "subagent_completion") { + throw new Error("expected correlated subagent completion queue entry"); + } + queueEntry.owner.deadlineAt = Date.now() - 1; + + expect(() => resolveCorrelatedSubagentDelivery(queueEntry)).toThrow( + SessionDeliveryDeadLetteredError, + ); + }); + + it("defers an unexpired generation whose logical owner has moved on", () => { + const { queueEntry, subagent } = records(); + subagent.delivery!.generation = 2; + subagentRuns.set(subagent.runId, subagent); + + expect(() => resolveCorrelatedSubagentDelivery(queueEntry)).toThrow( + SessionDeliveryDeferredError, + ); + }); + it("reloads a blocked text completion from SQLite before canonical owner redrive", async () => { await withEnvAsync({ OPENCLAW_STATE_DIR: tempDir }, async () => { closeOpenClawStateDatabaseForTest(); diff --git a/src/agents/subagents/completion/subagent-completion-delivery.ts b/src/agents/subagents/completion/subagent-completion-delivery.ts index 8450d1a82ce8..901e7019eff8 100644 --- a/src/agents/subagents/completion/subagent-completion-delivery.ts +++ b/src/agents/subagents/completion/subagent-completion-delivery.ts @@ -156,6 +156,11 @@ export function resolveCorrelatedSubagentDelivery( if (queued.kind !== "agentTurn" || queued.owner?.kind !== "subagent_completion") { return queued; } + if (Date.now() >= queued.owner.deadlineAt) { + throw new SessionDeliveryDeadLetteredError( + "correlated subagent completion delivery deadline expired", + ); + } const entry = subagentRuns.get(queued.owner.runId); if ( !entry || @@ -165,11 +170,6 @@ export function resolveCorrelatedSubagentDelivery( ) { throw new SessionDeliveryDeferredError("correlated subagent delivery owner mismatch"); } - if (Date.now() >= queued.owner.deadlineAt) { - throw new SessionDeliveryDeadLetteredError( - "correlated subagent completion delivery deadline expired", - ); - } return { ...queued, message: canonicalResultMessage(entry) }; } diff --git a/src/agents/subagents/registry/subagent-registry-lifecycle-common.ts b/src/agents/subagents/registry/subagent-registry-lifecycle-common.ts index 6ce9826d99fe..43177eb39149 100644 --- a/src/agents/subagents/registry/subagent-registry-lifecycle-common.ts +++ b/src/agents/subagents/registry/subagent-registry-lifecycle-common.ts @@ -48,8 +48,8 @@ export function createSubagentRegistryLifecycleCommon( clearTimeout(timer); } scheduledResumeTimers.clear(); - for (const timer of scheduledRequesterSettleWakeTimers.values()) { - clearTimeout(timer); + for (const scheduled of scheduledRequesterSettleWakeTimers.values()) { + clearTimeout(scheduled.timer); } scheduledRequesterSettleWakeTimers.clear(); pendingRequesterSettleWakeRearms.clear(); diff --git a/src/agents/subagents/registry/subagent-registry-lifecycle-contracts.ts b/src/agents/subagents/registry/subagent-registry-lifecycle-contracts.ts index baa9a362e36f..cff7b1ebcd25 100644 --- a/src/agents/subagents/registry/subagent-registry-lifecycle-contracts.ts +++ b/src/agents/subagents/registry/subagent-registry-lifecycle-contracts.ts @@ -57,7 +57,14 @@ export type SubagentRegistryLifecycleState = { scheduledResumeTimers: Set>; pendingRequesterSettleWakeRearms: Set; scheduledRequesterSettleWakeRuns: Set; - scheduledRequesterSettleWakeTimers: Map>; + scheduledRequesterSettleWakeTimers: Map< + string, + { + timer: ReturnType; + deadline: number; + rearmGeneration?: number; + } + >; terminalCompletionLocks: Map>; terminalGenerations: WeakMap; cleanupGenerations: WeakMap; diff --git a/src/agents/subagents/registry/subagent-registry-lifecycle-requester-wake.ts b/src/agents/subagents/registry/subagent-registry-lifecycle-requester-wake.ts index 6abfc0dda93e..11c99a6b3fef 100644 --- a/src/agents/subagents/registry/subagent-registry-lifecycle-requester-wake.ts +++ b/src/agents/subagents/registry/subagent-registry-lifecycle-requester-wake.ts @@ -103,7 +103,7 @@ export function createSubagentRegistryLifecycleRequesterWake( for (const [runId, entry] of entries) { const retryTimer = scheduledRequesterSettleWakeTimers.get(runId); if (retryTimer) { - clearTimeout(retryTimer); + clearTimeout(retryTimer.timer); scheduledRequesterSettleWakeTimers.delete(runId); } if (entry.requesterSettleWake === undefined || !params.runs.has(runId)) { @@ -183,17 +183,40 @@ export function createSubagentRegistryLifecycleRequesterWake( // cleanup parent reserves the root synchronously, so restart or suspend // cannot reach quiescence between scheduling and the wake's gateway turn. // Failures are logged only. + function retainScheduledRequesterSettleWakeTimer( + runId: string, + deadline: number, + rearmGeneration?: number, + ): boolean { + const scheduled = scheduledRequesterSettleWakeTimers.get(runId); + if (!scheduled) { + return false; + } + const hasNewerGeneration = + rearmGeneration !== undefined && + (scheduled.rearmGeneration === undefined || rearmGeneration > scheduled.rearmGeneration); + if (!hasNewerGeneration && deadline >= scheduled.deadline) { + return true; + } + clearTimeout(scheduled.timer); + scheduledRequesterSettleWakeTimers.delete(runId); + return false; + } + function scheduleRequesterSettleWakeRetry(runId: string, entry: SubagentRunRecord): void { const nextAttemptAt = entry.requesterSettleWake?.nextAttemptAt; - if ( - nextAttemptAt === undefined || - nextAttemptAt <= Date.now() || - scheduledRequesterSettleWakeTimers.has(runId) - ) { + if (nextAttemptAt === undefined || nextAttemptAt <= Date.now()) { + return; + } + const rearmGeneration = entry.requesterSettleWake?.rearmGeneration; + if (retainScheduledRequesterSettleWakeTimer(runId, nextAttemptAt, rearmGeneration)) { return; } const timer = setTimeout( () => { + if (scheduledRequesterSettleWakeTimers.get(runId)?.timer !== timer) { + return; + } scheduledRequesterSettleWakeTimers.delete(runId); const current = params.runs.get(runId); if (current === entry && current.requesterSettleWake) { @@ -203,7 +226,11 @@ export function createSubagentRegistryLifecycleRequesterWake( Math.max(0, nextAttemptAt - Date.now()), ); timer.unref?.(); - scheduledRequesterSettleWakeTimers.set(runId, timer); + scheduledRequesterSettleWakeTimers.set(runId, { + timer, + deadline: nextAttemptAt, + rearmGeneration, + }); } function scheduleRequesterSettleWake(runId: string, entry: SubagentRunRecord): void { @@ -216,12 +243,23 @@ export function createSubagentRegistryLifecycleRequesterWake( !hasSubagentRunEnded(entry) || !requesterSessionKey || (entry.requesterTurnRunId && entry.requesterTurnYielded === true) || - scheduledRequesterSettleWakeRuns.has(runId) || - scheduledRequesterSettleWakeTimers.has(runId) + scheduledRequesterSettleWakeRuns.has(runId) ) { return; } - if ((entry.requesterSettleWake?.nextAttemptAt ?? 0) > Date.now()) { + const now = Date.now(); + const nextAttemptAt = entry.requesterSettleWake?.nextAttemptAt; + const deadline = nextAttemptAt !== undefined && nextAttemptAt > now ? nextAttemptAt : now; + if ( + retainScheduledRequesterSettleWakeTimer( + runId, + deadline, + entry.requesterSettleWake?.rearmGeneration, + ) + ) { + return; + } + if (nextAttemptAt !== undefined && nextAttemptAt > now) { scheduleRequesterSettleWakeRetry(runId, entry); return; } diff --git a/src/agents/subagents/registry/subagent-registry-lifecycle.test.ts b/src/agents/subagents/registry/subagent-registry-lifecycle.test.ts index 6c194141ac29..83d1c10ce3db 100644 --- a/src/agents/subagents/registry/subagent-registry-lifecycle.test.ts +++ b/src/agents/subagents/registry/subagent-registry-lifecycle.test.ts @@ -4560,6 +4560,9 @@ describe("requester settle wake trigger", () => { }); await vi.advanceTimersByTimeAsync(0); expect(settleWake).toHaveBeenCalledTimes(1); + controller.resumeRequesterSettleWake(entry.runId, entry); + controller.resumeRequesterSettleWake(entry.runId, entry); + expect(vi.getTimerCount()).toBe(1); await vi.advanceTimersByTimeAsync(29_999); expect(settleWake).toHaveBeenCalledTimes(1); @@ -4572,6 +4575,61 @@ describe("requester settle wake trigger", () => { } }); + it("lets a fresh yield wake preempt a stale retry timer", async () => { + const entry = createRunEntry({ + endedAt: 4_000, + expectsCompletionMessage: true, + delivery: { status: "delivered" }, + requesterSettleWake: { + status: "pending", + attemptCount: 1, + nextAttemptAt: 120_000, + rearmGeneration: 1, + }, + }); + const settleWake = vi.fn( + async ( + params: Parameters< + LifecycleControllerParams["maybeWakeRequesterAfterAllChildrenSettled"] + >[0], + ) => { + params.completeBatch([entry.runId], entry.requesterSettleWake?.rearmGeneration); + return true; + }, + ); + const controller = createLifecycleController({ + entry, + maybeWakeRequesterAfterAllChildrenSettled: settleWake, + }); + + vi.useFakeTimers(); + vi.setSystemTime(0); + try { + controller.resumeRequesterSettleWake(entry.runId, entry); + expect(vi.getTimerCount()).toBe(1); + + entry.requesterTurnRunId = "run-requester"; + entry.requesterTurnYielded = true; + expect( + controller.settleRequesterTurnAfterSessionSpawns({ + requesterSessionKey: entry.requesterSessionKey, + requesterTurnRunId: "run-requester", + requesterYielded: true, + acceptedSessionSpawns: [{ runId: entry.runId, childSessionKey: entry.childSessionKey }], + }), + ).toBe(true); + await vi.advanceTimersByTimeAsync(0); + + expect(settleWake).toHaveBeenCalledOnce(); + expect(vi.getTimerCount()).toBe(0); + await vi.advanceTimersByTimeAsync(120_000); + expect(settleWake).toHaveBeenCalledOnce(); + } finally { + controller.clearScheduledResumeTimers(); + vi.useRealTimers(); + } + }); + it("does not re-arm coalesced batch rows whose retry deadline already passed", async () => { const state = { status: "pending" as const,