From f24e164f694cd447eb657ddcc34e62ffb794bcad Mon Sep 17 00:00:00 2001 From: licheer-zte Date: Thu, 13 Aug 2026 08:42:48 +0800 Subject: [PATCH] fix(recovery): reclaim proven-stale reply-only ownership with zero queued backlog (#122265) * fix(recovery): reclaim proven-stale reply-only ownership with zero queued backlog Stuck-session recovery kept reply-only ownership forever when the queued backlog was empty: isActiveRunProgressStale short-circuits to false at queueDepth 0, so stale active_reply_work was never reclaimed even after the durable session became killed. Evaluate staleness for reply-only ownership without the queue gate (the gate stays for run-handle paths), so proven-stale reply work expires through the existing abort-and-drain owner path while global-lane and deferred-maintenance exemptions and live reply work with recent progress are preserved. Closes #122227 * fix(recovery): keep maintenance phases out of zero-backlog stale reclaim The zero-backlog reclaim path (requireQueueBacklog: false) applied to every reply-only operation, including preflight_compacting and memory_flushing. Those phases are explicitly recognized as compaction and may honor a configured timeout above the stale threshold, so a valid long-running maintenance operation could be force-cleared early. Restore the queue-backlog guard for the maintenance phases so an unqueued compaction or memory flush is never reclaimed by this path; ordinary reply-only ownership keeps the zero-backlog expiry. Adds regressions for both maintenance phases. --- ...tic-stuck-session-recovery.runtime.test.ts | 115 ++++++++++++++++++ ...agnostic-stuck-session-recovery.runtime.ts | 21 +++- 2 files changed, 135 insertions(+), 1 deletion(-) diff --git a/src/logging/diagnostic-stuck-session-recovery.runtime.test.ts b/src/logging/diagnostic-stuck-session-recovery.runtime.test.ts index 2cfaeafc95a4..d9cf513b14ad 100644 --- a/src/logging/diagnostic-stuck-session-recovery.runtime.test.ts +++ b/src/logging/diagnostic-stuck-session-recovery.runtime.test.ts @@ -452,6 +452,121 @@ describe("stuck session recovery", () => { expect(mocks.resetCommandLane).not.toHaveBeenCalled(); }); + it("reclaims proven-stale reply-only ownership even with zero queued backlog", async () => { + mocks.resolveActiveEmbeddedRunSessionId.mockReturnValue("phantom-reply-session"); + mocks.resolveActiveEmbeddedRunHandleSessionId.mockReturnValue(undefined); + mocks.isEmbeddedAgentRunActive.mockReturnValue(true); + mocks.isEmbeddedAgentRunHandleActive.mockReturnValue(false); + mocks.abortEmbeddedAgentRun.mockReturnValue(true); + mocks.waitForEmbeddedAgentRunEnd.mockResolvedValue(true); + + const outcome = await recoverStuckDiagnosticSession({ + sessionId: "phantom-reply-session", + sessionKey: "agent:main:main", + ageMs: 720_000, + queueDepth: 0, + }); + + // Stale reply-only ownership must expire through the abort-and-drain owner + // path even when the queued backlog is empty; previously the zero-depth + // gate kept the lane forever (reason=active_reply_work). + expect(mocks.abortEmbeddedAgentRun).toHaveBeenCalledWith("phantom-reply-session"); + expect(mocks.waitForEmbeddedAgentRunEnd).toHaveBeenCalledWith("phantom-reply-session", 15_000); + expect(outcome).toMatchObject({ + status: "aborted", + action: "abort_embedded_run", + activeSessionId: "phantom-reply-session", + activeWorkKind: "embedded_run", + aborted: true, + drained: true, + }); + expect(warnLogMessages()).toEqual([ + "stuck session recovery reclaiming stale active reply work: sessionId=phantom-reply-session sessionKey=agent:main:main age=720s queueDepth=0 activeSessionId=phantom-reply-session", + "stuck session recovery: sessionId=phantom-reply-session sessionKey=agent:main:main age=720s action=abort_embedded_run aborted=true drained=true released=0", + "stuck session recovery outcome: status=aborted action=abort_embedded_run sessionId=phantom-reply-session sessionKey=agent:main:main activeSessionId=phantom-reply-session activeWorkKind=embedded_run lane=session:agent:main:main aborted=true drained=true forceCleared=false released=0", + ]); + }); + + it.each(["preflight_compacting", "memory_flushing"])( + "keeps zero-backlog maintenance phase %s out of the stale reclaim path", + async (phase) => { + mocks.resolveActiveEmbeddedRunSessionId.mockReturnValue("maintenance-reply-session"); + mocks.resolveActiveEmbeddedRunHandleSessionId.mockReturnValue(undefined); + mocks.isEmbeddedAgentRunActive.mockReturnValue(true); + mocks.isEmbeddedAgentRunHandleActive.mockReturnValue(false); + mocks.resolveEmbeddedAgentReplyRunPhase.mockReturnValue(phase); + mocks.getDiagnosticSessionActivitySnapshot.mockReturnValue({ + lastProgressAgeMs: 720_000, + }); + + const outcome = await recoverStuckDiagnosticSession({ + sessionId: "maintenance-reply-session", + sessionKey: "agent:main:main", + ageMs: 720_000, + queueDepth: 0, + }); + + // Preflight compaction and memory flush are recognized maintenance + // phases that may legitimately outlive the stale threshold (they honor a + // configured compaction timeout). The zero-backlog exemption must not + // turn a running maintenance operation into a reclaim target. + expect(outcome).toMatchObject({ + status: "skipped", + action: "keep_lane", + reason: "active_reply_work", + activeSessionId: "maintenance-reply-session", + }); + expect(mocks.abortEmbeddedAgentRun).not.toHaveBeenCalled(); + }, + ); + + it("keeps reply-only ownership with recent progress even with zero queued backlog", async () => { + mocks.resolveActiveEmbeddedRunSessionId.mockReturnValue("live-reply-session"); + mocks.resolveActiveEmbeddedRunHandleSessionId.mockReturnValue(undefined); + mocks.isEmbeddedAgentRunActive.mockReturnValue(true); + mocks.isEmbeddedAgentRunHandleActive.mockReturnValue(false); + mocks.getDiagnosticSessionActivitySnapshot.mockReturnValue({ lastProgressAgeMs: 1_000 }); + + const outcome = await recoverStuckDiagnosticSession({ + sessionId: "live-reply-session", + sessionKey: "agent:main:main", + ageMs: 720_000, + queueDepth: 0, + }); + + // Recent progress means the reply is genuinely active; the zero-depth + // exemption must not turn live work into a reclaim target. + expect(outcome).toMatchObject({ + status: "skipped", + action: "keep_lane", + reason: "active_reply_work", + activeSessionId: "live-reply-session", + }); + expect(mocks.abortEmbeddedAgentRun).not.toHaveBeenCalled(); + }); + + it("keeps the queue gate for active run handles with zero queued backlog", async () => { + mocks.resolveActiveEmbeddedRunHandleSessionId.mockReturnValue("session-1"); + mocks.isEmbeddedAgentRunHandleActive.mockReturnValue(true); + mocks.getDiagnosticSessionActivitySnapshot.mockReturnValue({ lastProgressAgeMs: 720_000 }); + + const outcome = await recoverStuckDiagnosticSession({ + sessionId: "session-1", + sessionKey: "agent:main:main", + ageMs: 720_000, + queueDepth: 0, + }); + + // Run-handle recovery keeps the queue gate: without a queued backlog the + // active run is presumed to be processing and must not be aborted. + expect(outcome).toMatchObject({ + status: "skipped", + action: "observe_only", + reason: "active_embedded_run", + }); + expect(mocks.abortEmbeddedAgentRun).not.toHaveBeenCalled(); + }); + it("releases the session lane when abort+drain succeeds but queued messages remain (ghost run + queued messages)", async () => { mocks.resolveActiveEmbeddedRunSessionId.mockReturnValue("ghost-run-session"); mocks.resolveActiveEmbeddedRunHandleSessionId.mockReturnValue(undefined); diff --git a/src/logging/diagnostic-stuck-session-recovery.runtime.ts b/src/logging/diagnostic-stuck-session-recovery.runtime.ts index 9b433d6ef81a..581627d5e0bf 100644 --- a/src/logging/diagnostic-stuck-session-recovery.runtime.ts +++ b/src/logging/diagnostic-stuck-session-recovery.runtime.ts @@ -62,8 +62,15 @@ function isActiveRunProgressStale(params: { sessionKey?: string; queueDepth?: number; staleAbortMs: number; + /** + * When false, staleness is evaluated even with a zero queued backlog. + * Run-handle recovery keeps the gate so an unqueued active run is not + * disturbed; reply-only ownership has no backlog to protect and must + * still expire when proven stale (phantom active reply work). + */ + requireQueueBacklog?: boolean; }): boolean { - if ((params.queueDepth ?? 0) <= 0) { + if ((params.queueDepth ?? 0) <= 0 && params.requireQueueBacklog !== false) { return false; } const activity = getDiagnosticSessionActivitySnapshot({ @@ -250,6 +257,18 @@ export async function recoverStuckDiagnosticSession( sessionKey: params.sessionKey, queueDepth: params.queueDepth, staleAbortMs: staleActiveProgressAbortMs, + // Reply-only ownership must expire when proven stale even with zero + // queued backlog; the queue gate exists to protect run handles that + // are actively draining queued turns, and there is no such backlog + // here to protect. Recognized maintenance phases are the exception: + // preflight compaction and memory flush are explicitly allowed to + // run longer than the stale threshold (they honor a configured + // compaction timeout), so they keep the queue-backlog guard and are + // never force-cleared early by this reclaim path. + requireQueueBacklog: + activeReplyPhase === "preflight_compacting" || activeReplyPhase === "memory_flushing" + ? undefined + : false, }); if (params.allowActiveAbort === true || reclaimStaleReplyWork) { if (reclaimStaleReplyWork) {