






















@@ -14,6 +14,7 @@ const mocks = vi.hoisted(() => ({
1414resolveActiveEmbeddedRunHandleSessionId: vi.fn(),
1515resolveEmbeddedSessionLane: vi.fn((key: string) => `session:${key}`),
1616waitForEmbeddedPiRunEnd: vi.fn(),
17+getDiagnosticSessionActivitySnapshot: vi.fn(),
1718diag: {
1819debug: vi.fn(),
1920warn: vi.fn(),
@@ -60,6 +61,10 @@ vi.mock("./diagnostic-runtime.js", () => ({
6061diagnosticLogger: mocks.diag,
6162}));
626364+vi.mock("./diagnostic-run-activity.js", () => ({
65+getDiagnosticSessionActivitySnapshot: mocks.getDiagnosticSessionActivitySnapshot,
66+}));
67+6368import {
6469testing,
6570recoverStuckDiagnosticSession,
@@ -85,6 +90,10 @@ function resetMocks() {
8590mocks.resolveActiveEmbeddedRunHandleSessionId.mockReset();
8691mocks.resolveEmbeddedSessionLane.mockClear();
8792mocks.waitForEmbeddedPiRunEnd.mockReset();
93+mocks.getDiagnosticSessionActivitySnapshot.mockReset();
94+// Default: no progress signal, so the staleness gate stays off unless a test
95+// opts in by returning a stale lastProgressAgeMs.
96+mocks.getDiagnosticSessionActivitySnapshot.mockReturnValue({});
8897mocks.diag.debug.mockReset();
8998mocks.diag.warn.mockReset();
9099}
@@ -121,6 +130,26 @@ describe("stuck session recovery", () => {
121130]);
122131});
123132133+it("reclaims a stale active embedded run with queued work and no forward progress (#85639)", async () => {
134+mocks.resolveActiveEmbeddedRunHandleSessionId.mockReturnValue("session-1");
135+mocks.getDiagnosticSessionActivitySnapshot.mockReturnValue({
136+lastProgressAgeMs: 10 * 60_000,
137+});
138+mocks.abortEmbeddedPiRun.mockReturnValue(true);
139+mocks.waitForEmbeddedPiRunEnd.mockResolvedValue(true);
140+141+const outcome = await recoverStuckDiagnosticSession({
142+sessionId: "session-1",
143+sessionKey: "agent:main:main",
144+ageMs: 180_000,
145+queueDepth: 1,
146+});
147+148+expect(mocks.abortEmbeddedPiRun).toHaveBeenCalledWith("session-1");
149+expect(outcome.status).toBe("aborted");
150+expect(warnLogMessages().some((m) => m.includes("reclaiming stale active run"))).toBe(true);
151+});
152+124153it("aborts an active embedded run when active abort recovery is enabled", async () => {
125154mocks.resolveActiveEmbeddedRunHandleSessionId.mockReturnValue("session-1");
126155mocks.abortEmbeddedPiRun.mockReturnValue(true);
@@ -292,6 +321,107 @@ describe("stuck session recovery", () => {
292321]);
293322});
294323324+it("reclaims stale leaked reply work with queued work and no forward progress (#85639)", async () => {
325+mocks.resolveActiveEmbeddedRunSessionId.mockReturnValue("queued-reply-session");
326+mocks.resolveActiveEmbeddedRunHandleSessionId.mockReturnValue(undefined);
327+mocks.isEmbeddedPiRunActive.mockReturnValue(true);
328+mocks.isEmbeddedPiRunHandleActive.mockReturnValue(false);
329+// The "active" run has made no forward progress for well past the staleness
330+// window — a leaked/dead handle, not genuine work.
331+mocks.getDiagnosticSessionActivitySnapshot.mockReturnValue({ lastProgressAgeMs: 10 * 60_000 });
332+mocks.abortEmbeddedPiRun.mockReturnValue(true);
333+mocks.waitForEmbeddedPiRunEnd.mockResolvedValue(true);
334+335+const outcome = await recoverStuckDiagnosticSession({
336+sessionId: "queued-reply-session",
337+sessionKey: "agent:main:main",
338+ageMs: 180_000,
339+queueDepth: 1,
340+});
341+342+// Reclaimed (aborted) instead of skipping with active_reply_work.
343+expect(mocks.abortEmbeddedPiRun).toHaveBeenCalledWith("queued-reply-session");
344+expect(outcome.status).not.toBe("skipped");
345+expect(warnLogMessages().some((m) => m.includes("reclaiming stale active reply work"))).toBe(
346+true,
347+);
348+});
349+350+it("honors an operator-raised stuck-session abort threshold for stale reclaim (#85639)", async () => {
351+mocks.resolveActiveEmbeddedRunSessionId.mockReturnValue("queued-reply-session");
352+mocks.resolveActiveEmbeddedRunHandleSessionId.mockReturnValue(undefined);
353+mocks.isEmbeddedPiRunActive.mockReturnValue(true);
354+mocks.isEmbeddedPiRunHandleActive.mockReturnValue(false);
355+mocks.abortEmbeddedPiRun.mockReturnValue(true);
356+mocks.waitForEmbeddedPiRunEnd.mockResolvedValue(true);
357+358+// Operator raised the abort threshold to 20 min to protect slow active work.
359+const raisedAbortMs = 20 * 60_000;
360+361+// Below the raised threshold (10 min): keep the lane, do not reclaim.
362+mocks.getDiagnosticSessionActivitySnapshot.mockReturnValue({ lastProgressAgeMs: 10 * 60_000 });
363+const kept = await recoverStuckDiagnosticSession({
364+sessionId: "queued-reply-session",
365+sessionKey: "agent:main:main",
366+ageMs: 180_000,
367+queueDepth: 1,
368+staleActiveProgressAbortMs: raisedAbortMs,
369+});
370+expect(mocks.abortEmbeddedPiRun).not.toHaveBeenCalled();
371+expect(kept.status).toBe("skipped");
372+373+// Past the raised threshold (25 min): reclaim.
374+mocks.getDiagnosticSessionActivitySnapshot.mockReturnValue({ lastProgressAgeMs: 25 * 60_000 });
375+const reclaimed = await recoverStuckDiagnosticSession({
376+sessionId: "queued-reply-session",
377+sessionKey: "agent:main:main",
378+ageMs: 180_000,
379+queueDepth: 1,
380+staleActiveProgressAbortMs: raisedAbortMs,
381+});
382+expect(mocks.abortEmbeddedPiRun).toHaveBeenCalledWith("queued-reply-session");
383+expect(reclaimed.status).not.toBe("skipped");
384+});
385+386+it("keeps the lane when active reply work is still progressing", async () => {
387+mocks.resolveActiveEmbeddedRunSessionId.mockReturnValue("queued-reply-session");
388+mocks.resolveActiveEmbeddedRunHandleSessionId.mockReturnValue(undefined);
389+mocks.isEmbeddedPiRunActive.mockReturnValue(true);
390+mocks.isEmbeddedPiRunHandleActive.mockReturnValue(false);
391+// Recent forward progress: a genuinely active run must not be reclaimed.
392+mocks.getDiagnosticSessionActivitySnapshot.mockReturnValue({ lastProgressAgeMs: 5_000 });
393+394+const outcome = await recoverStuckDiagnosticSession({
395+sessionId: "queued-reply-session",
396+sessionKey: "agent:main:main",
397+ageMs: 180_000,
398+queueDepth: 1,
399+});
400+401+expect(mocks.abortEmbeddedPiRun).not.toHaveBeenCalled();
402+expect(outcome.status).toBe("skipped");
403+expect(warnLogMessages().some((m) => m.includes("reason=active_reply_work"))).toBe(true);
404+});
405+406+it("does not reclaim stale reply work when no work is queued", async () => {
407+mocks.resolveActiveEmbeddedRunSessionId.mockReturnValue("queued-reply-session");
408+mocks.resolveActiveEmbeddedRunHandleSessionId.mockReturnValue(undefined);
409+mocks.isEmbeddedPiRunActive.mockReturnValue(true);
410+mocks.isEmbeddedPiRunHandleActive.mockReturnValue(false);
411+mocks.getDiagnosticSessionActivitySnapshot.mockReturnValue({ lastProgressAgeMs: 10 * 60_000 });
412+413+const outcome = await recoverStuckDiagnosticSession({
414+sessionId: "queued-reply-session",
415+sessionKey: "agent:main:main",
416+ageMs: 180_000,
417+queueDepth: 0,
418+});
419+420+expect(mocks.abortEmbeddedPiRun).not.toHaveBeenCalled();
421+expect(outcome.status).toBe("skipped");
422+expect(warnLogMessages().some((m) => m.includes("reason=active_reply_work"))).toBe(true);
423+});
424+295425it("aborts stale reply work without an embedded handle when active abort recovery is enabled", async () => {
296426mocks.resolveActiveEmbeddedRunSessionId.mockReturnValue("queued-reply-session");
297427mocks.resolveActiveEmbeddedRunHandleSessionId.mockReturnValue(undefined);
此内容由惯性聚合(RSS阅读器)自动聚合整理,仅供阅读参考。 原文来自 — 版权归原作者所有。