diff --git a/extensions/qa-lab/src/providers/mock-openai/mock-openai-contracts.ts b/extensions/qa-lab/src/providers/mock-openai/mock-openai-contracts.ts index 00676ebd8f48..0156577f7860 100644 --- a/extensions/qa-lab/src/providers/mock-openai/mock-openai-contracts.ts +++ b/extensions/qa-lab/src/providers/mock-openai/mock-openai-contracts.ts @@ -35,6 +35,7 @@ export type QaMockProviderDispatchResult = { failure?: QaMockProviderFailure; onResponseSent?: () => void; previewPauseMs?: number; + responsePauseMs?: number; }; export type StreamEvent = diff --git a/src/agents/embedded-agent-runner/runs.force-clear-terminal.test.ts b/src/agents/embedded-agent-runner/runs.force-clear-terminal.test.ts index 24594d0be724..9557ff55473e 100644 --- a/src/agents/embedded-agent-runner/runs.force-clear-terminal.test.ts +++ b/src/agents/embedded-agent-runner/runs.force-clear-terminal.test.ts @@ -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"; diff --git a/src/agents/embedded-agent-runner/runs.test.ts b/src/agents/embedded-agent-runner/runs.test.ts index cde30fcf7c31..c5af9576f835 100644 --- a/src/agents/embedded-agent-runner/runs.test.ts +++ b/src/agents/embedded-agent-runner/runs.test.ts @@ -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, diff --git a/src/agents/embedded-agent-runner/runs.ts b/src/agents/embedded-agent-runner/runs.ts index b029a7416f70..a0b5b47aa9c4 100644 --- a/src/agents/embedded-agent-runner/runs.ts +++ b/src/agents/embedded-agent-runner/runs.ts @@ -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 { 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((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((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((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 = { diff --git a/src/auto-reply/reply/reply-run-registry.test.ts b/src/auto-reply/reply/reply-run-registry.test.ts index c54516212ead..9425f4a44efb 100644 --- a/src/auto-reply/reply/reply-run-registry.test.ts +++ b/src/auto-reply/reply/reply-run-registry.test.ts @@ -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((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((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((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", diff --git a/src/auto-reply/reply/reply-run-registry.ts b/src/auto-reply/reply/reply-run-registry.ts index b0e272b8b8e7..8ce0834565ae 100644 --- a/src/auto-reply/reply/reply-run-registry.ts +++ b/src/auto-reply/reply/reply-run-registry.ts @@ -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; 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; + 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 | 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 { + 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((resolve) => { + timer = setTimeout(() => resolve(false), resolvedTimeoutMs); + timer.unref?.(); + }), + ]); + if (timer) { + clearTimeout(timer); + } + return settled; } export function expireStaleReplyRunBySessionId( sessionId: string, reason: ReplyOperationStaleReason, + options?: Parameters[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 diff --git a/test/e2e/qa-lab/runtime/gateway-repeated-request-recovery.e2e.test.ts b/test/e2e/qa-lab/runtime/gateway-repeated-request-recovery.e2e.test.ts index afe450eac0ac..96dda730f192 100644 --- a/test/e2e/qa-lab/runtime/gateway-repeated-request-recovery.e2e.test.ts +++ b/test/e2e/qa-lab/runtime/gateway-repeated-request-recovery.e2e.test.ts @@ -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> | undefined; @@ -187,6 +203,66 @@ async function waitForQueuedReply( : new Error(message, { cause: lastRetryableError }); } +async function readClassifiedMockRequests(mockBaseUrl: string): Promise { + return fetch(`${mockBaseUrl}/debug/requests`) + .then((response) => response.json() as Promise) + .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>["gateway"]; + mockBaseUrl: string | undefined; + sinceSeq: number; +}): Promise { + 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(