feat(FN-5076): merge fusion/fn-5076
This commit is contained in:
5
.changeset/FN-5076-insight-run-sweep.md
Normal file
5
.changeset/FN-5076-insight-run-sweep.md
Normal file
@@ -0,0 +1,5 @@
|
|||||||
|
---
|
||||||
|
"@runfusion/fusion": patch
|
||||||
|
---
|
||||||
|
|
||||||
|
Auto-recover stale pending/running insight runs at dashboard startup and on a periodic sweep so manual runs never hang indefinitely.
|
||||||
@@ -361,6 +361,7 @@ Intentional exclusions from shared snapshots:
|
|||||||
- true live conflicts continue returning HTTP 409 with structured payload details `{ code: "ACTIVE_RUN_CONFLICT", activeRunId, activeRunStatus, trigger }` so the dashboard can hydrate and display the existing active run instead of surfacing a raw backend exception
|
- true live conflicts continue returning HTTP 409 with structured payload details `{ code: "ACTIVE_RUN_CONFLICT", activeRunId, activeRunStatus, trigger }` so the dashboard can hydrate and display the existing active run instead of surfacing a raw backend exception
|
||||||
- `POST /api/insights/:id/create-task` remains a draft-payload endpoint (returns `suggestedTitle`/`suggestedDescription`); the dashboard `InsightsView` now uses that payload to create a real task through the normal app task-creation path (`column: triage`, `sourceType: dashboard_ui`, source metadata indicating insights origin)
|
- `POST /api/insights/:id/create-task` remains a draft-payload endpoint (returns `suggestedTitle`/`suggestedDescription`); the dashboard `InsightsView` now uses that payload to create a real task through the normal app task-creation path (`column: triage`, `sourceType: dashboard_ui`, source metadata indicating insights origin)
|
||||||
- Backed by `project_insights`, `project_insight_runs`, and `project_insight_run_events`
|
- Backed by `project_insights`, `project_insight_runs`, and `project_insight_run_events`
|
||||||
|
- **Architecture invariant:** stale `pending`/`running` insight runs auto-recover at dashboard startup and on periodic/drive-by sweeps; active-row conflicts must be evaluated by age plus live `activeRunControllers` ownership instead of assuming all active rows block forever.
|
||||||
|
|
||||||
### Research Runs
|
### Research Runs
|
||||||
|
|
||||||
@@ -710,6 +711,7 @@ Guardrails: this routine does **not** retry merges, does **not** apply to mixed/
|
|||||||
- `[self-healing]` — startup/maintenance recovery pass outcomes.
|
- `[self-healing]` — startup/maintenance recovery pass outcomes.
|
||||||
- `[worktree-metadata-reconcile]` — FN-4962 stale `task.worktree`/`task.branch` rebind-or-clear decisions and audit emission failures.
|
- `[worktree-metadata-reconcile]` — FN-4962 stale `task.worktree`/`task.branch` rebind-or-clear decisions and audit emission failures.
|
||||||
- `[scheduler]`, `[executor]`, `[merger]` — core execution/dispatch/merge lanes.
|
- `[scheduler]`, `[executor]`, `[merger]` — core execution/dispatch/merge lanes.
|
||||||
|
- `[insight-sweeper]` — startup/periodic/drive-by stale insight-run recovery outcomes and fail-soft sweep errors.
|
||||||
- `Notifier` (`notifier.ts`) — legacy ntfy compatibility shim (`NtfyNotifier`) plus shared ntfy helpers
|
- `Notifier` (`notifier.ts`) — legacy ntfy compatibility shim (`NtfyNotifier`) plus shared ntfy helpers
|
||||||
- Runtime ownership: `NtfyNotifier` no longer owns an independent task-lifecycle listener graph; `ProjectEngine` injects the canonical `NotificationService` instance so task lifecycle notifications (`task:moved`, `task:updated`, `task:merged`) are emitted through a single path.
|
- Runtime ownership: `NtfyNotifier` no longer owns an independent task-lifecycle listener graph; `ProjectEngine` injects the canonical `NotificationService` instance so task lifecycle notifications (`task:moved`, `task:updated`, `task:merged`) are emitted through a single path.
|
||||||
- Merge dedup safety: all merge-success → done code paths (direct merger completion, owned/no-op auto-finalize, mergeConfirmed fast-path, PR-strategy finalize, and merge-success self-healing finalizers) emit `store.emit("task:merged", result)` with a merged `MergeResult`. `NotificationService.notifiedEvents` remains the single dedup source of truth, so duplicate upstream emits still produce exactly one canonical `merged` ntfy lifecycle notification per task.
|
- Merge dedup safety: all merge-success → done code paths (direct merger completion, owned/no-op auto-finalize, mergeConfirmed fast-path, PR-strategy finalize, and merge-success self-healing finalizers) emit `store.emit("task:merged", result)` with a merged `MergeResult`. `NotificationService.notifiedEvents` remains the single dedup source of truth, so duplicate upstream emits still produce exactly one canonical `merged` ntfy lifecycle notification per task.
|
||||||
|
|||||||
9
docs/diagnostics.md
Normal file
9
docs/diagnostics.md
Normal file
@@ -0,0 +1,9 @@
|
|||||||
|
# Diagnostics
|
||||||
|
|
||||||
|
## Insight run sweeper (`[insight-sweeper]`)
|
||||||
|
|
||||||
|
The dashboard insight router runs stale-run recovery sweeps for `project_insight_runs` rows stuck in `pending`/`running` without a live controller owner.
|
||||||
|
|
||||||
|
- Recovery writes `terminalCause: "orphaned_active_run_recovered"` and lifecycle failure metadata (`failureClass: "non_retryable"`, `retryable: false`).
|
||||||
|
- Recovery appends both `warning` and `status_changed` events on `project_insight_run_events` with `metadata.recovery = "orphaned_active_run"`.
|
||||||
|
- `metadata.recoverySource` indicates where recovery occurred: `startup`, `periodic`, `drive_by`, or `manual`.
|
||||||
@@ -713,6 +713,120 @@ describe("InsightStore Run CRUD", () => {
|
|||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
|
describe("listStalePendingRuns", () => {
|
||||||
|
it("returns pending/running runs older than threshold", () => {
|
||||||
|
store.getDatabase().prepare(`
|
||||||
|
INSERT INTO project_insight_runs (id, projectId, trigger, status, summary, error, insightsCreated, insightsUpdated, inputMetadata, outputMetadata, createdAt, startedAt, completedAt)
|
||||||
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||||
|
`).run(
|
||||||
|
"INSR-OLD-PENDING",
|
||||||
|
"proj",
|
||||||
|
"manual",
|
||||||
|
"pending",
|
||||||
|
null,
|
||||||
|
null,
|
||||||
|
0,
|
||||||
|
0,
|
||||||
|
null,
|
||||||
|
null,
|
||||||
|
"2025-01-01T00:00:00.000Z",
|
||||||
|
null,
|
||||||
|
null,
|
||||||
|
);
|
||||||
|
|
||||||
|
store.getDatabase().prepare(`
|
||||||
|
INSERT INTO project_insight_runs (id, projectId, trigger, status, summary, error, insightsCreated, insightsUpdated, inputMetadata, outputMetadata, createdAt, startedAt, completedAt)
|
||||||
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||||
|
`).run(
|
||||||
|
"INSR-OLD-RUNNING",
|
||||||
|
"proj",
|
||||||
|
"schedule",
|
||||||
|
"running",
|
||||||
|
null,
|
||||||
|
null,
|
||||||
|
0,
|
||||||
|
0,
|
||||||
|
null,
|
||||||
|
null,
|
||||||
|
"2025-01-01T00:00:00.000Z",
|
||||||
|
"2025-01-01T12:00:00.000Z",
|
||||||
|
null,
|
||||||
|
);
|
||||||
|
|
||||||
|
const stale = store.listStalePendingRuns("2025-01-02T00:00:00.000Z");
|
||||||
|
expect(stale.map((run) => run.id)).toEqual(expect.arrayContaining(["INSR-OLD-PENDING", "INSR-OLD-RUNNING"]));
|
||||||
|
});
|
||||||
|
|
||||||
|
it("excludes terminal statuses", () => {
|
||||||
|
const terminal = store.createRun("proj", { trigger: "manual" });
|
||||||
|
store.updateRun(terminal.id, { status: "failed", error: "boom" });
|
||||||
|
|
||||||
|
const stale = store.listStalePendingRuns("9999-01-01T00:00:00.000Z");
|
||||||
|
expect(stale.some((run) => run.id === terminal.id)).toBe(false);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("honors projectId filter", () => {
|
||||||
|
const projectRun = store.createRun("proj-a", { trigger: "manual" });
|
||||||
|
store.createRun("proj-b", { trigger: "manual" });
|
||||||
|
|
||||||
|
const stale = store.listStalePendingRuns("9999-01-01T00:00:00.000Z", { projectId: "proj-a" });
|
||||||
|
expect(stale.map((run) => run.id)).toEqual([projectRun.id]);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("honors limit", () => {
|
||||||
|
for (let i = 0; i < 3; i++) {
|
||||||
|
store.createRun("proj", { trigger: "manual" });
|
||||||
|
}
|
||||||
|
|
||||||
|
const stale = store.listStalePendingRuns("9999-01-01T00:00:00.000Z", { limit: 2 });
|
||||||
|
expect(stale).toHaveLength(2);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("uses startedAt when present, otherwise createdAt", () => {
|
||||||
|
store.getDatabase().prepare(`
|
||||||
|
INSERT INTO project_insight_runs (id, projectId, trigger, status, summary, error, insightsCreated, insightsUpdated, inputMetadata, outputMetadata, createdAt, startedAt, completedAt)
|
||||||
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||||
|
`).run(
|
||||||
|
"INSR-STALE-CREATED",
|
||||||
|
"proj",
|
||||||
|
"manual",
|
||||||
|
"pending",
|
||||||
|
null,
|
||||||
|
null,
|
||||||
|
0,
|
||||||
|
0,
|
||||||
|
null,
|
||||||
|
null,
|
||||||
|
"2025-01-01T00:00:00.000Z",
|
||||||
|
null,
|
||||||
|
null,
|
||||||
|
);
|
||||||
|
|
||||||
|
store.getDatabase().prepare(`
|
||||||
|
INSERT INTO project_insight_runs (id, projectId, trigger, status, summary, error, insightsCreated, insightsUpdated, inputMetadata, outputMetadata, createdAt, startedAt, completedAt)
|
||||||
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||||
|
`).run(
|
||||||
|
"INSR-RECENT-START",
|
||||||
|
"proj",
|
||||||
|
"manual",
|
||||||
|
"running",
|
||||||
|
null,
|
||||||
|
null,
|
||||||
|
0,
|
||||||
|
0,
|
||||||
|
null,
|
||||||
|
null,
|
||||||
|
"2025-01-01T00:00:00.000Z",
|
||||||
|
"2025-01-03T00:00:00.000Z",
|
||||||
|
null,
|
||||||
|
);
|
||||||
|
|
||||||
|
const stale = store.listStalePendingRuns("2025-01-02T00:00:00.000Z");
|
||||||
|
expect(stale.map((run) => run.id)).toContain("INSR-STALE-CREATED");
|
||||||
|
expect(stale.map((run) => run.id)).not.toContain("INSR-RECENT-START");
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
describe("updateRun", () => {
|
describe("updateRun", () => {
|
||||||
it("updates mutable fields", () => {
|
it("updates mutable fields", () => {
|
||||||
const run = store.createRun("proj", { trigger: "manual" });
|
const run = store.createRun("proj", { trigger: "manual" });
|
||||||
|
|||||||
@@ -658,6 +658,33 @@ export class InsightStore extends EventEmitter<InsightStoreEvents> {
|
|||||||
return existingRow ? this.getRun(existingRow.id) : undefined;
|
return existingRow ? this.getRun(existingRow.id) : undefined;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
listStalePendingRuns(
|
||||||
|
olderThanIso: string,
|
||||||
|
options: {
|
||||||
|
projectId?: string;
|
||||||
|
limit?: number;
|
||||||
|
} = {},
|
||||||
|
): InsightRun[] {
|
||||||
|
const limit = Math.max(1, Math.floor(options.limit ?? 100));
|
||||||
|
const whereParts = ["status IN ('pending', 'running')", "COALESCE(startedAt, createdAt) <= ?"];
|
||||||
|
const params: (string | number)[] = [olderThanIso];
|
||||||
|
|
||||||
|
if (options.projectId) {
|
||||||
|
whereParts.push("projectId = ?");
|
||||||
|
params.push(options.projectId);
|
||||||
|
}
|
||||||
|
|
||||||
|
params.push(limit);
|
||||||
|
const rows = this.db.prepare(`
|
||||||
|
SELECT * FROM project_insight_runs
|
||||||
|
WHERE ${whereParts.join(" AND ")}
|
||||||
|
ORDER BY createdAt ASC, id ASC
|
||||||
|
LIMIT ?
|
||||||
|
`).all(...params) as Record<string, unknown>[];
|
||||||
|
|
||||||
|
return rows.map((row) => this.rowToRun(row));
|
||||||
|
}
|
||||||
|
|
||||||
createRunOrThrowConflict(projectId: string, input: InsightRunCreateInput): InsightRun {
|
createRunOrThrowConflict(projectId: string, input: InsightRunCreateInput): InsightRun {
|
||||||
const existing = this.findActiveRun(projectId, input.trigger);
|
const existing = this.findActiveRun(projectId, input.trigger);
|
||||||
if (existing) {
|
if (existing) {
|
||||||
|
|||||||
126
packages/dashboard/src/__tests__/insight-run-sweeper.test.ts
Normal file
126
packages/dashboard/src/__tests__/insight-run-sweeper.test.ts
Normal file
@@ -0,0 +1,126 @@
|
|||||||
|
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||||
|
import { InsightStore, createDatabase } from "@fusion/core";
|
||||||
|
import {
|
||||||
|
ORPHAN_GRACE_MS,
|
||||||
|
recoverOrphanedInsightRun,
|
||||||
|
startInsightRunSweeper,
|
||||||
|
sweepStaleInsightRuns,
|
||||||
|
} from "../insight-run-sweeper.js";
|
||||||
|
|
||||||
|
describe("insight-run-sweeper", () => {
|
||||||
|
let store: InsightStore;
|
||||||
|
let controllers: Map<string, AbortController>;
|
||||||
|
|
||||||
|
beforeEach(() => {
|
||||||
|
const db = createDatabase(":memory:");
|
||||||
|
store = new InsightStore(db);
|
||||||
|
controllers = new Map<string, AbortController>();
|
||||||
|
});
|
||||||
|
|
||||||
|
afterEach(() => {
|
||||||
|
vi.useRealTimers();
|
||||||
|
store.getDatabase().close();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("skips runs younger than graceMs", () => {
|
||||||
|
const run = store.createRun("proj", { trigger: "manual" });
|
||||||
|
const now = new Date(new Date(run.createdAt).getTime() + ORPHAN_GRACE_MS - 1);
|
||||||
|
|
||||||
|
const result = sweepStaleInsightRuns({
|
||||||
|
insightStore: store,
|
||||||
|
activeRunControllers: controllers,
|
||||||
|
now,
|
||||||
|
graceMs: ORPHAN_GRACE_MS,
|
||||||
|
source: "drive_by",
|
||||||
|
});
|
||||||
|
|
||||||
|
expect(result).toEqual({ scanned: 0, recovered: 0, skipped: 0 });
|
||||||
|
expect(store.getRun(run.id)?.status).toBe("pending");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("skips runs with active controllers", () => {
|
||||||
|
const run = store.createRun("proj", { trigger: "manual" });
|
||||||
|
controllers.set(run.id, new AbortController());
|
||||||
|
|
||||||
|
const result = sweepStaleInsightRuns({
|
||||||
|
insightStore: store,
|
||||||
|
activeRunControllers: controllers,
|
||||||
|
now: new Date(new Date(run.createdAt).getTime() + ORPHAN_GRACE_MS + 1000),
|
||||||
|
source: "startup",
|
||||||
|
});
|
||||||
|
|
||||||
|
expect(result).toEqual({ scanned: 1, recovered: 0, skipped: 1 });
|
||||||
|
expect(store.getRun(run.id)?.status).toBe("pending");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("recovers eligible runs and emits recovery metadata", () => {
|
||||||
|
const run = store.createRun("proj", { trigger: "manual" });
|
||||||
|
store.updateRun(run.id, {
|
||||||
|
status: "running",
|
||||||
|
startedAt: "2025-01-01T00:00:00.000Z",
|
||||||
|
});
|
||||||
|
|
||||||
|
const recoverResult = recoverOrphanedInsightRun({
|
||||||
|
insightStore: store,
|
||||||
|
run: store.getRun(run.id),
|
||||||
|
now: new Date("2025-01-01T01:00:00.000Z"),
|
||||||
|
activeRunControllers: controllers,
|
||||||
|
source: "periodic",
|
||||||
|
graceMs: ORPHAN_GRACE_MS,
|
||||||
|
});
|
||||||
|
|
||||||
|
expect(recoverResult).toEqual({ recovered: true });
|
||||||
|
|
||||||
|
const updated = store.getRun(run.id);
|
||||||
|
expect(updated?.status).toBe("failed");
|
||||||
|
expect(updated?.lifecycle.terminalCause).toBe("orphaned_active_run_recovered");
|
||||||
|
expect(updated?.lifecycle.failureClass).toBe("non_retryable");
|
||||||
|
expect(updated?.lifecycle.retryable).toBe(false);
|
||||||
|
|
||||||
|
const events = store.listRunEvents(run.id);
|
||||||
|
const warning = events.find((event) => event.type === "warning");
|
||||||
|
const statusChanged = events.find((event) => event.type === "status_changed");
|
||||||
|
|
||||||
|
expect(warning?.metadata?.recoverySource).toBe("periodic");
|
||||||
|
expect(statusChanged?.metadata?.recoverySource).toBe("periodic");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("returns accurate scanned/recovered/skipped counts", () => {
|
||||||
|
const recoverable = store.createRun("proj", { trigger: "manual" });
|
||||||
|
const withController = store.createRun("proj", { trigger: "schedule" });
|
||||||
|
controllers.set(withController.id, new AbortController());
|
||||||
|
|
||||||
|
const result = sweepStaleInsightRuns({
|
||||||
|
insightStore: store,
|
||||||
|
activeRunControllers: controllers,
|
||||||
|
now: new Date(new Date(recoverable.createdAt).getTime() + ORPHAN_GRACE_MS + 10_000),
|
||||||
|
source: "drive_by",
|
||||||
|
});
|
||||||
|
|
||||||
|
expect(result).toEqual({ scanned: 2, recovered: 1, skipped: 1 });
|
||||||
|
});
|
||||||
|
|
||||||
|
it("runs periodic sweeps and dispose stops interval", () => {
|
||||||
|
vi.useFakeTimers();
|
||||||
|
|
||||||
|
const first = store.createRun("proj", { trigger: "manual" });
|
||||||
|
const logger = { warn: vi.fn() };
|
||||||
|
const sweeper = startInsightRunSweeper({
|
||||||
|
insightStore: store,
|
||||||
|
activeRunControllers: controllers,
|
||||||
|
intervalMs: 1_000,
|
||||||
|
graceMs: 0,
|
||||||
|
logger,
|
||||||
|
});
|
||||||
|
|
||||||
|
vi.advanceTimersByTime(1_000);
|
||||||
|
expect(store.getRun(first.id)?.status).toBe("failed");
|
||||||
|
|
||||||
|
const second = store.createRun("proj", { trigger: "manual" });
|
||||||
|
sweeper.dispose();
|
||||||
|
vi.advanceTimersByTime(2_000);
|
||||||
|
|
||||||
|
expect(store.getRun(second.id)?.status).toBe("pending");
|
||||||
|
expect(logger.warn).not.toHaveBeenCalled();
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -1,4 +1,5 @@
|
|||||||
import { describe, it, expect, vi, beforeEach, afterEach } from "vitest";
|
import { describe, it, expect, vi, beforeEach, afterEach } from "vitest";
|
||||||
|
import express from "express";
|
||||||
import { mkdtempSync } from "node:fs";
|
import { mkdtempSync } from "node:fs";
|
||||||
import { rm } from "node:fs/promises";
|
import { rm } from "node:fs/promises";
|
||||||
import { tmpdir } from "node:os";
|
import { tmpdir } from "node:os";
|
||||||
@@ -8,6 +9,8 @@ import { TaskStore as TaskStoreClass } from "@fusion/core";
|
|||||||
import * as coreModule from "@fusion/core";
|
import * as coreModule from "@fusion/core";
|
||||||
import { request } from "../test-request.js";
|
import { request } from "../test-request.js";
|
||||||
import { createServer } from "../server.js";
|
import { createServer } from "../server.js";
|
||||||
|
import { createInsightsRouter } from "../insights-routes.js";
|
||||||
|
import { DEFAULT_SWEEP_INTERVAL_MS } from "../insight-run-sweeper.js";
|
||||||
|
|
||||||
const piMocks = vi.hoisted(() => ({
|
const piMocks = vi.hoisted(() => ({
|
||||||
createFnAgent: vi.fn(),
|
createFnAgent: vi.fn(),
|
||||||
@@ -68,6 +71,7 @@ vi.mock("../project-store-resolver.js", async () => {
|
|||||||
|
|
||||||
describe("Insights routes", () => {
|
describe("Insights routes", () => {
|
||||||
let rootA: string;
|
let rootA: string;
|
||||||
|
const disposableRouters: Array<{ __disposeSweeper?: () => void }> = [];
|
||||||
let rootB: string;
|
let rootB: string;
|
||||||
let storeA: TaskStore;
|
let storeA: TaskStore;
|
||||||
let storeB: TaskStore;
|
let storeB: TaskStore;
|
||||||
@@ -80,6 +84,15 @@ describe("Insights routes", () => {
|
|||||||
const parseResponseSpy = vi.spyOn(coreModule, "parseInsightExtractionResponse");
|
const parseResponseSpy = vi.spyOn(coreModule, "parseInsightExtractionResponse");
|
||||||
const mergeInsightsSpy = vi.spyOn(coreModule, "mergeInsights");
|
const mergeInsightsSpy = vi.spyOn(coreModule, "mergeInsights");
|
||||||
|
|
||||||
|
function createInsightsOnlyApp(store: TaskStore) {
|
||||||
|
const app = express();
|
||||||
|
app.use(express.json());
|
||||||
|
const router = createInsightsRouter(store) as ReturnType<typeof createInsightsRouter> & { __disposeSweeper?: () => void };
|
||||||
|
disposableRouters.push(router);
|
||||||
|
app.use("/api/insights", router);
|
||||||
|
return app;
|
||||||
|
}
|
||||||
|
|
||||||
beforeEach(async () => {
|
beforeEach(async () => {
|
||||||
vi.clearAllMocks();
|
vi.clearAllMocks();
|
||||||
|
|
||||||
@@ -120,6 +133,10 @@ describe("Insights routes", () => {
|
|||||||
});
|
});
|
||||||
|
|
||||||
afterEach(async () => {
|
afterEach(async () => {
|
||||||
|
vi.useRealTimers();
|
||||||
|
while (disposableRouters.length > 0) {
|
||||||
|
disposableRouters.pop()?.__disposeSweeper?.();
|
||||||
|
}
|
||||||
try {
|
try {
|
||||||
storeA.close();
|
storeA.close();
|
||||||
} catch {
|
} catch {
|
||||||
@@ -194,6 +211,96 @@ describe("Insights routes", () => {
|
|||||||
expect((paged.body as { runs: unknown[] }).runs).toHaveLength(1);
|
expect((paged.body as { runs: unknown[] }).runs).toHaveLength(1);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it("runs startup sweep when creating insights router", async () => {
|
||||||
|
const oneHourAgo = new Date(Date.now() - 60 * 60 * 1000).toISOString();
|
||||||
|
storeA.getDatabase().prepare(`
|
||||||
|
INSERT INTO project_insight_runs (id, projectId, trigger, status, summary, error, insightsCreated, insightsUpdated, inputMetadata, outputMetadata, createdAt, startedAt, completedAt)
|
||||||
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||||
|
`).run(
|
||||||
|
"INSR-STALE-STARTUP",
|
||||||
|
"",
|
||||||
|
"manual",
|
||||||
|
"pending",
|
||||||
|
null,
|
||||||
|
null,
|
||||||
|
0,
|
||||||
|
0,
|
||||||
|
null,
|
||||||
|
null,
|
||||||
|
oneHourAgo,
|
||||||
|
null,
|
||||||
|
null,
|
||||||
|
);
|
||||||
|
|
||||||
|
createInsightsOnlyApp(storeA);
|
||||||
|
|
||||||
|
const run = storeA.getInsightStore().getRun("INSR-STALE-STARTUP");
|
||||||
|
expect(run?.status).toBe("failed");
|
||||||
|
expect(run?.lifecycle.terminalCause).toBe("orphaned_active_run_recovered");
|
||||||
|
const events = storeA.getInsightStore().listRunEvents("INSR-STALE-STARTUP");
|
||||||
|
expect(events.some((event) => event.metadata?.recoverySource === "startup")).toBe(true);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("runs drive-by sweep on GET /api/insights/runs", async () => {
|
||||||
|
const oneHourAgo = new Date(Date.now() - 60 * 60 * 1000).toISOString();
|
||||||
|
storeA.getDatabase().prepare(`
|
||||||
|
INSERT INTO project_insight_runs (id, projectId, trigger, status, summary, error, insightsCreated, insightsUpdated, inputMetadata, outputMetadata, createdAt, startedAt, completedAt)
|
||||||
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||||
|
`).run(
|
||||||
|
"INSR-STALE-DRIVEBY",
|
||||||
|
"",
|
||||||
|
"manual",
|
||||||
|
"pending",
|
||||||
|
null,
|
||||||
|
null,
|
||||||
|
0,
|
||||||
|
0,
|
||||||
|
null,
|
||||||
|
null,
|
||||||
|
oneHourAgo,
|
||||||
|
null,
|
||||||
|
null,
|
||||||
|
);
|
||||||
|
|
||||||
|
const insightsApp = createInsightsOnlyApp(storeA);
|
||||||
|
const response = await request(insightsApp, "GET", "/api/insights/runs");
|
||||||
|
expect(response.status).toBe(200);
|
||||||
|
|
||||||
|
const run = storeA.getInsightStore().getRun("INSR-STALE-DRIVEBY");
|
||||||
|
expect(run?.status).toBe("failed");
|
||||||
|
const events = storeA.getInsightStore().listRunEvents("INSR-STALE-DRIVEBY");
|
||||||
|
expect(events.some((event) => event.metadata?.recoverySource === "drive_by")).toBe(true);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("runs periodic sweep and recover later stale rows", async () => {
|
||||||
|
vi.useFakeTimers();
|
||||||
|
|
||||||
|
const insightsApp = createInsightsOnlyApp(storeA);
|
||||||
|
const first = storeA.getInsightStore().createRun("", { trigger: "manual" });
|
||||||
|
storeA.getDatabase().prepare("UPDATE project_insight_runs SET createdAt = ? WHERE id = ?").run(
|
||||||
|
new Date(Date.now() - 60 * 60 * 1000).toISOString(),
|
||||||
|
first.id,
|
||||||
|
);
|
||||||
|
|
||||||
|
vi.advanceTimersByTime(DEFAULT_SWEEP_INTERVAL_MS + 100);
|
||||||
|
|
||||||
|
expect(storeA.getInsightStore().getRun(first.id)?.status).toBe("failed");
|
||||||
|
|
||||||
|
const second = storeA.getInsightStore().createRun("", { trigger: "manual" });
|
||||||
|
storeA.getDatabase().prepare("UPDATE project_insight_runs SET createdAt = ? WHERE id = ?").run(
|
||||||
|
new Date(Date.now() - 60 * 60 * 1000).toISOString(),
|
||||||
|
second.id,
|
||||||
|
);
|
||||||
|
|
||||||
|
vi.advanceTimersByTime(DEFAULT_SWEEP_INTERVAL_MS + 100);
|
||||||
|
|
||||||
|
expect(storeA.getInsightStore().getRun(second.id)?.status).toBe("failed");
|
||||||
|
const events = storeA.getInsightStore().listRunEvents(second.id);
|
||||||
|
expect(events.some((event) => event.metadata?.recoverySource === "periodic")).toBe(true);
|
||||||
|
|
||||||
|
expect(insightsApp).toBeTruthy();
|
||||||
|
});
|
||||||
|
|
||||||
it("GET /api/insights and /api/insights/runs resolve projectId-scoped stores", async () => {
|
it("GET /api/insights and /api/insights/runs resolve projectId-scoped stores", async () => {
|
||||||
const runA = storeA.getInsightStore().createRun("", { trigger: "manual" });
|
const runA = storeA.getInsightStore().createRun("", { trigger: "manual" });
|
||||||
const runB = storeB.getInsightStore().createRun("", { trigger: "manual" });
|
const runB = storeB.getInsightStore().createRun("", { trigger: "manual" });
|
||||||
|
|||||||
178
packages/dashboard/src/insight-run-sweeper.ts
Normal file
178
packages/dashboard/src/insight-run-sweeper.ts
Normal file
@@ -0,0 +1,178 @@
|
|||||||
|
import type { InsightRun, InsightStore } from "@fusion/core";
|
||||||
|
|
||||||
|
export const ORPHAN_GRACE_MS = 30_000;
|
||||||
|
export const DEFAULT_SWEEP_INTERVAL_MS = 5 * 60_000;
|
||||||
|
|
||||||
|
type RecoverySource = "startup" | "periodic" | "drive_by" | "manual";
|
||||||
|
|
||||||
|
type RecoverParams = {
|
||||||
|
insightStore: InsightStore;
|
||||||
|
run: InsightRun | null | undefined;
|
||||||
|
now: Date;
|
||||||
|
activeRunControllers: Map<string, AbortController>;
|
||||||
|
graceMs?: number;
|
||||||
|
source?: RecoverySource;
|
||||||
|
};
|
||||||
|
|
||||||
|
function getRunAgeMs(run: Pick<InsightRun, "startedAt" | "createdAt">, nowMs: number): number {
|
||||||
|
const anchor = run.startedAt ?? run.createdAt;
|
||||||
|
const anchorMs = Date.parse(anchor);
|
||||||
|
if (!Number.isFinite(anchorMs)) return 0;
|
||||||
|
return Math.max(0, nowMs - anchorMs);
|
||||||
|
}
|
||||||
|
|
||||||
|
export function recoverOrphanedInsightRun(params: RecoverParams): { recovered: boolean; reason?: string } {
|
||||||
|
const {
|
||||||
|
insightStore,
|
||||||
|
run,
|
||||||
|
now,
|
||||||
|
activeRunControllers,
|
||||||
|
graceMs = ORPHAN_GRACE_MS,
|
||||||
|
source = "manual",
|
||||||
|
} = params;
|
||||||
|
|
||||||
|
if (!run || !["pending", "running"].includes(run.status)) {
|
||||||
|
return { recovered: false, reason: "not_active_status" };
|
||||||
|
}
|
||||||
|
|
||||||
|
if (activeRunControllers.has(run.id)) {
|
||||||
|
return { recovered: false, reason: "has_live_controller" };
|
||||||
|
}
|
||||||
|
|
||||||
|
const ageMs = getRunAgeMs(run, now.getTime());
|
||||||
|
if (ageMs <= graceMs) {
|
||||||
|
return { recovered: false, reason: "within_grace_window" };
|
||||||
|
}
|
||||||
|
|
||||||
|
const nowIso = now.toISOString();
|
||||||
|
insightStore.appendRunEvent(run.id, {
|
||||||
|
type: "warning",
|
||||||
|
status: run.status,
|
||||||
|
classification: "non_retryable",
|
||||||
|
message: `Recovered orphaned active run after ${ageMs}ms without controller ownership`,
|
||||||
|
metadata: {
|
||||||
|
recovery: "orphaned_active_run",
|
||||||
|
ageMs,
|
||||||
|
graceMs,
|
||||||
|
hadController: false,
|
||||||
|
anchorTimestamp: run.startedAt ?? run.createdAt,
|
||||||
|
recoverySource: source,
|
||||||
|
},
|
||||||
|
});
|
||||||
|
|
||||||
|
const failed = insightStore.updateRun(run.id, {
|
||||||
|
status: "failed",
|
||||||
|
summary: "Recovered orphaned run",
|
||||||
|
error: "Run was marked active but had no live controller after grace period",
|
||||||
|
completedAt: nowIso,
|
||||||
|
lifecycle: {
|
||||||
|
...run.lifecycle,
|
||||||
|
terminalReason: "failed",
|
||||||
|
terminalCause: "orphaned_active_run_recovered",
|
||||||
|
failureClass: "non_retryable",
|
||||||
|
retryable: false,
|
||||||
|
},
|
||||||
|
});
|
||||||
|
|
||||||
|
if (!failed) {
|
||||||
|
return { recovered: false, reason: "update_failed" };
|
||||||
|
}
|
||||||
|
|
||||||
|
insightStore.appendRunEvent(run.id, {
|
||||||
|
type: "status_changed",
|
||||||
|
status: "failed",
|
||||||
|
classification: "non_retryable",
|
||||||
|
message: "Run marked failed after orphaned active-run recovery",
|
||||||
|
metadata: {
|
||||||
|
recovery: "orphaned_active_run",
|
||||||
|
recoverySource: source,
|
||||||
|
},
|
||||||
|
});
|
||||||
|
|
||||||
|
return { recovered: true };
|
||||||
|
}
|
||||||
|
|
||||||
|
export function sweepStaleInsightRuns(params: {
|
||||||
|
insightStore: InsightStore;
|
||||||
|
activeRunControllers: Map<string, AbortController>;
|
||||||
|
now?: Date;
|
||||||
|
graceMs?: number;
|
||||||
|
source: RecoverySource;
|
||||||
|
}): { scanned: number; recovered: number; skipped: number } {
|
||||||
|
const {
|
||||||
|
insightStore,
|
||||||
|
activeRunControllers,
|
||||||
|
now = new Date(),
|
||||||
|
graceMs = ORPHAN_GRACE_MS,
|
||||||
|
source,
|
||||||
|
} = params;
|
||||||
|
|
||||||
|
const thresholdIso = new Date(now.getTime() - graceMs).toISOString();
|
||||||
|
const staleRuns = insightStore.listStalePendingRuns(thresholdIso);
|
||||||
|
|
||||||
|
let recovered = 0;
|
||||||
|
let skipped = 0;
|
||||||
|
|
||||||
|
for (const run of staleRuns) {
|
||||||
|
if (activeRunControllers.has(run.id)) {
|
||||||
|
skipped += 1;
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
|
||||||
|
const result = recoverOrphanedInsightRun({
|
||||||
|
insightStore,
|
||||||
|
run,
|
||||||
|
now,
|
||||||
|
activeRunControllers,
|
||||||
|
graceMs,
|
||||||
|
source,
|
||||||
|
});
|
||||||
|
|
||||||
|
if (result.recovered) {
|
||||||
|
recovered += 1;
|
||||||
|
} else {
|
||||||
|
skipped += 1;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return {
|
||||||
|
scanned: staleRuns.length,
|
||||||
|
recovered,
|
||||||
|
skipped,
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
export function startInsightRunSweeper(params: {
|
||||||
|
insightStore: InsightStore;
|
||||||
|
activeRunControllers: Map<string, AbortController>;
|
||||||
|
intervalMs?: number;
|
||||||
|
graceMs?: number;
|
||||||
|
logger?: Pick<Console, "warn">;
|
||||||
|
}): { dispose: () => void } {
|
||||||
|
const {
|
||||||
|
insightStore,
|
||||||
|
activeRunControllers,
|
||||||
|
intervalMs = DEFAULT_SWEEP_INTERVAL_MS,
|
||||||
|
graceMs = ORPHAN_GRACE_MS,
|
||||||
|
logger,
|
||||||
|
} = params;
|
||||||
|
|
||||||
|
const timer = setInterval(() => {
|
||||||
|
try {
|
||||||
|
sweepStaleInsightRuns({
|
||||||
|
insightStore,
|
||||||
|
activeRunControllers,
|
||||||
|
graceMs,
|
||||||
|
source: "periodic",
|
||||||
|
});
|
||||||
|
} catch (error) {
|
||||||
|
logger?.warn?.("[insight-sweeper] periodic sweep failed", error);
|
||||||
|
}
|
||||||
|
}, intervalMs);
|
||||||
|
|
||||||
|
timer.unref?.();
|
||||||
|
|
||||||
|
return {
|
||||||
|
dispose: () => clearInterval(timer),
|
||||||
|
};
|
||||||
|
}
|
||||||
@@ -34,6 +34,13 @@ import {
|
|||||||
badRequest,
|
badRequest,
|
||||||
notFound,
|
notFound,
|
||||||
} from "./api-error.js";
|
} from "./api-error.js";
|
||||||
|
import {
|
||||||
|
DEFAULT_SWEEP_INTERVAL_MS,
|
||||||
|
ORPHAN_GRACE_MS,
|
||||||
|
recoverOrphanedInsightRun,
|
||||||
|
startInsightRunSweeper,
|
||||||
|
sweepStaleInsightRuns,
|
||||||
|
} from "./insight-run-sweeper.js";
|
||||||
import { createFnAgent, promptWithFallback } from "@fusion/engine";
|
import { createFnAgent, promptWithFallback } from "@fusion/engine";
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -88,79 +95,21 @@ const INSIGHT_CATEGORY_BY_MEMORY_CATEGORY: Record<MemoryInsightCategory, Insight
|
|||||||
};
|
};
|
||||||
|
|
||||||
const activeRunControllers = new Map<string, AbortController>();
|
const activeRunControllers = new Map<string, AbortController>();
|
||||||
const ORPHAN_GRACE_MS = 30_000;
|
|
||||||
|
|
||||||
function getRunAgeMs(run: { startedAt: string | null; createdAt: string }, nowMs: number): number {
|
|
||||||
const anchor = run.startedAt ?? run.createdAt;
|
|
||||||
const anchorMs = Date.parse(anchor);
|
|
||||||
if (!Number.isFinite(anchorMs)) return 0;
|
|
||||||
return Math.max(0, nowMs - anchorMs);
|
|
||||||
}
|
|
||||||
|
|
||||||
function maybeRecoverOrphanedActiveRun(params: {
|
function maybeRecoverOrphanedActiveRun(params: {
|
||||||
insightStore: InsightStore;
|
insightStore: InsightStore;
|
||||||
run: ReturnType<InsightStore["getRun"]>;
|
run: ReturnType<InsightStore["getRun"]>;
|
||||||
trigger: InsightRunTrigger;
|
|
||||||
now: Date;
|
now: Date;
|
||||||
}): boolean {
|
}): boolean {
|
||||||
const { insightStore, run, trigger, now } = params;
|
const { insightStore, run, now } = params;
|
||||||
if (!run || !["pending", "running"].includes(run.status)) {
|
return recoverOrphanedInsightRun({
|
||||||
return false;
|
insightStore,
|
||||||
}
|
run,
|
||||||
|
now,
|
||||||
if (activeRunControllers.has(run.id)) {
|
activeRunControllers,
|
||||||
return false;
|
source: "manual",
|
||||||
}
|
graceMs: ORPHAN_GRACE_MS,
|
||||||
|
}).recovered;
|
||||||
const ageMs = getRunAgeMs(run, now.getTime());
|
|
||||||
if (ageMs <= ORPHAN_GRACE_MS) {
|
|
||||||
return false;
|
|
||||||
}
|
|
||||||
|
|
||||||
const nowIso = now.toISOString();
|
|
||||||
insightStore.appendRunEvent(run.id, {
|
|
||||||
type: "warning",
|
|
||||||
status: run.status,
|
|
||||||
classification: "non_retryable",
|
|
||||||
message: `Recovered orphaned active run after ${ageMs}ms without controller ownership`,
|
|
||||||
metadata: {
|
|
||||||
recovery: "orphaned_active_run",
|
|
||||||
trigger,
|
|
||||||
ageMs,
|
|
||||||
graceMs: ORPHAN_GRACE_MS,
|
|
||||||
hadController: false,
|
|
||||||
anchorTimestamp: run.startedAt ?? run.createdAt,
|
|
||||||
},
|
|
||||||
});
|
|
||||||
|
|
||||||
const failed = insightStore.updateRun(run.id, {
|
|
||||||
status: "failed",
|
|
||||||
summary: "Recovered orphaned run",
|
|
||||||
error: "Run was marked active but had no live controller after grace period",
|
|
||||||
completedAt: nowIso,
|
|
||||||
lifecycle: {
|
|
||||||
...run.lifecycle,
|
|
||||||
terminalReason: "failed",
|
|
||||||
terminalCause: "orphaned_active_run_recovered",
|
|
||||||
failureClass: "non_retryable",
|
|
||||||
retryable: false,
|
|
||||||
},
|
|
||||||
});
|
|
||||||
|
|
||||||
if (failed) {
|
|
||||||
insightStore.appendRunEvent(run.id, {
|
|
||||||
type: "status_changed",
|
|
||||||
status: "failed",
|
|
||||||
classification: "non_retryable",
|
|
||||||
message: "Run marked failed after orphaned active-run recovery",
|
|
||||||
metadata: {
|
|
||||||
recovery: "orphaned_active_run",
|
|
||||||
},
|
|
||||||
});
|
|
||||||
return true;
|
|
||||||
}
|
|
||||||
|
|
||||||
return false;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
async function withAbort<T>(signal: AbortSignal, task: Promise<T>): Promise<T> {
|
async function withAbort<T>(signal: AbortSignal, task: Promise<T>): Promise<T> {
|
||||||
@@ -299,6 +248,32 @@ export function createInsightsRouter(store: TaskStore): Router {
|
|||||||
const router = Router();
|
const router = Router();
|
||||||
const requestContext = new AsyncLocalStorage<TaskStore>();
|
const requestContext = new AsyncLocalStorage<TaskStore>();
|
||||||
|
|
||||||
|
const rootInsightStore = typeof (store as { getInsightStore?: () => InsightStore }).getInsightStore === "function"
|
||||||
|
? (store as { getInsightStore: () => InsightStore }).getInsightStore()
|
||||||
|
: undefined;
|
||||||
|
|
||||||
|
if (rootInsightStore) {
|
||||||
|
try {
|
||||||
|
sweepStaleInsightRuns({
|
||||||
|
insightStore: rootInsightStore,
|
||||||
|
activeRunControllers,
|
||||||
|
graceMs: ORPHAN_GRACE_MS,
|
||||||
|
source: "startup",
|
||||||
|
});
|
||||||
|
} catch (error) {
|
||||||
|
console.warn("[insight-sweeper] startup sweep failed", error);
|
||||||
|
}
|
||||||
|
|
||||||
|
const { dispose: disposeSweeper } = startInsightRunSweeper({
|
||||||
|
insightStore: rootInsightStore,
|
||||||
|
activeRunControllers,
|
||||||
|
intervalMs: DEFAULT_SWEEP_INTERVAL_MS,
|
||||||
|
graceMs: ORPHAN_GRACE_MS,
|
||||||
|
logger: console,
|
||||||
|
});
|
||||||
|
(router as Router & { __disposeSweeper?: () => void }).__disposeSweeper = disposeSweeper;
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Middleware to capture the appropriate store for this request.
|
* Middleware to capture the appropriate store for this request.
|
||||||
* Uses projectId from query/body to get the scoped store if provided,
|
* Uses projectId from query/body to get the scoped store if provided,
|
||||||
@@ -426,7 +401,6 @@ export function createInsightsRouter(store: TaskStore): Router {
|
|||||||
maybeRecoverOrphanedActiveRun({
|
maybeRecoverOrphanedActiveRun({
|
||||||
insightStore,
|
insightStore,
|
||||||
run: existingActiveRun,
|
run: existingActiveRun,
|
||||||
trigger,
|
|
||||||
now: new Date(),
|
now: new Date(),
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
@@ -514,6 +488,17 @@ export function createInsightsRouter(store: TaskStore): Router {
|
|||||||
options.offset = offset;
|
options.offset = offset;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
try {
|
||||||
|
sweepStaleInsightRuns({
|
||||||
|
insightStore: store,
|
||||||
|
activeRunControllers,
|
||||||
|
graceMs: ORPHAN_GRACE_MS,
|
||||||
|
source: "drive_by",
|
||||||
|
});
|
||||||
|
} catch (error) {
|
||||||
|
console.warn("[insight-sweeper] drive-by sweep failed", error);
|
||||||
|
}
|
||||||
|
|
||||||
const runs = store.listRuns(options);
|
const runs = store.listRuns(options);
|
||||||
res.json({ runs });
|
res.json({ runs });
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
@@ -527,6 +512,18 @@ export function createInsightsRouter(store: TaskStore): Router {
|
|||||||
try {
|
try {
|
||||||
const id = String(req.params.id);
|
const id = String(req.params.id);
|
||||||
const store = getInsightStore();
|
const store = getInsightStore();
|
||||||
|
|
||||||
|
try {
|
||||||
|
sweepStaleInsightRuns({
|
||||||
|
insightStore: store,
|
||||||
|
activeRunControllers,
|
||||||
|
graceMs: ORPHAN_GRACE_MS,
|
||||||
|
source: "drive_by",
|
||||||
|
});
|
||||||
|
} catch (error) {
|
||||||
|
console.warn("[insight-sweeper] drive-by sweep failed", error);
|
||||||
|
}
|
||||||
|
|
||||||
const run = store.getRun(id);
|
const run = store.getRun(id);
|
||||||
if (!run) {
|
if (!run) {
|
||||||
throw notFound(`Run not found: ${id}`);
|
throw notFound(`Run not found: ${id}`);
|
||||||
|
|||||||
Reference in New Issue
Block a user