fix(diagnostics): defer watchdog abort during preflight (#125929)

This commit is contained in:
Heming Zeng
2026-08-26 02:24:27 +08:00
committed by GitHub
parent 4b44cfce5f
commit 9bf766fa89
8 changed files with 349 additions and 33 deletions
+6 -5
View File
@@ -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<ReplyOperation, "phase" | "lastActivityAtMs"> | undefined {
const operation = resolveActiveReplyOperationForSessionId(sessionId);
return operation
? { phase: operation.phase, lastActivityAtMs: operation.lastActivityAtMs }
: undefined;
}
export function isEmbeddedAgentRunHandleActive(sessionId: string): boolean {
@@ -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();
@@ -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
@@ -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",
@@ -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,
@@ -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<void>((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<void>((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 {
@@ -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);
@@ -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,