fix(agents): defer recovered followups until drain

This commit is contained in:
joshavant
2026-08-05 23:27:10 -05:00
committed by Josh Avant
parent 7bcdb5823e
commit 5cc2fa6249
7 changed files with 373 additions and 39 deletions
@@ -35,6 +35,7 @@ export type QaMockProviderDispatchResult = {
failure?: QaMockProviderFailure;
onResponseSent?: () => void;
previewPauseMs?: number;
responsePauseMs?: number;
};
export type StreamEvent =
@@ -1,6 +1,11 @@
import path from "node:path";
import { afterEach, beforeEach, describe, expect, it } from "vitest";
import { useAutoCleanupTempDirTracker } from "../../../test/helpers/temp-dir.js";
import {
createReplyOperation,
isReplyRunActiveForSessionId,
runAfterReplyOperationClear,
} from "../../auto-reply/reply/reply-run-registry.js";
import { testing as replyRunTesting } from "../../auto-reply/reply/reply-run-registry.test-support.js";
import { clearRuntimeConfigSnapshot, setRuntimeConfigSnapshot } from "../../config/io.js";
import { loadSessionEntry, upsertSessionEntry } from "../../config/sessions/session-accessor.js";
@@ -44,6 +49,53 @@ describe("force-clear terminal state persistence", () => {
replyRunTesting.resetReplyRunRegistry();
});
it("delays stale-owner followups until the old reply owner settles", async () => {
const sessionKey = "agent:main:reply-stuck-followup";
const sessionId = "session-reply-stuck-followup";
const operation = createReplyOperation({ sessionKey, sessionId, resetTriggered: false });
const handle = createRunHandle();
operation.attachBackend({
kind: "embedded",
cancel: handle.abort,
isStreaming: handle.isStreaming,
});
operation.setPhase("running");
setActiveEmbeddedRun(sessionId, handle, sessionKey);
const followupObservedActiveHandle: boolean[] = [];
runAfterReplyOperationClear(operation, () => {
followupObservedActiveHandle.push(isEmbeddedAgentRunHandleActive(sessionId));
});
const recovery = abortAndDrainEmbeddedAgentRun({
sessionId,
sessionKey,
reason: "stuck_recovery",
forceClear: true,
settleMs: 100,
});
expect(isReplyRunActiveForSessionId(sessionId)).toBe(false);
expect(followupObservedActiveHandle).toEqual([]);
clearActiveEmbeddedRun(sessionId, handle, sessionKey);
let recoverySettled = false;
void recovery.then(() => {
recoverySettled = true;
});
await Promise.resolve();
expect(recoverySettled).toBe(false);
expect(followupObservedActiveHandle).toEqual([]);
operation.complete();
await expect(recovery).resolves.toEqual({
aborted: true,
drained: true,
forceCleared: false,
});
await Promise.resolve();
expect(followupObservedActiveHandle).toEqual([false]);
});
it("persists killed status after a force-cleared run", async () => {
const sessionKey = "agent:main:main";
const sessionId = "session-1";
@@ -239,6 +239,7 @@ describe("embedded-agent runner run registry", () => {
sessionId: "session-reply-stuck",
resetTriggered: false,
});
cancel.mockImplementation(() => operation.complete());
operation.attachBackend({
kind: "embedded",
cancel,
+70 -31
View File
@@ -20,6 +20,7 @@ import {
resolveReplyRunPhaseForSessionId,
type ReplyOperation,
type ReplyOperationPhase,
waitForReplyOperationOwnerSettlement,
waitForReplyRunEndBySessionId,
} from "../../auto-reply/reply/reply-run-registry.js";
import { getRuntimeConfig } from "../../config/io.js";
@@ -875,47 +876,85 @@ export async function abortAndDrainEmbeddedAgentRun(params: {
reason?: string;
}): Promise<AbortAndDrainEmbeddedAgentRunResult> {
const settleMs = params.settleMs ?? 15_000;
const settleDeadline = Date.now() + settleMs;
const embeddedRunHandle = ACTIVE_EMBEDDED_RUNS.get(params.sessionId);
const replyOperation = resolveActiveReplyOperationForSessionId(params.sessionId);
let releaseStaleExpiryBarrier: (() => void) | undefined;
const staleExpiryBarrier =
params.reason === "stuck_recovery"
? new Promise<void>((resolve) => {
releaseStaleExpiryBarrier = resolve;
})
: undefined;
// Recovery is a staleness expiry: stamp run_stalled on the reply operation
// BEFORE any handle abort, or the run loop's abort handler re-enters
// abortByUser and misattributes the watchdog kill to the user.
const expiredReplyRun =
params.reason === "stuck_recovery" &&
expireStaleReplyRunBySessionId(params.sessionId, "stuck_recovery");
if (expiredReplyRun && !ACTIVE_EMBEDDED_RUNS.has(params.sessionId)) {
// Reply expiry aborts synchronously and clears registry ownership. Let the
// command lane observe that abort before recovery decides whether to reset it.
await new Promise<void>((resolve) => {
setImmediate(resolve);
expireStaleReplyRunBySessionId(params.sessionId, "stuck_recovery", {
afterClearBarrier: staleExpiryBarrier,
followupAdmissionBarrierTimeout: settleMs + 1_000,
});
const drained = await waitForEmbeddedAgentRunEnd(params.sessionId, settleMs);
return { aborted: true, drained, forceCleared: false };
}
const aborted = abortEmbeddedAgentRun(params.sessionId) || expiredReplyRun;
const drained = aborted ? await waitForEmbeddedAgentRunEnd(params.sessionId, settleMs) : false;
const persistenceSnapshot =
params.forceClear === true && params.sessionKey
? tryLoadForceClearSessionSnapshot(params.sessionKey)
: undefined;
const forceCleared =
params.forceClear === true && (!aborted || !drained)
? forceClearEmbeddedAgentRun(
params.sessionId,
embeddedRunHandle,
replyOperation,
params.sessionKey,
params.reason,
)
const waitForExpiredOwnerSettlement = async () => {
if (!expiredReplyRun || !replyOperation) {
return true;
}
const settled = await waitForReplyOperationOwnerSettlement(
replyOperation,
Math.max(100, settleDeadline - Date.now()),
);
if (!settled) {
diag.warn(
`stuck recovery: reply owner settlement timed out sessionId=${params.sessionId} settleMs=${settleMs}`,
);
}
return settled;
};
try {
if (expiredReplyRun && !ACTIVE_EMBEDDED_RUNS.has(params.sessionId)) {
// Reply expiry aborts synchronously and clears registry ownership. Let the
// command lane observe that abort before recovery decides whether to reset it.
await new Promise<void>((resolve) => {
setImmediate(resolve);
});
const embeddedDrained = await waitForEmbeddedAgentRunEnd(params.sessionId, settleMs);
const ownerSettled = await waitForExpiredOwnerSettlement();
const drained = embeddedDrained && ownerSettled;
return { aborted: true, drained, forceCleared: false };
}
const aborted = abortEmbeddedAgentRun(params.sessionId) || expiredReplyRun;
const embeddedDrained = aborted
? await waitForEmbeddedAgentRunEnd(params.sessionId, settleMs)
: false;
if (forceCleared && params.sessionKey && persistenceSnapshot) {
await persistForceClearedEmbeddedRunTerminalState({
...persistenceSnapshot,
sessionId: params.sessionId,
sessionKey: params.sessionKey,
});
const ownerSettled = await waitForExpiredOwnerSettlement();
const drained = embeddedDrained && ownerSettled;
const persistenceSnapshot =
params.forceClear === true && params.sessionKey
? tryLoadForceClearSessionSnapshot(params.sessionKey)
: undefined;
const forceCleared =
params.forceClear === true && (!aborted || !drained)
? forceClearEmbeddedAgentRun(
params.sessionId,
embeddedRunHandle,
replyOperation,
params.sessionKey,
params.reason,
)
: false;
if (forceCleared && params.sessionKey && persistenceSnapshot) {
await persistForceClearedEmbeddedRunTerminalState({
...persistenceSnapshot,
sessionId: params.sessionId,
sessionKey: params.sessionKey,
});
}
return { aborted, drained, forceCleared };
} finally {
// Queue drains registered on the stale owner must not start while its
// backend can still claim the same session and requeue the adopted turn.
releaseStaleExpiryBarrier?.();
}
return { aborted, drained, forceCleared };
}
type ForceClearSessionSnapshot = {
@@ -31,6 +31,7 @@ import {
runAfterReplyOperationClear,
resolveActiveReplyRunSessionId,
resolveReplyRunPhaseForSessionId,
waitForReplyOperationOwnerSettlement,
waitForReplyRunEndBySessionId,
} from "./reply-run-registry.js";
import { testing } from "./reply-run-registry.test-support.js";
@@ -332,6 +333,84 @@ describe("reply run registry", () => {
});
});
it("keeps owner settlement pending after stale expiry through its completion barrier", async () => {
const operation = createTestReplyOperation({ sessionId: "session-stale-owner" });
operation.setPhase("running");
expect(expireStaleReplyOperation(operation, "stuck_recovery")).toBe(true);
expect(replyRunRegistry.isActive("agent:main:main")).toBe(false);
const settlement = waitForReplyOperationOwnerSettlement(operation, 1_000);
let settled = false;
void settlement.then((value) => {
settled = value;
});
await Promise.resolve();
expect(settled).toBe(false);
let releaseCompletion: () => void = () => {};
const completionBarrier = new Promise<void>((resolve) => {
releaseCompletion = resolve;
});
operation.completeWithAfterClearBarrier(completionBarrier);
await Promise.resolve();
expect(settled).toBe(false);
releaseCompletion();
await expect(settlement).resolves.toBe(true);
});
it("installs stale recovery barrier before synchronous cancel completion", async () => {
const operation = createTestReplyOperation({ sessionId: "session-sync-cancel" });
operation.setPhase("running");
operation.attachBackend({
kind: "embedded",
cancel: () => operation.complete(),
isStreaming: () => true,
});
let releaseRecovery: () => void = () => {};
const recoveryBarrier = new Promise<void>((resolve) => {
releaseRecovery = resolve;
});
const afterClear = vi.fn();
runAfterReplyOperationClear(operation, afterClear);
expect(
expireStaleReplyOperation(operation, "stuck_recovery", {
afterClearBarrier: recoveryBarrier,
}),
).toBe(true);
expect(afterClear).not.toHaveBeenCalled();
releaseRecovery();
await vi.waitFor(() => {
expect(afterClear).toHaveBeenCalledWith("session-sync-cancel");
});
});
it("keeps late after-clear registration behind an active stale barrier", async () => {
const operation = createTestReplyOperation({ sessionId: "session-late-callback" });
operation.setPhase("running");
let releaseRecovery: () => void = () => {};
const recoveryBarrier = new Promise<void>((resolve) => {
releaseRecovery = resolve;
});
expect(
expireStaleReplyOperation(operation, "stuck_recovery", {
afterClearBarrier: recoveryBarrier,
}),
).toBe(true);
const afterClear = vi.fn();
runAfterReplyOperationClear(operation, afterClear);
expect(afterClear).not.toHaveBeenCalled();
releaseRecovery();
await vi.waitFor(() => {
expect(afterClear).toHaveBeenCalledWith("session-late-callback");
});
});
it("keeps later after-clear work behind earlier delivery barriers", async () => {
const first = createTestReplyOperation({
sessionId: "first-session",
+71 -5
View File
@@ -20,6 +20,7 @@ import { diagnosticLogger as diag } from "../../logging/diagnostic-runtime.js";
import type { MediaFact } from "../../media/media-facts.js";
import type { PromptImageOrderEntry } from "../../media/prompt-image-order.js";
import type { UserTurnTranscriptRecorder } from "../../sessions/user-turn-transcript.types.js";
import { createDeferred } from "../../shared/deferred.js";
import { resolveGlobalSingleton } from "../../shared/global-singleton.js";
import { resolveTimerTimeoutMs } from "../../shared/number-coercion.js";
import type {
@@ -209,6 +210,8 @@ export type ReplyOperation = {
* Dispatch uses this while a user-visible failure payload still needs delivery.
*/
retainFailureUntilComplete(): void;
/** Settles after the lifecycle owner's final delivery/persistence barrier. */
readonly ownerSettlement?: Promise<void>;
complete(): void;
/**
* Complete the operation, clear active-run state, then run follow-up work.
@@ -387,9 +390,13 @@ const afterClearCallbacksByOperation = new WeakMap<
ReplyOperation,
Set<(sessionId: string) => void>
>();
type ReplyOperationStaleExpiryOptions = {
afterClearBarrier?: PromiseLike<unknown>;
followupAdmissionBarrierTimeout?: number | ReplyFollowupAdmissionBarrierTimeoutPolicy;
};
const expireReplyOperationByOperation = new WeakMap<
ReplyOperation,
(reason: ReplyOperationStaleReason) => boolean
(reason: ReplyOperationStaleReason, options?: ReplyOperationStaleExpiryOptions) => boolean
>();
function getAttachedBackend(operation: ReplyOperation): ReplyBackendHandle | undefined {
@@ -444,6 +451,11 @@ export function runAfterReplyOperationClear(
afterClear: (sessionId: string) => void,
): void {
if (replyRunState.activeRunsByKey.get(operation.key) !== operation) {
const barrier = replyRunState.followupAdmissionBarriersByKey.get(operation.key);
if (barrier) {
void barrier.settled.then(() => afterClear(barrier.sessionId));
return;
}
afterClear(operation.sessionId);
return;
}
@@ -616,9 +628,19 @@ export function createReplyOperation(params: {
let staleExpiryReason: ReplyOperationStaleReason | undefined;
let result: ReplyOperationResult | null = null;
let stateCleared = false;
let clearBarrierSettlement: Promise<void> | undefined;
let retainFailureUntilComplete = false;
let terminalRecovery = false;
let acceptedSteeredInboundAudio = false;
const ownerSettlement = createDeferred();
let ownerSettled = false;
const settleOwner = () => {
if (ownerSettled) {
return;
}
ownerSettled = true;
ownerSettlement.resolve(undefined);
};
const startedAtMs = Date.now();
const lifecycleGeneration = getAgentEventLifecycleGeneration();
let lastActivityAtMs = startedAtMs;
@@ -679,6 +701,7 @@ export function createReplyOperation(params: {
void registeredBarrier.settled.then(() =>
flushReplyOperationAfterClear(operation, registeredBarrier.sessionId),
);
clearBarrierSettlement = registeredBarrier.settled;
};
const abortInternally = (reason?: unknown) => {
@@ -917,12 +940,14 @@ export function createReplyOperation(params: {
retainFailureUntilComplete() {
retainFailureUntilComplete = true;
},
ownerSettlement: ownerSettlement.promise,
complete() {
if (!result) {
setResult({ kind: "completed" });
phase = "completed";
}
clearState();
settleOwner();
},
completeThen(afterClear) {
runAfterReplyOperationClear(operation, afterClear);
@@ -933,7 +958,19 @@ export function createReplyOperation(params: {
setResult({ kind: "completed" });
phase = "completed";
}
const wasAlreadyCleared = stateCleared;
clearState(barrier, timeoutMs);
// This barrier owns dispatch delivery and terminal persistence. Stale
// expiry may have already cleared the slot, but recovery must still wait
// for that old owner's durable work before admitting a queued turn.
const completionSettlement = wasAlreadyCleared
? waitForReplyBarrierSettlement(barrier, timeoutMs)
: clearBarrierSettlement;
if (completionSettlement) {
void completionSettlement.then(settleOwner);
} else {
settleOwner();
}
},
fail(code, cause) {
abortFrozenOperations.add(operation);
@@ -987,7 +1024,7 @@ export function createReplyOperation(params: {
},
};
expireReplyOperationByOperation.set(operation, (reason) => {
expireReplyOperationByOperation.set(operation, (reason, options) => {
if (replyRunState.activeRunsByKey.get(currentSessionKey) !== operation) {
return false;
}
@@ -1003,6 +1040,10 @@ export function createReplyOperation(params: {
setResult({ kind: "failed", code: "run_stalled" });
phase = "failed";
}
// Install the recovery fence before backend cancellation. Cancel can
// synchronously re-enter complete(), which must not flush queued turns
// while the old lifecycle owner is still finalizing.
clearState(options?.afterClearBarrier, options?.followupAdmissionBarrierTimeout);
getAttachedBackend(operation)?.cancel("superseded");
abortInternally(createAbortError("Reply operation expired as stale"));
diag.warn(
@@ -1010,7 +1051,6 @@ export function createReplyOperation(params: {
result,
)} ageMs=${Date.now() - lastActivityAtMs} ranForMs=${Date.now() - startedAtMs}`,
);
clearState();
return true;
});
const finalizationLease = replyRunSettle.createReplyRunFinalizationLease({
@@ -1118,16 +1158,42 @@ export function createReplyOperation(params: {
export function expireStaleReplyOperation(
operation: ReplyOperation,
reason: ReplyOperationStaleReason,
options?: ReplyOperationStaleExpiryOptions,
): boolean {
return expireReplyOperationByOperation.get(operation)?.(reason) ?? false;
return expireReplyOperationByOperation.get(operation)?.(reason, options) ?? false;
}
/** Wait for the old lifecycle owner's terminal work after stale expiry clears its slot. */
export async function waitForReplyOperationOwnerSettlement(
operation: ReplyOperation,
timeoutMs: number,
): Promise<boolean> {
const settlement = operation.ownerSettlement;
if (!settlement) {
return true;
}
const resolvedTimeoutMs = resolveTimerTimeoutMs(timeoutMs, 100, 100);
let timer: NodeJS.Timeout | undefined;
const settled = await Promise.race([
settlement.then(() => true),
new Promise<boolean>((resolve) => {
timer = setTimeout(() => resolve(false), resolvedTimeoutMs);
timer.unref?.();
}),
]);
if (timer) {
clearTimeout(timer);
}
return settled;
}
export function expireStaleReplyRunBySessionId(
sessionId: string,
reason: ReplyOperationStaleReason,
options?: Parameters<typeof expireStaleReplyOperation>[2],
): boolean {
const operation = resolveReplyRunForCurrentSessionId(sessionId);
return operation ? expireStaleReplyOperation(operation, reason) : false;
return operation ? expireStaleReplyOperation(operation, reason, options) : false;
}
// lastActivityAtMs is refreshed by agent events only; timers and user-message
@@ -10,6 +10,8 @@ type StabilityEvent = {
reason?: unknown;
outcome?: unknown;
ageMs?: unknown;
queueDepth?: unknown;
source?: unknown;
};
type StabilitySnapshot = {
@@ -32,6 +34,20 @@ type GatewayChatHistory = {
messages?: GatewayChatMessage[];
};
type MockRequestSnapshot = {
cursor?: unknown;
prompt?: unknown;
outcome?: unknown;
errorCode?: unknown;
};
type ClassifiedMockRequest = {
cursor: unknown;
prompt: "recovery" | "queued" | "other" | "missing";
outcome: unknown;
errorCode: unknown;
};
const RECOVERY_PROMPT =
"Repeated request recovery Gateway QA check. Keep attempting without producing a reply.";
const QUEUED_PROMPT =
@@ -39,7 +55,7 @@ const QUEUED_PROMPT =
const QUEUED_REPLY_MARKER = "GATEWAY_REPEATED_REQUEST_QUEUED_OK";
const RECOVERY_REASON = "repeated_model_requests_without_progress";
const PRODUCTION_RECOVERY_BOUND_MS = 360_000;
const HISTORY_RETRY_TIMEOUT_MS = 10_000;
const HISTORY_RETRY_TIMEOUT_MS = 60_000;
const HISTORY_RETRY_INTERVAL_MS = 250;
let harness: Awaited<ReturnType<typeof startQaLiveLaneGateway>> | undefined;
@@ -187,6 +203,66 @@ async function waitForQueuedReply(
: new Error(message, { cause: lastRetryableError });
}
async function readClassifiedMockRequests(mockBaseUrl: string): Promise<ClassifiedMockRequest[]> {
return fetch(`${mockBaseUrl}/debug/requests`)
.then((response) => response.json() as Promise<MockRequestSnapshot[]>)
.then((records) =>
records.map(({ cursor, prompt, outcome, errorCode }) => ({
cursor,
prompt:
typeof prompt === "string"
? prompt.includes(QUEUED_PROMPT)
? "queued"
: prompt.includes(RECOVERY_PROMPT)
? "recovery"
: "other"
: "missing",
outcome,
errorCode,
})),
);
}
async function readFailureEvidence(params: {
gateway: Awaited<ReturnType<typeof startQaLiveLaneGateway>>["gateway"];
mockBaseUrl: string | undefined;
sinceSeq: number;
}): Promise<string> {
const events = (await readStability(params.gateway, params.sinceSeq)).events ?? [];
const stability = events
.filter(
(event) =>
typeof event.type === "string" &&
(event.type.startsWith("session.") ||
event.type === "message.queued" ||
event.type === "model.call.started"),
)
.map(({ type, action, reason, outcome, ageMs, queueDepth, source }) => ({
type,
action,
reason,
outcome,
ageMs,
queueDepth,
source,
}));
const requests = params.mockBaseUrl
? await readClassifiedMockRequests(params.mockBaseUrl).catch((error: unknown) => [
{ requestEvidenceError: String(error) },
])
: [];
const gatewayLogs = params.gateway
.logs()
.split("\n")
.filter((line) =>
/followup queue|reply run stale takeover|stuck session recovery|queue: active session/iu.test(
line,
),
)
.slice(-100);
return JSON.stringify({ stability, requests, gatewayLogs });
}
describe("Gateway repeated-request recovery", () => {
it(
"aborts the real stalled owner once and releases one queued followup",
@@ -271,7 +347,7 @@ describe("Gateway repeated-request recovery", () => {
]);
expect(
events.filter((event) => event.type === "model.call.started").length,
).toBeGreaterThanOrEqual(3);
).toBeGreaterThanOrEqual(4);
const activeTerminal = (await gateway.call(
"agent.wait",
@@ -287,8 +363,28 @@ describe("Gateway repeated-request recovery", () => {
)) as GatewayChatRun;
expect(queuedTerminal.status).toBe("ok");
const history = await waitForQueuedReply(gateway, sessionKey);
const history = await waitForQueuedReply(gateway, sessionKey).catch(
async (error: unknown) => {
const evidence = await readFailureEvidence({
gateway,
mockBaseUrl: harness?.mock?.baseUrl,
sinceSeq: baselineSeq,
});
throw new Error(`${String(error)}; evidence=${evidence}`, { cause: error });
},
);
expect(historyContainsQueuedReply(history)).toBe(true);
const mockBaseUrl = harness?.mock?.baseUrl;
if (!mockBaseUrl) {
throw new Error("mock provider request evidence unavailable");
}
const requests = await readClassifiedMockRequests(mockBaseUrl);
expect(
requests.filter((request) => request.prompt === "recovery").length,
).toBeGreaterThanOrEqual(4);
expect(requests.filter((request) => request.prompt === "queued")).toEqual([
expect.objectContaining({ outcome: "success" }),
]);
const finalEvents = (await readStability(gateway, baselineSeq)).events ?? [];
expect(