From 9bf766fa8991dfcde16cf97018272ace62ffd426 Mon Sep 17 00:00:00 2001 From: Heming Zeng Date: Wed, 26 Aug 2026 02:24:27 +0800 Subject: [PATCH] fix(diagnostics): defer watchdog abort during preflight (#125929) --- src/agents/embedded-agent-runner/runs.ts | 11 +- .../reply/agent-runner-memory.test.ts | 5 + .../reply/reply-run-registry.registry.ts | 7 - .../reply/reply-run-registry.test.ts | 4 - src/auto-reply/reply/reply-run-registry.ts | 2 - ...stuck-session-recovery.integration.test.ts | 180 ++++++++++++++++++ ...tic-stuck-session-recovery.runtime.test.ts | 150 ++++++++++++++- ...agnostic-stuck-session-recovery.runtime.ts | 23 ++- 8 files changed, 349 insertions(+), 33 deletions(-) diff --git a/src/agents/embedded-agent-runner/runs.ts b/src/agents/embedded-agent-runner/runs.ts index ddb0efdb62b3..10b1e8b5075f 100644 --- a/src/agents/embedded-agent-runner/runs.ts +++ b/src/agents/embedded-agent-runner/runs.ts @@ -18,10 +18,8 @@ import { resolveActiveReplyOperationForSessionId, resolveActiveReplyRunSessionId, resolveReplyBackendQueueMessageMismatch, - resolveReplyRunPhaseForSessionId, supersedeReplyRunByRunId, type ReplyOperation, - type ReplyOperationPhase, waitForReplyOperationOwnerSettlement, waitForReplyRunEndBySessionId, } from "../../auto-reply/reply/reply-run-registry.js"; @@ -802,10 +800,13 @@ export function isEmbeddedAgentRunInProgress(sessionId: string): boolean { return resolveEmbeddedAgentRunProgressState(sessionId) !== undefined; } -export function resolveEmbeddedAgentReplyRunPhase( +export function resolveEmbeddedReplyActivity( sessionId: string, -): ReplyOperationPhase | undefined { - return resolveReplyRunPhaseForSessionId(sessionId); +): Pick | undefined { + const operation = resolveActiveReplyOperationForSessionId(sessionId); + return operation + ? { phase: operation.phase, lastActivityAtMs: operation.lastActivityAtMs } + : undefined; } export function isEmbeddedAgentRunHandleActive(sessionId: string): boolean { diff --git a/src/auto-reply/reply/agent-runner-memory.test.ts b/src/auto-reply/reply/agent-runner-memory.test.ts index 0f9b7f9575eb..a5610ae79127 100644 --- a/src/auto-reply/reply/agent-runner-memory.test.ts +++ b/src/auto-reply/reply/agent-runner-memory.test.ts @@ -2339,6 +2339,11 @@ describe("runMemoryFlushIfNeeded", () => { const compactCall = requireCompactEmbeddedAgentSessionCall(); expect(compactCall.contextTokenBudget).toBe(200_000); expect(replyOperation.setPhase).toHaveBeenCalledWith("preflight_compacting"); + expect( + replyOperation.setPhase.mock.invocationCallOrder[0] ?? Number.POSITIVE_INFINITY, + ).toBeLessThan( + compactEmbeddedAgentSessionMock.mock.invocationCallOrder[0] ?? Number.NEGATIVE_INFINITY, + ); expect(replyOperation.updateSessionId).not.toHaveBeenCalled(); expect(incrementCompactionCountMock).not.toHaveBeenCalled(); expect(refreshQueuedFollowupSessionMock).not.toHaveBeenCalled(); diff --git a/src/auto-reply/reply/reply-run-registry.registry.ts b/src/auto-reply/reply/reply-run-registry.registry.ts index 0c133dfdcb72..cb2db27f0b50 100644 --- a/src/auto-reply/reply/reply-run-registry.registry.ts +++ b/src/auto-reply/reply/reply-run-registry.registry.ts @@ -11,7 +11,6 @@ import { replyMessageInjectionTargetOperation, replyRunInterruptTargetOperation, type ReplyOperation, - type ReplyOperationPhase, type ReplyRunInterruptTarget, type ReplyRunRegistry, } from "./reply-run-registry.contracts.js"; @@ -232,12 +231,6 @@ export function isReplyRunActiveForSessionId(sessionId: string): boolean { return resolveReplyRunForCurrentSessionId(sessionId) !== undefined; } -export function resolveReplyRunPhaseForSessionId( - sessionId: string, -): ReplyOperationPhase | undefined { - return resolveReplyRunForCurrentSessionId(sessionId)?.phase; -} - export function isReplyRunAbortableForCompaction(sessionId: string): boolean { const operation = resolveReplyRunForCurrentSessionId(sessionId); // Manual compaction uses this as a coordination gate: a finalizing run still diff --git a/src/auto-reply/reply/reply-run-registry.test.ts b/src/auto-reply/reply/reply-run-registry.test.ts index 28931d96015e..8a27d2367bc4 100644 --- a/src/auto-reply/reply/reply-run-registry.test.ts +++ b/src/auto-reply/reply/reply-run-registry.test.ts @@ -42,7 +42,6 @@ import { runAfterReplyOperationClear, resolveActiveReplyRunSessionId, resolveActiveReplyOperationForSessionId, - resolveReplyRunPhaseForSessionId, waitForReplyOperationOwnerSettlement, waitForReplyRunEndBySessionId, waitForReplyRunSuccessorAdmission, @@ -342,9 +341,6 @@ describe("reply run registry", () => { operation.markWaitingForDeferredMaintenance(); expect(operation.phase).toBe("waiting_for_deferred_maintenance"); - expect(resolveReplyRunPhaseForSessionId("session-wait")).toBe( - "waiting_for_deferred_maintenance", - ); expect( getDiagnosticSessionActivitySnapshot({ sessionId: "session-wait", diff --git a/src/auto-reply/reply/reply-run-registry.ts b/src/auto-reply/reply/reply-run-registry.ts index 99217a1902fb..e0efc2075f8d 100644 --- a/src/auto-reply/reply/reply-run-registry.ts +++ b/src/auto-reply/reply/reply-run-registry.ts @@ -12,7 +12,6 @@ export type { ReplyMessageInjectionAttempt, ReplyMessageInjectionTarget, ReplyOperation, - ReplyOperationPhase, ReplyTurnKind, } from "./reply-run-registry.contracts.js"; export { @@ -43,7 +42,6 @@ export { resolveActiveReplyOperationForSessionId, resolveActiveReplyRunSessionId, resolveActiveReplyRunThreadId, - resolveReplyRunPhaseForSessionId, supersedeReplyRunByRunId, waitForReplyOperationOwnerSettlement, waitForReplyRunEndBySessionId, diff --git a/src/logging/diagnostic-stuck-session-recovery.integration.test.ts b/src/logging/diagnostic-stuck-session-recovery.integration.test.ts index 8ff89b4bda69..d7a243cf015c 100644 --- a/src/logging/diagnostic-stuck-session-recovery.integration.test.ts +++ b/src/logging/diagnostic-stuck-session-recovery.integration.test.ts @@ -14,6 +14,7 @@ import { testing as replyRunTesting } from "../auto-reply/reply/reply-run-regist import { onDiagnosticEvent, resetDiagnosticEventsForTest, + setDiagnosticsEnabledForProcess, type DiagnosticEventPayload, } from "../infra/diagnostic-events.js"; import { enqueueCommandInLane, getQueueSize, resetCommandLane } from "../process/command-queue.js"; @@ -25,6 +26,7 @@ import { markDiagnosticRunProgress, } from "./diagnostic-run-activity.js"; import { markDiagnosticModelStartedForTest } from "./diagnostic-run-activity.test-support.js"; +import { logMessageQueuedWithBacklogPolicy } from "./diagnostic-runtime.js"; import { recoverStuckDiagnosticSession } from "./diagnostic-stuck-session-recovery.runtime.js"; import { logSessionStateChange, startDiagnosticHeartbeat } from "./diagnostic.js"; import { resetDiagnosticStateForTest } from "./diagnostic.test-support.js"; @@ -340,6 +342,184 @@ describe("stuck session recovery integration", () => { expect(getQueueSize(lane)).toBe(0); }); + it("keeps queued preflight compaction alive until its configured safety timeout", async () => { + const sessionKey = "agent:main:active-preflight"; + const sessionId = "active-preflight-session"; + const lane = resolveEmbeddedSessionLane(sessionKey); + const operation = createReplyOperation({ sessionKey, sessionId, resetTriggered: false }); + operation.setPhase("preflight_compacting"); + let markActiveStarted!: () => void; + const activeStarted = new Promise((resolve) => { + markActiveStarted = resolve; + }); + const active = enqueueCommandInLane( + lane, + () => + new Promise<"aborted">((resolve) => { + markActiveStarted(); + operation.abortSignal.addEventListener( + "abort", + () => { + operation.complete(); + resolve("aborted"); + }, + { once: true }, + ); + }), + { warnAfterMs: Number.MAX_SAFE_INTEGER }, + ); + const queued = enqueueCommandInLane(lane, async () => "drained", { + warnAfterMs: Number.MAX_SAFE_INTEGER, + }); + await activeStarted; + + const outcome = await recoverStuckDiagnosticSession({ + sessionId, + sessionKey, + // The session can be old even though it only just entered preflight. + ageMs: 30 * 60_000, + queueDepth: 1, + allowActiveAbort: true, + compactionSafetyTimeoutMs: 10 * 60_000, + }); + + expect(outcome).toMatchObject({ + status: "skipped", + action: "keep_lane", + reason: "active_reply_work", + activeSessionId: sessionId, + }); + expect(operation.abortSignal.aborted).toBe(false); + await expectPendingAfterEventLoopTurn(active); + await expectPendingAfterEventLoopTurn(queued); + expect(getQueueSize(lane)).toBe(2); + + // Intentional cancellation remains owned by the reply operation. + expect(operation.abortByUser()).toBe(true); + await expect(active).resolves.toBe("aborted"); + await expect(queued).resolves.toBe("drained"); + }); + + it("keeps fresh preflight compaction through the queued-session heartbeat watchdog", async () => { + vi.useFakeTimers(); + const events: DiagnosticEventPayload[] = []; + const unsubscribe = onDiagnosticEvent((event) => events.push(event)); + try { + const sessionKey = "agent:main:heartbeat-preflight"; + const sessionId = "heartbeat-preflight-session"; + const lane = resolveEmbeddedSessionLane(sessionKey); + const startMs = Date.parse("2026-08-18T12:00:00Z"); + vi.setSystemTime(startMs); + setDiagnosticsEnabledForProcess(true); + logSessionStateChange({ sessionId, sessionKey, state: "processing" }); + logMessageQueuedWithBacklogPolicy({ sessionId, sessionKey, source: "test" }, true); + markDiagnosticEmbeddedRunStarted({ sessionId, sessionKey }); + + // The diagnostic owner is old, but preflight starts only now. + vi.setSystemTime(startMs + 12 * 60_000); + const operation = createReplyOperation({ sessionKey, sessionId, resetTriggered: false }); + operation.setPhase("preflight_compacting"); + let markActiveStarted!: () => void; + const activeStarted = new Promise((resolve) => { + markActiveStarted = resolve; + }); + const active = enqueueCommandInLane( + lane, + () => + new Promise<"aborted">((resolve) => { + markActiveStarted(); + operation.abortSignal.addEventListener( + "abort", + () => { + operation.complete(); + resolve("aborted"); + }, + { once: true }, + ); + }), + { warnAfterMs: Number.MAX_SAFE_INTEGER }, + ); + const queued = enqueueCommandInLane(lane, async () => "drained", { + warnAfterMs: Number.MAX_SAFE_INTEGER, + }); + await activeStarted; + + startDiagnosticHeartbeat( + { + diagnostics: { enabled: true }, + agents: { defaults: { compaction: { timeoutSeconds: 600 } } }, + }, + { + recoverStuckSession: recoverStuckDiagnosticSession, + testTimings: { stuckSessionWarnMs: 30_000, stuckSessionAbortMs: 90_000 }, + }, + ); + await vi.advanceTimersByTimeAsync(90_000); + await Promise.resolve(); + + expect(events).toContainEqual( + expect.objectContaining({ + type: "session.recovery.requested", + sessionId, + sessionKey, + queueDepth: 1, + allowActiveAbort: true, + }), + ); + expect(events).toContainEqual( + expect.objectContaining({ + type: "session.recovery.completed", + sessionId, + sessionKey, + status: "skipped", + action: "keep_lane", + outcomeReason: "active_reply_work", + }), + ); + expect(operation.abortSignal.aborted).toBe(false); + expect(getQueueSize(lane)).toBe(2); + + // Restart remains an intentional, immediate cancellation source. + expect(operation.abortForRestart()).toBe(true); + await expect(active).resolves.toBe("aborted"); + await expect(queued).resolves.toBe("drained"); + } finally { + unsubscribe(); + resetDiagnosticStateForTest(); + await vi.runOnlyPendingTimersAsync(); + vi.useRealTimers(); + } + }); + + it("keeps supersession cancellation immediate during protected preflight", async () => { + const sessionKey = "agent:main:superseded-preflight"; + const sessionId = "superseded-preflight-session"; + const operation = createReplyOperation({ sessionKey, sessionId, resetTriggered: false }); + operation.setPhase("preflight_compacting"); + + await expect( + recoverStuckDiagnosticSession({ + sessionId, + sessionKey, + ageMs: 30 * 60_000, + queueDepth: 1, + allowActiveAbort: true, + compactionSafetyTimeoutMs: 10 * 60_000, + }), + ).resolves.toMatchObject({ + status: "skipped", + action: "keep_lane", + reason: "active_reply_work", + }); + + expect(operation.supersede()).toBe(true); + expect(operation.abortSignal.aborted).toBe(true); + expect(operation.result).toEqual({ + kind: "aborted", + code: "aborted_for_supersession", + }); + }); + it("keeps queued lane work behind reply-only force-clear settlement", async () => { vi.useFakeTimers(); try { diff --git a/src/logging/diagnostic-stuck-session-recovery.runtime.test.ts b/src/logging/diagnostic-stuck-session-recovery.runtime.test.ts index 4ac6804eac0d..1771f80a1777 100644 --- a/src/logging/diagnostic-stuck-session-recovery.runtime.test.ts +++ b/src/logging/diagnostic-stuck-session-recovery.runtime.test.ts @@ -20,7 +20,7 @@ const mocks = vi.hoisted(() => ({ resolveActiveEmbeddedRunSessionIdBySessionFile: vi.fn(), resolveActiveEmbeddedRunHandleSessionId: vi.fn(), resolveActiveEmbeddedRunHandleSessionIdBySessionFile: vi.fn(), - resolveEmbeddedAgentReplyRunPhase: vi.fn(), + resolveEmbeddedReplyActivity: vi.fn(), resolveEmbeddedSessionLane: vi.fn((key: string) => `session:${key}`), waitForEmbeddedAgentRunEnd: vi.fn(), getDiagnosticSessionActivitySnapshot: vi.fn(), @@ -58,7 +58,7 @@ vi.mock("../agents/embedded-agent-runner/runs.js", () => ({ resolveActiveEmbeddedRunHandleSessionId: mocks.resolveActiveEmbeddedRunHandleSessionId, resolveActiveEmbeddedRunHandleSessionIdBySessionFile: mocks.resolveActiveEmbeddedRunHandleSessionIdBySessionFile, - resolveEmbeddedAgentReplyRunPhase: mocks.resolveEmbeddedAgentReplyRunPhase, + resolveEmbeddedReplyActivity: mocks.resolveEmbeddedReplyActivity, waitForEmbeddedAgentRunEnd: mocks.waitForEmbeddedAgentRunEnd, })); @@ -103,7 +103,7 @@ function resetMocks() { mocks.resolveActiveEmbeddedRunSessionIdBySessionFile.mockReset(); mocks.resolveActiveEmbeddedRunHandleSessionId.mockReset(); mocks.resolveActiveEmbeddedRunHandleSessionIdBySessionFile.mockReset(); - mocks.resolveEmbeddedAgentReplyRunPhase.mockReset(); + mocks.resolveEmbeddedReplyActivity.mockReset(); mocks.resolveEmbeddedSessionLane.mockClear(); mocks.waitForEmbeddedAgentRunEnd.mockReset(); mocks.getDiagnosticSessionActivitySnapshot.mockReset(); @@ -401,7 +401,10 @@ describe("stuck session recovery", () => { it("keeps the lane while reply work waits for deferred maintenance", async () => { mocks.resolveActiveEmbeddedRunSessionId.mockReturnValue("queued-reply-session"); mocks.resolveActiveEmbeddedRunHandleSessionId.mockReturnValue(undefined); - mocks.resolveEmbeddedAgentReplyRunPhase.mockReturnValue("waiting_for_deferred_maintenance"); + mocks.resolveEmbeddedReplyActivity.mockReturnValue({ + phase: "waiting_for_deferred_maintenance", + lastActivityAtMs: Date.now(), + }); mocks.isEmbeddedAgentRunActive.mockReturnValue(true); mocks.isEmbeddedAgentRunHandleActive.mockReturnValue(false); @@ -430,7 +433,10 @@ describe("stuck session recovery", () => { it("keeps a reply queued on the global lane instead of reclaiming it", async () => { mocks.resolveActiveEmbeddedRunSessionId.mockReturnValue("queued-reply-session"); mocks.resolveActiveEmbeddedRunHandleSessionId.mockReturnValue(undefined); - mocks.resolveEmbeddedAgentReplyRunPhase.mockReturnValue("waiting_for_global_lane"); + mocks.resolveEmbeddedReplyActivity.mockReturnValue({ + phase: "waiting_for_global_lane", + lastActivityAtMs: Date.now(), + }); mocks.isEmbeddedAgentRunActive.mockReturnValue(true); const outcome = await recoverStuckDiagnosticSession({ @@ -494,7 +500,10 @@ describe("stuck session recovery", () => { mocks.resolveActiveEmbeddedRunHandleSessionId.mockReturnValue(undefined); mocks.isEmbeddedAgentRunActive.mockReturnValue(true); mocks.isEmbeddedAgentRunHandleActive.mockReturnValue(false); - mocks.resolveEmbeddedAgentReplyRunPhase.mockReturnValue(phase); + mocks.resolveEmbeddedReplyActivity.mockReturnValue({ + phase, + lastActivityAtMs: Date.now(), + }); mocks.getDiagnosticSessionActivitySnapshot.mockReturnValue({ lastProgressAgeMs: 720_000, }); @@ -556,7 +565,10 @@ describe("stuck session recovery", () => { ); mocks.isEmbeddedAgentRunActive.mockReturnValue(true); mocks.isEmbeddedAgentRunHandleActive.mockReturnValue(hasEmbeddedHandle); - mocks.resolveEmbeddedAgentReplyRunPhase.mockReturnValue(phase); + mocks.resolveEmbeddedReplyActivity.mockReturnValue({ + phase, + lastActivityAtMs: Date.now() - ageMs, + }); mocks.getDiagnosticSessionActivitySnapshot.mockReturnValue({ lastProgressAgeMs: ageMs }); mocks.abortEmbeddedAgentRun.mockReturnValue(true); mocks.waitForEmbeddedAgentRunEnd.mockResolvedValue(true); @@ -579,6 +591,130 @@ describe("stuck session recovery", () => { }, ); + it.each([ + { + name: "keeps fresh queued preflight despite an old session attention age", + replyActivityAgeMs: 1_000, + queueDepth: 1, + compactionSafetyTimeoutMs: 600_000, + expectedAbort: false, + }, + { + name: "clamps a future preflight activity clock instead of treating it as stale", + replyActivityAgeMs: -1_000, + queueDepth: 1, + compactionSafetyTimeoutMs: 600_000, + expectedAbort: false, + }, + { + name: "keeps zero-backlog preflight one millisecond before timeout plus settle", + replyActivityAgeMs: 614_999, + queueDepth: 0, + compactionSafetyTimeoutMs: 600_000, + expectedAbort: false, + }, + { + name: "recovers queued preflight exactly at timeout plus settle", + replyActivityAgeMs: 615_000, + queueDepth: 1, + compactionSafetyTimeoutMs: 600_000, + expectedAbort: true, + }, + { + name: "recovers queued preflight after timeout plus settle", + replyActivityAgeMs: 615_001, + queueDepth: 1, + compactionSafetyTimeoutMs: 600_000, + expectedAbort: true, + }, + { + name: "keeps preflight before the default stale recovery floor", + replyActivityAgeMs: 299_999, + queueDepth: 1, + expectedAbort: false, + }, + { + name: "recovers preflight at the default stale recovery floor", + replyActivityAgeMs: 300_000, + queueDepth: 1, + expectedAbort: true, + }, + { + name: "uses the default stale recovery floor for an invalid compaction timeout", + replyActivityAgeMs: 299_999, + queueDepth: 1, + compactionSafetyTimeoutMs: 0, + expectedAbort: false, + }, + ])( + "$name", + async ({ replyActivityAgeMs, queueDepth, compactionSafetyTimeoutMs, expectedAbort }) => { + const now = 1_800_000; + const dateNow = vi.spyOn(Date, "now").mockReturnValue(now); + try { + mocks.resolveActiveEmbeddedRunSessionId.mockReturnValue("preflight-session"); + mocks.resolveActiveEmbeddedRunHandleSessionId.mockReturnValue(undefined); + mocks.isEmbeddedAgentRunActive.mockReturnValue(true); + mocks.resolveEmbeddedReplyActivity.mockReturnValue({ + phase: "preflight_compacting", + lastActivityAtMs: now - replyActivityAgeMs, + }); + mocks.abortEmbeddedAgentRun.mockReturnValue(true); + mocks.waitForEmbeddedAgentRunEnd.mockResolvedValue(true); + mocks.resetCommandLane.mockReturnValue(0); + + const outcome = await recoverStuckDiagnosticSession({ + sessionId: "preflight-session", + sessionKey: "agent:main:main", + // Deliberately older than every preflight clock in this table. + ageMs: 30 * 60_000, + queueDepth, + allowActiveAbort: true, + compactionSafetyTimeoutMs, + }); + + expect(mocks.abortEmbeddedAgentRun).toHaveBeenCalledTimes(expectedAbort ? 1 : 0); + expect(outcome).toMatchObject( + expectedAbort + ? { status: "aborted", action: "abort_embedded_run" } + : { status: "skipped", action: "keep_lane", reason: "active_reply_work" }, + ); + } finally { + dateNow.mockRestore(); + } + }, + ); + + it("uses reply activity rather than session age for memory flushing", async () => { + const now = Date.now(); + mocks.resolveActiveEmbeddedRunSessionId.mockReturnValue("memory-flush-session"); + mocks.resolveActiveEmbeddedRunHandleSessionId.mockReturnValue(undefined); + mocks.isEmbeddedAgentRunActive.mockReturnValue(true); + mocks.resolveEmbeddedReplyActivity.mockReturnValue({ + phase: "memory_flushing", + lastActivityAtMs: now, + }); + mocks.abortEmbeddedAgentRun.mockReturnValue(true); + mocks.waitForEmbeddedAgentRunEnd.mockResolvedValue(true); + mocks.resetCommandLane.mockReturnValue(0); + + const outcome = await recoverStuckDiagnosticSession({ + sessionId: "memory-flush-session", + sessionKey: "agent:main:main", + ageMs: 30 * 60_000, + queueDepth: 1, + allowActiveAbort: true, + compactionSafetyTimeoutMs: 600_000, + }); + + expect(mocks.abortEmbeddedAgentRun).not.toHaveBeenCalled(); + expect(outcome).toMatchObject({ + status: "skipped", + action: "keep_lane", + reason: "active_reply_work", + }); + }); + it("keeps reply-only ownership with recent progress even with zero queued backlog", async () => { mocks.resolveActiveEmbeddedRunSessionId.mockReturnValue("live-reply-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 8ae3db8e9865..d1960088c29b 100644 --- a/src/logging/diagnostic-stuck-session-recovery.runtime.ts +++ b/src/logging/diagnostic-stuck-session-recovery.runtime.ts @@ -4,7 +4,7 @@ import { abortAndDrainEmbeddedAgentRun, isEmbeddedAgentRunActive, isEmbeddedAgentRunHandleActive, - resolveEmbeddedAgentReplyRunPhase, + resolveEmbeddedReplyActivity, resolveActiveEmbeddedRunSessionId, resolveActiveEmbeddedRunSessionIdBySessionFile, resolveActiveEmbeddedRunHandleSessionId, @@ -177,22 +177,29 @@ export async function recoverStuckDiagnosticSession( let forceCleared = false; const staleActiveProgressAbortMs = resolveStaleActiveProgressAbortMs(params); const staleActiveLaneTaskReleaseMs = resolveStaleActiveLaneTaskReleaseMs(params); - const activeReplyPhase = activeWorkSessionId - ? resolveEmbeddedAgentReplyRunPhase(activeWorkSessionId) + const activeReplyActivity = activeWorkSessionId + ? resolveEmbeddedReplyActivity(activeWorkSessionId) + : undefined; + const activeReplyPhase = activeReplyActivity?.phase; + // Phase changes refresh the reply operation's activity clock. Session + // attention age may predate maintenance, so it cannot own this timeout. + const activeReplyAgeMs = activeReplyActivity + ? Math.max(0, Date.now() - activeReplyActivity.lastActivityAtMs) : undefined; const maintenancePhase = activeReplyPhase === "preflight_compacting" || activeReplyPhase === "memory_flushing"; + const activeMaintenanceProtected = + maintenancePhase && + activeReplyAgeMs !== undefined && + activeReplyAgeMs < staleActiveLaneTaskReleaseMs; - if ( - activeReplyPhase === "waiting_for_global_lane" || - (maintenancePhase && params.ageMs < staleActiveLaneTaskReleaseMs) - ) { + if (activeReplyPhase === "waiting_for_global_lane" || activeMaintenanceProtected) { // Queued replies and configured maintenance own their lane until their // producer finishes or the existing compaction safety window expires. return reportRecoveryOutcome({ status: "skipped", action: "keep_lane", - reason: maintenancePhase ? "active_reply_work" : "global_lane_wait", + reason: activeMaintenanceProtected ? "active_reply_work" : "global_lane_wait", sessionId: params.sessionId, sessionKey: params.sessionKey, activeSessionId: activeWorkSessionId,