fix(engine): bound spec-drift reconciliation so boot cannot exhaust the connection pool

Fusion wedged on "starting" and never brought the engine up. The dashboard bound
the migration holding server on 4040, then every query behind it failed with
"sorry, too many clients already", so the card never progressed and the
supervisor crash-looped.

Root cause: each spec-drift reconcile costs a DEDICATED PostgreSQL connection.
persist -> appendSpecDriftReport -> withPlanningLifecycleLock opens its own
postgres(directUrl, { max: 1 }) session, because the planning advisory lock is
session-scoped and deliberately fences a stale report against a newer plan.

enqueue() released every id straight into its own microtask, and project-engine
enqueues every task at runtime-boundary setup (listTasks includeArchived). On a
1,082-task project that opened ~1,082 lock sessions simultaneously against
max_connections = 500. The cluster saturated ~25s into boot and stayed saturated.

The flat 1s retry then made it self-sustaining rather than transient: once
saturated, every task failed for the same shared reason and re-armed in lockstep
once per second, re-opening the whole fleet of sessions and pinning the very
resource it was waiting on. Measured 4,777 lock sessions in 17 seconds.

Fix, contained to the reconciler — the advisory lock and its fencing semantics
are load-bearing and unchanged:
- concurrency bound (maxConcurrent, default 4) drained by a fair
  insertion-ordered pump, so fan-out can no longer exceed a known connection cost
- per-task in-flight dedupe; two passes on one task would contend on that task's
  own advisory lock while holding two connections
- exponential backoff with jitter capped at 60s, and retries re-enter through
  enqueue so a retry storm is bounded by the same limit as a first pass

Verified against the real 1,082-task project: connections stay flat at 3-10
across a 70s boot that previously reached 1,109 and saturated, and the engine
boots through to executing tasks and shuts down cleanly.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
gsxdsm
2026-08-10 16:44:19 -07:00
parent 67cdb33a92
commit 55fd20e5f8
3 changed files with 181 additions and 9 deletions

View File

@@ -0,0 +1,7 @@
---
"@runfusion/fusion": patch
---
summary: Fix Fusion hanging on startup — spec-drift reconciliation exhausted the database connection pool.
category: fix
dev: `SpecDriftReconciler.enqueue` released every id into its own microtask, and each reconcile costs a DEDICATED PostgreSQL connection (`appendSpecDriftReport` -> `withPlanningLifecycleLock` opens its own session-scoped `max: 1` client). `project-engine.ts` enqueues every task at runtime-boundary setup, so a 1,082-task project opened ~1,082 lock sessions at once against `max_connections = 500`, saturating the cluster ~25s into boot; every later query then failed with "sorry, too many clients already" and the dashboard wedged on "starting" behind the migration holding server. The flat 1s retry made it self-sustaining (measured 4,777 lock sessions in 17s). Adds a concurrency bound (`maxConcurrent`, default 4) with a fair insertion-ordered pump, per-task in-flight dedupe, and exponential backoff with jitter capped at 60s; retries re-enter through `enqueue` so they are bounded too. Verified against the real project: connections stay flat at 3-10 across a 70s boot that previously hit 1,109 and saturated.

View File

@@ -112,6 +112,90 @@ describe("SpecDriftReconciler", () => {
await expect(new SpecDriftReconciler(createStoreSpecDriftRepository(unavailableStore)).reconcile("FN-UNAVAILABLE")).resolves.toMatchObject({ alignment: "unavailable" }); await expect(new SpecDriftReconciler(createStoreSpecDriftRepository(unavailableStore)).reconcile("FN-UNAVAILABLE")).resolves.toMatchObject({ alignment: "unavailable" });
}); });
/*
FNXC:SpecDrift 2026-08-10-18:32 (connection-exhaustion incident):
Each reconcile costs a DEDICATED PostgreSQL connection — `appendSpecDriftReport` takes the
session-scoped planning advisory lock, which opens its own `max: 1` session. So concurrency here is
a connection count, not just a scheduling detail.
The regression: `project-engine.ts` enqueues every task at runtime-boundary setup
(`listTasks({ includeArchived: true })`) and `enqueue` released each id into its own microtask. On a
1,082-task project that opened ~1,082 lock sessions at once against `max_connections = 500`; the
cluster saturated ~25s into boot, every later query failed with "sorry, too many clients already",
and Fusion wedged on "starting" behind the migration holding server. Measured 4,777 lock sessions in
17 seconds. This asserts the bound that makes a large project survive boot.
*/
it("bounds concurrent reconciles so a whole-project enqueue cannot exhaust connections", async () => {
let active = 0;
let peak = 0;
let completed = 0;
let release: (() => void) | undefined;
const gate = new Promise<void>((resolve) => { release = resolve; });
const reconciler = new SpecDriftReconciler({
snapshot: async () => ({ latestLock: lock, currentPlan: evidence, approvedPlanFingerprint: "approved" }),
persist: async () => {
active += 1;
peak = Math.max(peak, active);
await gate;
active -= 1;
completed += 1;
},
}, { maxConcurrent: 4 });
for (let i = 0; i < 200; i += 1) reconciler.enqueue(`FN-${i}`);
// Let every queued id get a chance to start before releasing the gate.
await new Promise<void>((resolve) => setTimeout(resolve, 0));
expect(peak).toBeLessThanOrEqual(4);
release?.();
await vi.waitFor(() => expect(completed).toBe(200), { timeout: 5_000 });
// Bounded, but still drains everything — throughput is preserved, only fan-out is capped.
expect(peak).toBeLessThanOrEqual(4);
});
it("does not run the same task twice concurrently", async () => {
// Two sessions on one task would contend on its own advisory lock while holding two connections.
let active = 0;
let peak = 0;
let release: (() => void) | undefined;
const gate = new Promise<void>((resolve) => { release = resolve; });
const reconciler = new SpecDriftReconciler({
snapshot: async () => ({ latestLock: lock, currentPlan: evidence, approvedPlanFingerprint: "approved" }),
persist: async () => { active += 1; peak = Math.max(peak, active); await gate; active -= 1; },
}, { maxConcurrent: 4 });
reconciler.enqueue("FN-SAME");
reconciler.enqueue("FN-SAME");
reconciler.enqueue("FN-SAME");
await new Promise<void>((resolve) => setTimeout(resolve, 0));
expect(peak).toBe(1);
release?.();
});
it("backs a persistent outage off exponentially instead of re-firing every second", async () => {
vi.useFakeTimers();
let attempts = 0;
const reconciler = new SpecDriftReconciler({
snapshot: async () => ({ latestLock: lock, currentPlan: evidence, approvedPlanFingerprint: "approved" }),
persist: async () => { attempts += 1; throw new Error("persistent outage"); },
});
await expect(reconciler.reconcile("FN-BACKOFF")).rejects.toThrow("persistent outage");
expect(attempts).toBe(1);
// First retry lands inside the base window (jittered to 0.5-1x).
await vi.advanceTimersByTimeAsync(1_000);
expect(attempts).toBe(2);
// The second must NOT fire at another flat 1s — that lockstep re-fire is what kept the
// connection pool exhausted once every task started failing for the same shared reason.
await vi.advanceTimersByTimeAsync(1_000);
expect(attempts).toBe(2);
await vi.advanceTimersByTimeAsync(2_000);
expect(attempts).toBe(3);
vi.useRealTimers();
});
it("does not leak queued writes after stop", async () => { it("does not leak queued writes after stop", async () => {
const reconciler = new SpecDriftReconciler({ snapshot: async () => ({ latestLock: lock, currentPlan: evidence }), persist: async () => { throw new Error("must not write"); } }); const reconciler = new SpecDriftReconciler({ snapshot: async () => ({ latestLock: lock, currentPlan: evidence }), persist: async () => { throw new Error("must not write"); } });
reconciler.enqueue("FN-QUEUED"); reconciler.enqueue("FN-QUEUED");

View File

@@ -13,7 +13,29 @@ export interface SpecDriftRepository {
onPersisted?(taskId: string, report: DriftReport): Promise<void>; onPersisted?(taskId: string, report: DriftReport): Promise<void>;
} }
const RETRY_DELAY_MS = 1_000; /*
FNXC:SpecDrift 2026-08-10-18:32 (connection-exhaustion incident):
Reconciliation is CONCURRENCY-BOUNDED because each run costs a dedicated PostgreSQL connection, not a
pooled one. `persist` -> `appendSpecDriftReport` -> `withPlanningLifecycleLock` opens its own
`postgres(directUrl, { max: 1 })` session, because the planning advisory lock is session-scoped and
deliberately fences a stale report against a newer plan (see `appendSpecDriftReportWhilePlanningLocked`).
Unbounded fan-out therefore converts directly into unbounded connections. At runtime boundary setup
`project-engine.ts` enqueues EVERY task (`listTasks({ includeArchived: true })`) and `enqueue` released
each one straight into a `queueMicrotask`, so a 1,082-task project opened ~1,082 lock sessions at once
against `max_connections = 500`. The cluster saturated ~25s into boot; every later query then failed
with "sorry, too many clients already", including the engine's own startup, so Fusion wedged on
"starting" behind the migration holding server and never recovered.
The flat 1s retry made it self-sustaining rather than transient: once saturated, every task failed and
re-armed at a fixed 1s, so ~1,082 tasks re-opened ~1,082 sessions every second indefinitely (measured:
4,777 lock sessions in 17s). Backoff is exponential + jittered so a persistent outage decays instead of
pinning the resource it is waiting on.
*/
const RETRY_BASE_DELAY_MS = 1_000;
const RETRY_MAX_DELAY_MS = 60_000;
/** Max simultaneous reconciles, i.e. max simultaneous planning-lock sessions this component holds. */
const DEFAULT_MAX_CONCURRENT_RECONCILES = 4;
/** /**
* FNXC:SpecDrift 2026-08-10-09:28: * FNXC:SpecDrift 2026-08-10-09:28:
@@ -55,21 +77,37 @@ export function createStoreSpecDriftRepository(
export class SpecDriftReconciler { export class SpecDriftReconciler {
private stopped = false; private stopped = false;
private readonly retryTimers = new Map<string, ReturnType<typeof setTimeout>>(); private readonly retryTimers = new Map<string, ReturnType<typeof setTimeout>>();
/** Waiting to run. Insertion-ordered, so a burst is drained fairly rather than by task id. */
private readonly queuedTaskIds = new Set<string>(); private readonly queuedTaskIds = new Set<string>();
public constructor(private readonly repository: SpecDriftRepository) {} /** Running right now — each one holds a dedicated planning-lock connection. */
private readonly inFlightTaskIds = new Set<string>();
private readonly retryAttempts = new Map<string, number>();
private readonly maxConcurrent: number;
public constructor(
private readonly repository: SpecDriftRepository,
options: { maxConcurrent?: number } = {},
) {
this.maxConcurrent = Math.max(1, options.maxConcurrent ?? DEFAULT_MAX_CONCURRENT_RECONCILES);
}
/** /**
* FNXC:SpecDrift 2026-08-09-18:32: * FNXC:SpecDrift 2026-08-09-18:32:
* Live task events can arrive in one transaction-sized burst. Queue one fresh comparison per * Live task events can arrive in one transaction-sized burst. Queue one fresh comparison per
* task rather than making event delivery a polling loop; persistence still fences the snapshot. * task rather than making event delivery a polling loop; persistence still fences the snapshot.
*
* FNXC:SpecDrift 2026-08-10-18:32:
* The queue is now DRAINED by a bounded pump instead of releasing every id into its own
* microtask. The previous version deleted the dedupe key before the work ran, so it coalesced
* only ids still waiting in the same tick — a startup sweep over every task, or repeated
* `task:updated` events, still produced one concurrent run (and one connection) per task.
*/ */
enqueue(taskId: string): void { enqueue(taskId: string): void {
if (this.stopped || this.queuedTaskIds.has(taskId)) return; if (this.stopped) return;
// Already waiting, or already running with a re-run implied by the in-flight pass.
if (this.queuedTaskIds.has(taskId)) return;
this.queuedTaskIds.add(taskId); this.queuedTaskIds.add(taskId);
queueMicrotask(() => { this.pump();
this.queuedTaskIds.delete(taskId);
void this.reconcile(taskId).catch(() => undefined);
});
} }
stop(): void { stop(): void {
@@ -77,6 +115,34 @@ export class SpecDriftReconciler {
for (const timer of this.retryTimers.values()) clearTimeout(timer); for (const timer of this.retryTimers.values()) clearTimeout(timer);
this.retryTimers.clear(); this.retryTimers.clear();
this.queuedTaskIds.clear(); this.queuedTaskIds.clear();
this.retryAttempts.clear();
}
/**
* Start queued work up to the concurrency limit.
*
* Skips ids already in flight: a second reconcile for the same task would contend on the same
* per-task advisory lock while holding a second connection, which is the worst of both costs.
* Such an id stays queued and is picked up when its current pass finishes.
*/
private pump(): void {
if (this.stopped) return;
while (this.inFlightTaskIds.size < this.maxConcurrent) {
let next: string | undefined;
for (const candidate of this.queuedTaskIds) {
if (!this.inFlightTaskIds.has(candidate)) { next = candidate; break; }
}
if (next === undefined) return;
this.queuedTaskIds.delete(next);
this.inFlightTaskIds.add(next);
const taskId = next;
void this.reconcile(taskId)
.catch(() => undefined)
.finally(() => {
this.inFlightTaskIds.delete(taskId);
this.pump();
});
}
} }
async reconcile(taskId: string): Promise<DriftReport | undefined> { async reconcile(taskId: string): Promise<DriftReport | undefined> {
@@ -92,6 +158,7 @@ export class SpecDriftReconciler {
const retry = this.retryTimers.get(taskId); const retry = this.retryTimers.get(taskId);
if (retry) clearTimeout(retry); if (retry) clearTimeout(retry);
this.retryTimers.delete(taskId); this.retryTimers.delete(taskId);
this.retryAttempts.delete(taskId);
return report; return report;
} catch (error) { } catch (error) {
this.scheduleRetry(taskId); this.scheduleRetry(taskId);
@@ -99,12 +166,26 @@ export class SpecDriftReconciler {
} }
} }
/**
* FNXC:SpecDrift 2026-08-10-18:32:
* Exponential backoff with jitter, replacing a flat 1s. The failure this guards against is a
* SHARED-resource outage: every task fails for the same reason at the same time, so a fixed delay
* re-fires the whole fleet in lockstep once per second and keeps the resource exhausted. Retries
* re-enter through `enqueue`, so they are also subject to the concurrency bound — a retry storm
* can never open more connections than a first pass.
*/
private scheduleRetry(taskId: string): void { private scheduleRetry(taskId: string): void {
if (this.stopped || this.retryTimers.has(taskId)) return; if (this.stopped || this.retryTimers.has(taskId)) return;
const attempt = (this.retryAttempts.get(taskId) ?? 0) + 1;
this.retryAttempts.set(taskId, attempt);
const backoff = Math.min(RETRY_MAX_DELAY_MS, RETRY_BASE_DELAY_MS * 2 ** (attempt - 1));
const delay = backoff / 2 + Math.random() * (backoff / 2);
const timer = setTimeout(() => { const timer = setTimeout(() => {
this.retryTimers.delete(taskId); this.retryTimers.delete(taskId);
void this.reconcile(taskId).catch(() => undefined); this.enqueue(taskId);
}, RETRY_DELAY_MS); }, delay);
// Never hold the event loop open for a retry — this is best-effort background repair.
timer.unref?.();
this.retryTimers.set(taskId, timer); this.retryTimers.set(taskId, timer);
} }
} }