feat(FN-5223): anchor staleness and stall detectors to engine activation ti
The merge introduces an engine-activation timestamp as the staleness floor for task age calculations, replacing arbitrary wall-clock thresholds with a runtime-relative anchor. Step 1 adds settings defaults, Steps 2–4 wire the floor helper through project engine, in-process runtime, and task store hy Fusion-Task-Id: FN-5223
This commit is contained in:
committed by
gsxdsm
parent
62f11e69d3
commit
c60045df22
@@ -1696,6 +1696,35 @@ describe("ProjectEngine paused in-review auto-merge behavior", () => {
|
||||
await engine.stop();
|
||||
});
|
||||
|
||||
it("stamps engineActiveSinceMs on global and engine unpause transitions", async () => {
|
||||
const mockStore = createMockStore({ ...baseSettings, autoMerge: true });
|
||||
mocks.currentStore = mockStore.store;
|
||||
const engine = createEngine();
|
||||
|
||||
await engine.start();
|
||||
mockStore.store.updateSettings.mockClear();
|
||||
|
||||
await mockStore.emitSettingsUpdated(
|
||||
{ ...baseSettings, autoMerge: true, globalPause: false },
|
||||
{ ...baseSettings, autoMerge: true, globalPause: true },
|
||||
);
|
||||
|
||||
await mockStore.emitSettingsUpdated(
|
||||
{ ...baseSettings, autoMerge: true, enginePaused: false },
|
||||
{ ...baseSettings, autoMerge: true, enginePaused: true },
|
||||
);
|
||||
|
||||
const activationStampCalls = mockStore.store.updateSettings.mock.calls.filter(
|
||||
([patch]) => patch && typeof patch === "object" && "engineActiveSinceMs" in patch,
|
||||
);
|
||||
expect(activationStampCalls).toHaveLength(2);
|
||||
for (const [patch] of activationStampCalls) {
|
||||
expect((patch as { engineActiveSinceMs: unknown }).engineActiveSinceMs).toEqual(expect.any(Number));
|
||||
}
|
||||
|
||||
await engine.stop();
|
||||
});
|
||||
|
||||
it("resumes deferred startup recovery on engine unpause", async () => {
|
||||
const mockStore = createMockStore(baseSettings);
|
||||
mocks.currentStore = mockStore.store;
|
||||
|
||||
@@ -0,0 +1,134 @@
|
||||
import { EventEmitter } from "node:events";
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import type { Task, TaskStore } from "@fusion/core";
|
||||
import { SelfHealingManager } from "../../self-healing.js";
|
||||
import { StuckTaskDetector } from "../../stuck-task-detector.js";
|
||||
|
||||
function makeTask(overrides: Partial<Task> = {}): Task {
|
||||
return {
|
||||
id: "FN-5223-RI",
|
||||
title: "t",
|
||||
description: "d",
|
||||
column: "in-review",
|
||||
paused: false,
|
||||
status: undefined,
|
||||
steps: [{ name: "s", status: "done" as const }],
|
||||
workflowStepResults: [],
|
||||
dependencies: [],
|
||||
log: [],
|
||||
updatedAt: "2026-01-01T00:00:00.000Z",
|
||||
columnMovedAt: "2026-01-01T00:00:00.000Z",
|
||||
...overrides,
|
||||
} as Task;
|
||||
}
|
||||
|
||||
function createStore(tasks: Task[], settingsOverrides: Record<string, unknown> = {}): TaskStore & EventEmitter {
|
||||
const emitter = new EventEmitter() as TaskStore & EventEmitter;
|
||||
(emitter as any).getSettings = vi.fn().mockResolvedValue({
|
||||
autoMerge: true,
|
||||
globalPause: false,
|
||||
enginePaused: false,
|
||||
inReviewStalledThresholdMs: 60_000,
|
||||
stalePausedReviewThresholdMs: 60_000,
|
||||
stalePausedTodoThresholdMs: 60_000,
|
||||
taskStuckTimeoutMs: 60_000,
|
||||
engineActivationGraceMs: 300_000,
|
||||
...settingsOverrides,
|
||||
});
|
||||
(emitter as any).updateSettings = vi.fn().mockImplementation(async (patch: Record<string, unknown>) => {
|
||||
const current = await (emitter as any).getSettings();
|
||||
(emitter as any).getSettings = vi.fn().mockResolvedValue({ ...current, ...patch });
|
||||
});
|
||||
(emitter as any).listTasks = vi.fn().mockImplementation(async ({ column }: { column?: string } = {}) => {
|
||||
if (!column) return tasks;
|
||||
return tasks.filter((t) => t.column === column);
|
||||
});
|
||||
(emitter as any).logEntry = vi.fn().mockResolvedValue(undefined);
|
||||
(emitter as any).updateTask = vi.fn().mockImplementation(async (taskId: string, updates: Partial<Task>) => {
|
||||
const task = tasks.find((t) => t.id === taskId);
|
||||
if (task) Object.assign(task, updates);
|
||||
return task;
|
||||
});
|
||||
(emitter as any).moveTask = vi.fn().mockResolvedValue(undefined);
|
||||
(emitter as any).recordRunAuditEvent = vi.fn().mockResolvedValue(undefined);
|
||||
return emitter;
|
||||
}
|
||||
|
||||
describe("FN-5223 reliability interactions: engineActiveSince floor", () => {
|
||||
beforeEach(() => vi.useFakeTimers());
|
||||
afterEach(() => vi.useRealTimers());
|
||||
|
||||
it("FN-5223 suppresses in-review stalled signal immediately after activation and restores after grace+threshold", async () => {
|
||||
const now = new Date("2026-01-01T01:00:00.000Z");
|
||||
vi.setSystemTime(now);
|
||||
const task = makeTask({ id: "FN-5223-Q1" });
|
||||
const store = createStore([task], { engineActiveSinceMs: now.getTime() });
|
||||
const manager = new SelfHealingManager(store, { rootDir: "/tmp/repo" });
|
||||
|
||||
expect(await manager.surfaceInReviewStalled()).toBe(0);
|
||||
|
||||
vi.setSystemTime(new Date(now.getTime() + 7 * 60_000));
|
||||
expect(await manager.surfaceInReviewStalled()).toBe(1);
|
||||
manager.stop();
|
||||
});
|
||||
|
||||
it("FN-5223 pause/unpause stamping prevents immediate stale paused review surfacing", async () => {
|
||||
const now = new Date("2026-01-01T01:00:00.000Z");
|
||||
vi.setSystemTime(now);
|
||||
const task = makeTask({ id: "FN-5223-Q2", paused: true });
|
||||
const store = createStore([task], { engineActiveSinceMs: now.getTime() });
|
||||
const manager = new SelfHealingManager(store, { rootDir: "/tmp/repo" });
|
||||
|
||||
await (store as any).updateSettings({ globalPause: true });
|
||||
await (store as any).updateSettings({ globalPause: false, engineActiveSinceMs: Date.now() });
|
||||
|
||||
expect(await manager.surfaceStalePausedReviews()).toBe(0);
|
||||
manager.stop();
|
||||
});
|
||||
|
||||
it("FN-5223 composes with globalPause/enginePaused cycle gate", async () => {
|
||||
vi.setSystemTime(new Date("2026-01-01T01:00:00.000Z"));
|
||||
const task = makeTask({ id: "FN-5223-Q3" });
|
||||
const store = createStore([task], { globalPause: true, engineActiveSinceMs: Date.now() - 60_000 });
|
||||
const manager = new SelfHealingManager(store, { rootDir: "/tmp/repo" });
|
||||
|
||||
expect(await manager.surfaceInReviewStalled()).toBe(0);
|
||||
manager.stop();
|
||||
});
|
||||
|
||||
it("FN-5223 grace 0 disables warmup", async () => {
|
||||
const now = new Date("2026-01-01T01:00:00.000Z");
|
||||
vi.setSystemTime(now);
|
||||
const task = makeTask({ id: "FN-5223-Q4" });
|
||||
const store = createStore([task], {
|
||||
inReviewStalledThresholdMs: 1,
|
||||
engineActivationGraceMs: 0,
|
||||
engineActiveSinceMs: now.getTime() - 10_000,
|
||||
});
|
||||
const manager = new SelfHealingManager(store, { rootDir: "/tmp/repo" });
|
||||
|
||||
expect(await manager.surfaceInReviewStalled()).toBe(1);
|
||||
manager.stop();
|
||||
});
|
||||
|
||||
it("FN-5223 stuck detector resume hook refreshes tracked timestamps independent of persisted clock", async () => {
|
||||
const now = new Date("2026-01-01T01:00:00.000Z");
|
||||
vi.setSystemTime(now);
|
||||
const store = createStore([]);
|
||||
const detector = new StuckTaskDetector(store);
|
||||
const session = { dispose: vi.fn() };
|
||||
|
||||
detector.trackTask("FN-5223-Q5", session);
|
||||
const trackedBeforePause = (detector as any).tracked.get("FN-5223-Q5");
|
||||
expect(trackedBeforePause.lastActivity).toBe(now.getTime());
|
||||
|
||||
detector.pause();
|
||||
vi.setSystemTime(new Date(now.getTime() + 2 * 60_000));
|
||||
detector.resume();
|
||||
|
||||
const trackedAfterResume = (detector as any).tracked.get("FN-5223-Q5");
|
||||
expect(trackedAfterResume.lastActivity).toBe(now.getTime() + 2 * 60_000);
|
||||
expect(trackedAfterResume.lastProgressAt).toBe(now.getTime() + 2 * 60_000);
|
||||
expect(trackedAfterResume.activitySinceProgress).toBe(0);
|
||||
});
|
||||
});
|
||||
@@ -4702,6 +4702,24 @@ describe("SelfHealingManager", () => {
|
||||
managerWithRecovery.stop();
|
||||
});
|
||||
|
||||
it("suppresses transient-merge stall surfacing when engine activation floor is recent", async () => {
|
||||
vi.setSystemTime(new Date("2026-01-01T00:10:00.000Z"));
|
||||
const managerWithRecovery = new SelfHealingManager(store, { rootDir: "/tmp/test-project" });
|
||||
(store.getSettings as ReturnType<typeof vi.fn>).mockResolvedValue({
|
||||
taskStuckTimeoutMs: 60_000,
|
||||
autoMerge: true,
|
||||
engineActiveSinceMs: Date.parse("2026-01-01T00:10:00.000Z"),
|
||||
engineActivationGraceMs: 300_000,
|
||||
});
|
||||
(store.listTasks as ReturnType<typeof vi.fn>).mockResolvedValue([
|
||||
staleMergingTask({ mergeDetails: { mergeConfirmed: true } }),
|
||||
]);
|
||||
|
||||
expect(await managerWithRecovery.surfaceInReviewStalls()).toBe(0);
|
||||
expect(store.logEntry).not.toHaveBeenCalled();
|
||||
managerWithRecovery.stop();
|
||||
});
|
||||
|
||||
it("skips entirely when autoMerge is disabled", async () => {
|
||||
vi.setSystemTime(new Date("2026-01-01T00:10:00.000Z"));
|
||||
const managerWithRecovery = new SelfHealingManager(store, { rootDir: "/tmp/test-project" });
|
||||
@@ -4998,6 +5016,22 @@ describe("SelfHealingManager", () => {
|
||||
managerWithRecovery.stop();
|
||||
});
|
||||
|
||||
it("suppresses quiet in-review surfacing when engine activation floor is recent", async () => {
|
||||
vi.setSystemTime(new Date("2026-01-02T01:00:00.000Z"));
|
||||
const managerWithRecovery = new SelfHealingManager(store, { rootDir: "/tmp/test-project" });
|
||||
(store.getSettings as ReturnType<typeof vi.fn>).mockResolvedValue({
|
||||
inReviewStalledThresholdMs: 24 * 60 * 60_000,
|
||||
autoMerge: true,
|
||||
engineActiveSinceMs: Date.parse("2026-01-02T01:00:00.000Z"),
|
||||
engineActivationGraceMs: 300_000,
|
||||
});
|
||||
(store.listTasks as ReturnType<typeof vi.fn>).mockResolvedValue([inReviewTask()]);
|
||||
|
||||
expect(await managerWithRecovery.surfaceInReviewStalled()).toBe(0);
|
||||
expect(store.logEntry).not.toHaveBeenCalled();
|
||||
managerWithRecovery.stop();
|
||||
});
|
||||
|
||||
it("skips for recent activity, paused, global pause, engine pause, autoMerge off, threshold off, executing, and active merge", async () => {
|
||||
vi.setSystemTime(new Date("2026-01-02T01:00:00.000Z"));
|
||||
const managerWithRecovery = new SelfHealingManager(store, {
|
||||
@@ -5089,6 +5123,21 @@ describe("SelfHealingManager", () => {
|
||||
managerWithRecovery.stop();
|
||||
});
|
||||
|
||||
it("suppresses stale paused review surfacing when engine activation floor is recent", async () => {
|
||||
vi.setSystemTime(new Date("2026-01-02T01:00:00.000Z"));
|
||||
const managerWithRecovery = new SelfHealingManager(store, { rootDir: "/tmp/test-project" });
|
||||
(store.getSettings as ReturnType<typeof vi.fn>).mockResolvedValue({
|
||||
stalePausedReviewThresholdMs: 24 * 60 * 60_000,
|
||||
engineActiveSinceMs: Date.parse("2026-01-02T01:00:00.000Z"),
|
||||
engineActivationGraceMs: 300_000,
|
||||
});
|
||||
(store.listTasks as ReturnType<typeof vi.fn>).mockResolvedValue([pausedReviewTask()]);
|
||||
|
||||
expect(await managerWithRecovery.surfaceStalePausedReviews()).toBe(0);
|
||||
expect(store.logEntry).not.toHaveBeenCalled();
|
||||
managerWithRecovery.stop();
|
||||
});
|
||||
|
||||
it("logs disposition recommendation when threshold met", async () => {
|
||||
vi.setSystemTime(new Date("2026-01-02T01:00:00.000Z"));
|
||||
const managerWithRecovery = new SelfHealingManager(store, { rootDir: "/tmp/test-project" });
|
||||
@@ -5187,6 +5236,21 @@ describe("SelfHealingManager", () => {
|
||||
managerWithRecovery.stop();
|
||||
});
|
||||
|
||||
it("suppresses stale paused todo surfacing when engine activation floor is recent", async () => {
|
||||
vi.setSystemTime(new Date("2026-01-02T01:00:00.000Z"));
|
||||
const managerWithRecovery = new SelfHealingManager(store, { rootDir: "/tmp/test-project" });
|
||||
(store.getSettings as ReturnType<typeof vi.fn>).mockResolvedValue({
|
||||
stalePausedTodoThresholdMs: 24 * 60 * 60_000,
|
||||
engineActiveSinceMs: Date.parse("2026-01-02T01:00:00.000Z"),
|
||||
engineActivationGraceMs: 300_000,
|
||||
});
|
||||
(store.listTasks as ReturnType<typeof vi.fn>).mockResolvedValue([pausedTodoTask()]);
|
||||
|
||||
expect(await managerWithRecovery.surfaceStalePausedTodos()).toBe(0);
|
||||
expect(store.logEntry).not.toHaveBeenCalled();
|
||||
managerWithRecovery.stop();
|
||||
});
|
||||
|
||||
it("skips under threshold and for unpaused/non-todo tasks", async () => {
|
||||
vi.setSystemTime(new Date("2026-01-01T00:10:00.000Z"));
|
||||
const managerWithRecovery = new SelfHealingManager(store, { rootDir: "/tmp/test-project" });
|
||||
|
||||
@@ -2416,6 +2416,14 @@ export class ProjectEngine {
|
||||
);
|
||||
}
|
||||
|
||||
try {
|
||||
await store.updateSettings({ engineActiveSinceMs: Date.now() });
|
||||
} catch (err: unknown) {
|
||||
runtimeLog.warn(
|
||||
`${source}: failed to stamp engineActiveSinceMs: ${err instanceof Error ? err.message : String(err)}`,
|
||||
);
|
||||
}
|
||||
|
||||
if (settings.globalPause || settings.enginePaused || !settings.autoMerge) {
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -20,6 +20,7 @@ const {
|
||||
mockResumeOrphaned,
|
||||
mockTaskStoreSettings,
|
||||
mockTaskStoreGetTask,
|
||||
mockTaskStoreUpdateSettings,
|
||||
mockMessageStoreSetHook,
|
||||
mockSchedulerConfigurePrMonitoring,
|
||||
mockIsGitRepository,
|
||||
@@ -37,6 +38,7 @@ const {
|
||||
mockResumeOrphaned: vi.fn().mockResolvedValue(undefined),
|
||||
mockTaskStoreSettings: {} as Record<string, unknown>,
|
||||
mockTaskStoreGetTask: vi.fn().mockResolvedValue(null),
|
||||
mockTaskStoreUpdateSettings: vi.fn().mockResolvedValue(undefined),
|
||||
mockMessageStoreSetHook: vi.fn(),
|
||||
mockSchedulerConfigurePrMonitoring: vi.fn(),
|
||||
mockIsGitRepository: vi.fn().mockResolvedValue(true),
|
||||
@@ -72,6 +74,7 @@ vi.mock("@fusion/core", async () => {
|
||||
self.updateTask = vi.fn().mockImplementation(async (taskId: string, patch: Record<string, unknown>) => ({ id: taskId, ...patch }));
|
||||
self.moveTask = vi.fn().mockResolvedValue(undefined);
|
||||
self.getSettings = vi.fn().mockImplementation(async () => structuredClone(mockTaskStoreSettings));
|
||||
self.updateSettings = mockTaskStoreUpdateSettings;
|
||||
self.getMissionStore = vi.fn().mockReturnValue({
|
||||
listMissions: vi.fn().mockReturnValue([]),
|
||||
getMissionWithHierarchy: vi.fn().mockReturnValue(null),
|
||||
@@ -286,6 +289,19 @@ describe("InProcessRuntime", () => {
|
||||
expect(runtime.getStatus()).toBe("active");
|
||||
}, 30000);
|
||||
|
||||
it("stamps engineActiveSinceMs during runtime start", async () => {
|
||||
const before = Date.now();
|
||||
await runtime.start();
|
||||
const after = Date.now();
|
||||
|
||||
expect(mockTaskStoreUpdateSettings).toHaveBeenCalledWith(
|
||||
expect.objectContaining({ engineActiveSinceMs: expect.any(Number) }),
|
||||
);
|
||||
const stamp = (mockTaskStoreUpdateSettings.mock.calls.at(-1)?.[0] as { engineActiveSinceMs: number }).engineActiveSinceMs;
|
||||
expect(stamp).toBeGreaterThanOrEqual(before);
|
||||
expect(stamp).toBeLessThanOrEqual(after);
|
||||
});
|
||||
|
||||
it("does not spawn real git subprocesses during start()", async () => {
|
||||
const execSpy = vi.spyOn(childProcess, "exec");
|
||||
const execFileSpy = vi.spyOn(childProcess, "execFile");
|
||||
|
||||
@@ -777,6 +777,13 @@ export class InProcessRuntime
|
||||
// 14. Start MissionAutopilot background polling
|
||||
this.missionAutopilot?.start();
|
||||
|
||||
try {
|
||||
await this.taskStore.updateSettings({ engineActiveSinceMs: Date.now() });
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? error.message : String(error);
|
||||
runtimeLog.warn(`Failed to stamp engineActiveSinceMs on runtime start: ${message}`);
|
||||
}
|
||||
|
||||
this.setStatus("active");
|
||||
runtimeLog.log(`InProcessRuntime started for project ${this.config.projectId}`);
|
||||
} catch (error) {
|
||||
|
||||
@@ -4130,6 +4130,8 @@ export class SelfHealingManager {
|
||||
executingTaskIds,
|
||||
staleMergingMinAgeMs: this.options.staleMergingStatusMinAgeMs ?? DEFAULT_STALE_MERGING_STATUS_MIN_AGE_MS,
|
||||
maxAutoMergeRetries: MAX_AUTO_MERGE_RETRIES,
|
||||
engineActiveSinceMs: settings.engineActiveSinceMs,
|
||||
engineActivationGraceMs: settings.engineActivationGraceMs,
|
||||
});
|
||||
if (!signal) continue;
|
||||
|
||||
@@ -4266,6 +4268,8 @@ export class SelfHealingManager {
|
||||
autoMerge: true,
|
||||
activeMergeTaskId,
|
||||
executingTaskIds,
|
||||
engineActiveSinceMs: settings.engineActiveSinceMs,
|
||||
engineActivationGraceMs: settings.engineActivationGraceMs,
|
||||
});
|
||||
if (!signal) continue;
|
||||
|
||||
@@ -4315,7 +4319,12 @@ export class SelfHealingManager {
|
||||
|
||||
for (const task of tasks) {
|
||||
if (task.paused !== true) continue;
|
||||
const signal = getStalePausedReviewSignal(task, { now: cycleStartMs, thresholdMs });
|
||||
const signal = getStalePausedReviewSignal(task, {
|
||||
now: cycleStartMs,
|
||||
thresholdMs,
|
||||
engineActiveSinceMs: settings.engineActiveSinceMs,
|
||||
engineActivationGraceMs: settings.engineActivationGraceMs,
|
||||
});
|
||||
if (!signal) continue;
|
||||
if (Date.parse(task.updatedAt) >= cycleStartMs) continue;
|
||||
|
||||
@@ -4361,7 +4370,12 @@ export class SelfHealingManager {
|
||||
|
||||
for (const task of tasks) {
|
||||
if (task.paused !== true) continue;
|
||||
const signal = getStalePausedTodoSignal(task, { now: cycleStartMs, thresholdMs });
|
||||
const signal = getStalePausedTodoSignal(task, {
|
||||
now: cycleStartMs,
|
||||
thresholdMs,
|
||||
engineActiveSinceMs: settings.engineActiveSinceMs,
|
||||
engineActivationGraceMs: settings.engineActivationGraceMs,
|
||||
});
|
||||
if (!signal) continue;
|
||||
if (Date.parse(task.updatedAt) >= cycleStartMs) continue;
|
||||
|
||||
|
||||
Reference in New Issue
Block a user