





















@@ -1487,6 +1487,91 @@ describe("stuck session diagnostics threshold", () => {
14871487}
14881488});
148914891490+it("does not re-emit session.recovery.requested when generation bumps mid-flight (idle-queued stall)", async () => {
1491+const events: DiagnosticEventPayload[] = [];
1492+// Pin recover() to an unresolved Promise so the in-flight window spans two
1493+// heartbeat ticks (production awaits abort/drain, settleMs up to 15s). Same
1494+// seam as the already_in_flight dedup test above.
1495+let resolveRecovery:
1496+| ((outcome: {
1497+status: "skipped";
1498+action: "observe_only";
1499+reason: "already_in_flight";
1500+sessionId: string;
1501+sessionKey: string;
1502+}) => void)
1503+| undefined;
1504+const recoverStuckSession = vi.fn(
1505+() =>
1506+new Promise<{
1507+status: "skipped";
1508+action: "observe_only";
1509+reason: "already_in_flight";
1510+sessionId: string;
1511+sessionKey: string;
1512+}>((resolve) => {
1513+resolveRecovery = resolve;
1514+}),
1515+);
1516+const unsubscribe = onDiagnosticEvent((event) => {
1517+events.push(event);
1518+});
1519+try {
1520+startDiagnosticHeartbeat(
1521+{
1522+diagnostics: {
1523+enabled: true,
1524+stuckSessionWarnMs: 30_000,
1525+stuckSessionAbortMs: 60_000,
1526+},
1527+},
1528+{ recoverStuckSession },
1529+);
1530+// idle-queued-recoverable-stall setup: embedded run ownership + idle + queued.
1531+logSessionStateChange({ sessionId: "s1", sessionKey: "main", state: "processing" });
1532+markDiagnosticEmbeddedRunStarted({ sessionId: "s1", sessionKey: "main" });
1533+logSessionStateChange({ sessionId: "s1", sessionKey: "main", state: "idle" });
1534+1535+// T1 tick: lastProgressAgeMs > staleMs -> recoveryEligible -> coordinator key
1536+// add, requested #1 emitted, recover() in-flight (pending).
1537+vi.advanceTimersByTime(59_000);
1538+logMessageQueued({ sessionId: "s1", sessionKey: "main", source: "t1" });
1539+vi.advanceTimersByTime(1_000);
1540+await Promise.resolve();
1541+expect(recoverStuckSession).toHaveBeenCalledTimes(1);
1542+1543+// New message queued during the in-flight window -> state.generation +1 (but
1544+// lastProgressAt not refreshed). Next tick the coordinator key becomes S:G+1.
1545+logMessageQueued({ sessionId: "s1", sessionKey: "main", source: "t2-followup" });
1546+1547+// T2 tick (30s later): same session still idle-queued-stall -> re-classified.
1548+vi.advanceTimersByTime(30_000);
1549+await Promise.resolve();
1550+} finally {
1551+resolveRecovery?.({
1552+status: "skipped",
1553+action: "observe_only",
1554+reason: "already_in_flight",
1555+sessionId: "s1",
1556+sessionKey: "main",
1557+});
1558+await Promise.resolve();
1559+unsubscribe();
1560+}
1561+1562+const requestedEvents = events.filter(
1563+(event) => event.type === "session.recovery.requested",
1564+);
1565+// Before the fix (RED): coordinator key = `${ref}:${generation}`, so S:G and
1566+// S:G+1 are distinct -> a second requested event is emitted ->
1567+// requestedEvents.length === 2. The runtime sees the same ref and skips the
1568+// actual recovery as already_in_flight, leaving only a duplicate event.
1569+// After the fix (GREEN): coordinator key is ref-only -> S:G+1 collides with the
1570+// in-flight S -> the coordinator also absorbs it as already_in_flight ->
1571+// requestedEvents.length === 1 (both dedup layers share granularity).
1572+expect(requestedEvents).toHaveLength(1);
1573+});
1574+14901575it("reports long-running sessions separately when active work is making progress", () => {
14911576const events: DiagnosticEventPayload[] = [];
14921577const recoverStuckSession = vi.fn();
此内容由惯性聚合(RSS阅读器)自动聚合整理,仅供阅读参考。 原文来自 — 版权归原作者所有。