fix(agents): stop orphan delivery polling and delayed yield wakes (#121782)

* fix(agents): dead-letter expired orphan completion deliveries

* fix(agents): let a fresh yield wake preempt a stale retry timer
This commit is contained in:
Peter Steinberger
2026-08-10 18:17:20 -07:00
committed by GitHub
parent dd2aedf08f
commit dfb4b3658e
6 changed files with 148 additions and 19 deletions
@@ -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();
@@ -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) };
}
@@ -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();
@@ -57,7 +57,14 @@ export type SubagentRegistryLifecycleState = {
scheduledResumeTimers: Set<ReturnType<typeof setTimeout>>;
pendingRequesterSettleWakeRearms: Set<string>;
scheduledRequesterSettleWakeRuns: Set<string>;
scheduledRequesterSettleWakeTimers: Map<string, ReturnType<typeof setTimeout>>;
scheduledRequesterSettleWakeTimers: Map<
string,
{
timer: ReturnType<typeof setTimeout>;
deadline: number;
rearmGeneration?: number;
}
>;
terminalCompletionLocks: Map<string, Promise<void>>;
terminalGenerations: WeakMap<SubagentRunRecord, number>;
cleanupGenerations: WeakMap<SubagentRunRecord, number>;
@@ -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;
}
@@ -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,