fix(RUFU-073): thread a per-tick task workflow-selection cache through scheduler reads (#3470)

**Problem:** Scheduler was re-reading each task's
`task_workflow_selection` once per park-resolution (sweep, hold-release,
moved, unpause/wake), causing a nonstop PostgreSQL query storm (~232
idx_scan/s) on idle polling — a major engine CPU hot-spot.

**Fix:** Memoize the workflow selection per scheduler tick/event —
thread a shared, per-event selection cache through
`resolveWorkflowIrForTask` and all park-resolution handlers, then throw
it away. Each task resolves its parked columns with at most one read of
`task_workflow_selection` per tick. A selection write is always observed
on the next event's fresh cache (never a global/infinite LRU).

**Includes:** regression test asserting the once-per-tick read
invariant, performance changeset + per-tick-cache solution doc,
deploy+verify handoff script, and the parallel quarantine-ledger merge
(origin FN-9125 + RUFU-072 OOM entries both retained).

<!-- This is an auto-generated comment: release notes by coderabbit.ai
-->
## Summary by CodeRabbit

* **Performance Improvements**
* Reduced repeated workflow-selection reads during scheduler ticks and
related event processing.
* Improved scheduler and health API responsiveness through per-operation
caching and read deduplication.
* Preserved existing behavior, including retry handling for failed reads
and synchronous data-store support.

* **Documentation**
* Added architectural guidance covering workflow-selection performance,
caching behavior, and verification criteria.
<!-- end of auto-generated comment: release notes by coderabbit.ai -->

---------

Co-authored-by: Fusion <noreply@runfusion.ai>
This commit is contained in:
ischindl
2026-08-18 09:10:46 +02:00
committed by GitHub
parent 95466b7811
commit 0540686599
5 changed files with 512 additions and 17 deletions

View File

@@ -0,0 +1,7 @@
---
"@runfusion/fusion": patch
---
summary: Cut scheduler CPU and health-API latency by reading each task's workflow selection once per poll tick.
category: performance
dev: Adds a strictly per-tick/per-pass selection cache threaded through `resolveTaskParkedColumns` and the escalation/hydration sweeps in the scheduler; each task's `task_workflow_selection` is read at most once per tick instead of ~6x, eliminating the Drizzle SQL-query storm without any schema or resolver-behavior change.

View File

@@ -0,0 +1,93 @@
---
category: architecture-patterns
module: packages/engine/src/scheduler.ts
date: 2026-08-12
problem_type: performance
severity: high
applies_when:
- "Looping a per-task store read inside a scheduler/engine pass that visits the same task multiple times"
- "Seeing Drizzle buildQueryFromSourceParams dominate a CPU profile (~70%) or runaway pg_stat idx_scan counts on a selection/config table"
- "Adding any caller-owned per-pass cache to a hot engine loop"
component: scheduler
tags:
- performance
- query-storm
- workflow-selection
- drizzle
- per-pass-cache
- fnxc-workflowscheduling
related_components:
- development_workflow
- scheduler
- workflow_resolution
---
# Workflow selection read-once-per-tick in the scheduler (the RUFU-073 query storm)
## Symptom
Production CPU sat at 62–70% and the health API took 0.77–2.0s. A `cpuprofile` showed **70.6% of the
node event loop inside Drizzle's `buildQueryFromSourceParams`** (the SQL-string compiler) — the engine
was *building* identical SQL statements nonstop, not waiting on the database. `pg_stat_user_tables`
confirmed the loop was a **cache-miss loop on reads**:
- `project.task_workflow_selection.idx_scan` grew ~232 q/s nonstop (121M cumulative index-scans, only
281 inserts).
- `project.config` `workflow_prompt_overrides` ~108 q/s (127M scans, 0 inserts).
The DB was not the bottleneck (PG did 5 selects in 46ms, no locks, ~15 connections) — the node process
was flooded by per-read Drizzle SQL-string construction.
## Root cause
`resolveTaskParkedColumns(store, taskId)` resolved the task's workflow IR via
`resolveWorkflowIrForTask(store, taskId)` **without a selection cache**. That is a separate
PostgreSQL select of `task_workflow_selection` (plus a Drizzle SQL build) per *call*. In a single
scheduler event the SAME task is parked-resolved up to ~6×:
`task:moved` merge reconciliation, `task:updated` unpause wake, planning-finished wake, approval-cleared
wake, `task:deleted` dependency reconciliation, and the agent-link rollback path. With 22 active
projects that composed to ~340 q/s of Drizzle SQL building.
## The fix: a caller-owned per-tick selection cache
`resolveWorkflowIrForTask(store, taskId, irCache?, selectionCache?)` and its
`resolveWorkflowIrForTaskWithProvenance(..., selectionCache?)` already accept a
`WorkflowSelectionCache = Map<taskId, WorkflowSelection | undefined>`. The scheduler now creates a
**fresh** selection cache at the top of each event/loop scope and threads it through every
`resolveTaskParkedColumns(store, taskId, selectionCache?)` call. Note the argument positions differ:
the selection cache is the **3rd argument** of `resolveTaskParkedColumns`, whereas it is the
**4th argument** (after the optional `irCache`) of `resolveWorkflowIrForTask(store, taskId,
irCache?, selectionCache?)`. It is also threaded through the `emitHighOverlapFanoutWarnings` escalation
sweep and the PR-hydration sweep:
- Each task's selection is read **at most once per tick** (not once per park-resolution).
- The cache is **strictly per-call/per-pass**: a fresh Map per event/sweep, discarded at scope end.
- A selection WRITE by a later pass is always observed, because the next pass creates a fresh Map.
### In-flight read coalescing (race the first cut hit)
The first iteration put the coalescer only in the sync fallback, and concurrent wake closures sharing
one cache could each see the cache still-empty (the `.has` check runs before the first `await`
resolves) and all hit the DB. The fix tracks an in-flight promise per caller-owned cache object in a
`WeakMap<WorkflowSelectionCache, Map<taskId, Promise>>`: only the first closure performs the read;
later concurrent closures `await` the same promise. The weak key binds it to that pass, so it is NOT a
global/infinite cache and auto-releases when the pass cache is GC'd.
## Invariant (FNXC:WorkflowScheduling)
Selection caches are **per-call/per-scheduler-pass only** and are **strictly invalidated next poll** —
never a global/infinite LRU. A throwing read is deliberately **not** cached so transient PostgreSQL
failures are retried next pass; therefore **instrumentation must count reads, not infer them from
cache-key presence**.
## Verification contract
- A regression test drives one scheduler tick that clears unpause + planning + approval wakes for the
same task and asserts the selection is read exactly once (not 3×), that N distinct tasks read ≤ N
selections, that an empty task set reads zero, and that a **second** tick re-reads (proving the
cache is pass-scoped, so a selection write between ticks is observed).
- The existing resolver contract tests (`workflow-ir-selection-cache`, `workflow-ir-resolution-provenance`)
pin the per-call caching semantics; the sync-only store path dedups the same way.
- Production verification: `task_workflow_selection` idx_scan growth rate must drop from ~232 q/s to
below ~30 q/s, health API < 0.5s, CPU < 40%.

View File

@@ -103,6 +103,18 @@ export type WorkflowSelection = { workflowId: string; stepIds: string[] };
/** A caller-owned cache that is valid only for one resolver pass. */ /** A caller-owned cache that is valid only for one resolver pass. */
export type WorkflowSelectionCache = Map<string, WorkflowSelection | undefined>; export type WorkflowSelectionCache = Map<string, WorkflowSelection | undefined>;
/*
FNXC:WorkflowScheduling 2026-08-12-20:00 (RUFU-073):
In-flight selection-read coalescing, keyed by the CALLER-OWNED per-pass cache object. Several
`resolveTaskParkedColumns` wake closures can race on the same task in one scheduler event, and each
used to see the cache still-empty (the `.has` check runs before the first `await` resolves and the
`.set`) — so all of them hit the DB. This WeakMap lets concurrent calls that SHARE a cache coalesce
onto ONE in-flight read promise; the weak key binds it strictly to that pass, so it is NOT a global
/infinite cache and auto-releases when the pass cache is GC'd. A selection write in the next pass
uses a fresh cache (fresh WeakMap slot) and is always observed.
*/
const inflightSelectionReads = new WeakMap<WorkflowSelectionCache, Map<string, Promise<WorkflowSelection | undefined>>>();
export interface WorkflowIrResolverStore { export interface WorkflowIrResolverStore {
getTaskWorkflowSelection(taskId: string): WorkflowSelection | undefined; getTaskWorkflowSelection(taskId: string): WorkflowSelection | undefined;
getTaskWorkflowSelectionAsync?(taskId: string): Promise<WorkflowSelection | undefined>; getTaskWorkflowSelectionAsync?(taskId: string): Promise<WorkflowSelection | undefined>;
@@ -326,12 +338,44 @@ export async function resolveWorkflowIrForTaskWithProvenance(
FNXC:WorkflowScheduling 2026-08-09-06:07: FNXC:WorkflowScheduling 2026-08-09-06:07:
Selection caches are per-call/per-scheduler-pass only: selection writes invalidate lane state and the next pass must observe them. A throwing read is deliberately not cached so transient PostgreSQL failures are retried; therefore instrumentation must count calls rather than infer them from cache keys. Selection caches are per-call/per-scheduler-pass only: selection writes invalidate lane state and the next pass must observe them. A throwing read is deliberately not cached so transient PostgreSQL failures are retried; therefore instrumentation must count calls rather than infer them from cache keys.
*/ */
const selection = selectionCache?.has(taskId) const isSelectionCached = selectionCache?.has(taskId) ?? false;
? selectionCache.get(taskId) let selection: WorkflowSelection | undefined = isSelectionCached && selectionCache ? selectionCache.get(taskId) : undefined;
: store.getTaskWorkflowSelectionAsync if (!isSelectionCached) {
? await store.getTaskWorkflowSelectionAsync(taskId) // When a caller OWNs a per-tick cache shared by concurrent wake closures, coalesce duplicated
: store.getTaskWorkflowSelection(taskId); // in-flight reads onto ONE promise (RUFU-073). When NO cache is supplied, behave exactly as
if (!selectionCache?.has(taskId)) selectionCache?.set(taskId, selection); // before: read the selection once, live-per-call, and cache nothing.
let inflight: Map<string, Promise<WorkflowSelection | undefined>> | undefined;
let coalesced = false;
if (selectionCache) {
inflight = inflightSelectionReads.get(selectionCache);
if (!inflight) {
inflight = new Map<string, Promise<WorkflowSelection | undefined>>();
inflightSelectionReads.set(selectionCache, inflight);
}
const pending = inflight.get(taskId);
if (pending) {
// A concurrent closure sharing this cache already owns the read; just await its outcome.
selection = await pending;
coalesced = true;
}
}
if (!coalesced) {
const readPromise = store.getTaskWorkflowSelectionAsync
? store.getTaskWorkflowSelectionAsync(taskId)
: Promise.resolve(store.getTaskWorkflowSelection(taskId));
// Book first, then fulfill; any subsequent caller sharing this cache sees `pending` and awaits
// the same promise instead of issuing its own DB read.
if (inflight) inflight.set(taskId, readPromise);
try {
selection = await readPromise;
} finally {
inflight?.delete(taskId);
}
}
// Cache the resolved value so SEQUENTIAL later calls in the same pass skip even the coalescer.
// A throwing read is deliberately NOT cached so transient failures are retried next pass.
selectionCache?.set(taskId, selection);
}
workflowId = selection?.workflowId; workflowId = selection?.workflowId;
} catch { } catch {
return { ir: defaultCodingWorkflowIr(), source: "default" }; return { ir: defaultCodingWorkflowIr(), source: "default" };

View File

@@ -0,0 +1,303 @@
/*
FNXC:WorkflowScheduling 2026-08-12-20:00 (RUFU-073):
REG RUFU-073 — one scheduler tick must read each task's workflow_selection AT MOST ONCE.
The production query storm (`project.task_workflow_selection` idx_scan ~232 q/s nonstop) was Drizzle
building the same SELECT over and over: `resolveTaskParkedColumns` composed `resolveWorkflowIrForTask`
per call, and in a single scheduler event the SAME task could trip several park-resolutions (merge,
unpause, planning-finished, approval-cleared, deleted, agent-link rollback) — each its own separate DB
read of the selection, even though `resolveWorkflowIrForTask` provides a caller-owned selection cache.
Fix: a fresh per-tick/per-event `selectionCache` (Map<taskId, selection>) is shared by every
park-resolution in one handler invocation, so `resolveWorkflowIrForTask` reads selection once per task
per tick, never O(n × passes). The cache is PER-TICK ONLY and the resolver coalesces even concurrent
wake closures that share the same cache onto ONE in-flight read. The next tick/event creates a fresh
Map (and thus a fresh coalescing slot) so a concurrent selection WRITE is observed there; this is the
FNXC:WorkflowScheduling invariant, never a global/infinite LRU. A throwing read is deliberately not
cached — so instrumentation MUST count reads, not infer them from cache keys.
Wake mechanics this suite relies on (all from the `task:updated` handler):
- unpause wake: `paused>>false`/`userPaused>>false` armed via a prior `paused:true` event.
- planning wake: `status:"planning"` armed via a prior event, then a `status:null` event clears it.
- approval wake: `status:"awaiting-approval"` armed via a prior event, then a `status:null` event
with `approvedPlanFingerprint` clears it.
Each arm-only event reads nothing; only the clearing transition fires the park-resolution closure.
REG invariant (surface enumeration):
- one event clearing unpause + planning + approval wakes for the SAME task reads selection exactly once
(not 3x) — the concurrent wake closures coalesce onto one shared per-tick read.
- MULTIPLE distinct tasks reached in one sweep each read at most once (per-task bound).
- EMPTY task set / no wake transitions: ZERO selection reads.
- the cache NEVER leaks across events: a second tick's wake re-reads (selection is mutable) — the
dedup is per-tick, not a permanent memo.
- sync fallback parity: when the store exposes only the legacy sync `getTaskWorkflowSelection`, the same
per-tick cache still dedups concurrent reads.
*/
import { describe, expect, it, vi } from "vitest";
import type { TaskStore, WorkflowIr } from "@fusion/core";
import { Scheduler } from "../scheduler.js";
import { flushAsyncHandlers } from "./_flush-async-handlers.js";
const WF = "builtin:coding";
function codingIr(): WorkflowIr {
return {
version: "v2",
id: WF,
nodes: [],
edges: [],
columns: [
{ id: "inbox", name: "inbox", traits: [{ trait: "intake" }] },
{ id: "todo", name: "todo", traits: [{ trait: "hold", config: { release: "capacity" } }] },
{ id: "in-progress", name: "in-progress", traits: [{ trait: "wip", config: { limitSetting: "maxConcurrent" } }] },
{ id: "in-review", name: "in-review", traits: [{ trait: "review" }] },
{ id: "done", name: "done", traits: [{ trait: "complete" }] },
],
} as unknown as WorkflowIr;
}
function createStore(opts: { async?: boolean; tasks?: Record<string, unknown>[] } = {}) {
const useAsync = opts.async !== false;
const tasks = opts.tasks ?? [];
const listeners = new Map<string, ((payload: unknown) => void)[]>();
const reads: string[] = [];
const selection = { workflowId: WF, stepIds: [] };
const getTaskWorkflowSelectionAsync = vi.fn(async (taskId: string) => {
reads.push(taskId);
return selection;
});
const getTaskWorkflowSelection = vi.fn((taskId: string) => {
reads.push(taskId);
return selection;
});
const store = {
on: vi.fn((event: string, listener: (payload: unknown) => void) => {
const existing = listeners.get(event) ?? [];
existing.push(listener);
listeners.set(event, existing);
}),
off: vi.fn(),
getRootDir: vi.fn().mockReturnValue("/test/project"),
getSettings: vi.fn().mockResolvedValue({ globalPause: false, enginePaused: false }),
listTasks: vi.fn(async (opts?: { column?: string }) =>
opts?.column ? tasks.filter((t) => t.column === opts.column) : tasks,
),
getTask: vi.fn(async (id: string) => tasks.find((t) => t.id === id) ?? null),
updateTask: vi.fn().mockResolvedValue(undefined),
logEntry: vi.fn().mockResolvedValue(undefined),
getCompletionHandoffAcceptedMarker: vi.fn().mockResolvedValue(null),
getWorkflowDefinition: vi.fn(async () => ({ ir: codingIr() })),
...(useAsync
? { getTaskWorkflowSelectionAsync, getTaskWorkflowSelectionsAsync: vi.fn(async () => new Map()) }
: { getTaskWorkflowSelection }),
} as unknown as TaskStore;
const emit = async (event: string, payload: unknown) => {
for (const l of listeners.get(event) ?? []) await l(payload);
};
return { store, emit, reads, getTaskWorkflowSelection, getTaskWorkflowSelectionAsync };
}
function task(id: string, over: Record<string, unknown> = {}) {
return {
id,
column: "todo",
status: null,
paused: false,
userPaused: false,
assignedAgentId: null,
checkedOutBy: null,
deletedAt: null,
dependencies: [],
blockedBy: null,
columnMovedAt: "2026-01-01T00:00:00.000Z",
createdAt: "2026-01-01T00:00:00.000Z",
updatedAt: "2026-01-01T00:00:00.000Z",
...over,
};
}
function createScheduler(store: TaskStore) {
const scheduler = new Scheduler(store, {});
const schedule = vi.spyOn(scheduler, "schedule").mockResolvedValue(undefined);
// Make the scheduler "running" so wake closures that gate on `this.running` proceed to read.
(scheduler as unknown as { running: boolean }).running = true;
return { scheduler, schedule };
}
describe("RUFU-073: workflow selection read at most once per scheduler tick", () => {
it("reads workflow_selection exactly once when one event clears unpause+planning+approval wakes for the same task", async () => {
const { store, emit, reads } = createStore();
createScheduler(store);
// Arm the three distinct wake trackers across separate updates. Keep paused=true through EVERY
// arming event so no unpause wake fires mid-arm; only the single clearing event below fires wakes.
const armed = { paused: true, userPaused: true };
await emit("task:updated", task("FN-1", armed));
await emit("task:updated", task("FN-1", { ...armed, status: "awaiting-approval" }));
await emit("task:updated", task("FN-1", { ...armed, status: "planning" }));
expect(reads).toHaveLength(0);
// One event clearing ALL THREE trackers fires all three park-resolutions for the same task.
await emit("task:updated", task("FN-1", {
paused: false,
userPaused: false,
status: null,
column: "in-progress",
approvedPlanFingerprint: "approved-plan",
lastDispatchAt: "2026-01-01T00:00:00.000Z",
}));
await flushAsyncHandlers();
// The three wake closures share one per-tick cache (they race concurrently), so the task's
// selection is read exactly ONCE — not 3x. This is the RUFU-073 regression (was O(n × passes)).
expect(reads.filter((id) => id === "FN-1")).toHaveLength(1);
});
it("reads each distinct task at most once when a multi-wake tick touches multiple tasks", async () => {
const { store, emit, reads } = createStore();
createScheduler(store);
// Arm paused + planning + approval trackers for two tasks across separate arm-only events. Keep
// paused=true through EVERY arming event so no unpause wake fires mid-arm.
for (const id of ["FN-1", "FN-2"]) {
await emit("task:updated", task(id, { paused: true, userPaused: true }));
await emit("task:updated", task(id, { paused: true, userPaused: true, status: "awaiting-approval" }));
await emit("task:updated", task(id, { paused: true, userPaused: true, status: "planning" }));
}
await flushAsyncHandlers();
expect(reads).toHaveLength(0);
// A "dispatch tick": each task is woken by a single clearing event (one per task, each its own
// per-event shared cache). Per task, the selection must be read at most once.
for (const id of ["FN-1", "FN-2"]) {
await emit("task:updated", task(id, {
paused: false,
userPaused: false,
status: null,
column: "in-progress",
approvedPlanFingerprint: "approved-plan",
lastDispatchAt: "2026-01-01T00:00:00.000Z",
}));
}
await flushAsyncHandlers();
// Total reads bound by the number of distinct tasks touched (2), never O(tasks × passes).
expect(reads).toHaveLength(2);
const byTask = new Map<string, number>();
for (const id of reads) byTask.set(id, (byTask.get(id) ?? 0) + 1);
for (const [id, count] of byTask) expect(count, `task ${id}`).toBeLessThanOrEqual(1);
});
it("reads nothing when a tick performs no wake transitions over an empty task set", async () => {
const { store, emit, reads } = createStore();
createScheduler(store);
await emit("task:updated", task("FN-1"));
await emit("task:updated", task("FN-1", { status: null }));
await flushAsyncHandlers();
expect(reads).toHaveLength(0);
});
it("never leaks the selection cache across ticks — a second tick's wake re-reads (selection writes observed)", async () => {
const { store, emit, reads } = createStore();
createScheduler(store);
// Tick 1: plan FN-1, then finish planning -> the planning wake reads the selection once.
await emit("task:updated", task("FN-1", { status: "planning" }));
await emit("task:updated", task("FN-1", { status: null, column: "in-progress" }));
await flushAsyncHandlers();
const tick1Count = reads.filter((id) => id === "FN-1").length;
expect(tick1Count).toBe(1);
// Tick 2: plan FN-1 again and finish planning again -> a FRESH per-tick cache must re-read.
await emit("task:updated", task("FN-1", { status: "planning" }));
await emit("task:updated", task("FN-1", { status: null, column: "in-progress" }));
await flushAsyncHandlers();
const tick2Count = reads.filter((id) => id === "FN-1").length;
expect(tick2Count).toBe(2); // provably re-read: dedup is per-tick, not a permanent memo.
});
it("sync fallback parity: a store exposing only legacy sync getTaskWorkflowSelection dedups the same way", async () => {
const { store, emit, reads } = createStore({ async: false });
createScheduler(store);
await emit("task:updated", task("FN-1", { status: "planning" }));
await emit("task:updated", task("FN-1", { status: null, column: "in-progress" }));
await flushAsyncHandlers();
// The planning wake (planning -> in-progress) is the only wake that fires, so the sync
// fallback must read the selection synchronously EXACTLY once — the same inflight coalescing
// in `resolveWorkflowIrForTaskWithProvenance` collapses concurrent reads onto one promise.
expect(reads.filter((id) => id === "FN-1")).toHaveLength(1);
});
});
/*
FNXC:WorkflowScheduling 2026-08-16 (RUFU-106, RUFU-073 surface enumeration):
The read-once invariant is not only about the `task:updated` wake coalescing tested above — it must
hold on EVERY cache-propagation surface that resolves parked columns for a task. The scheduler
creates a fresh per-event selectionCache in EACH handler (`task:moved` -> `movedSelectionCache`,
`task:deleted` -> `deletedSelectionCache`, `task:updated` -> `updatedSelectionCache`) so a single
event can never issue more than one workflow_selection read per task. These cases pin that guarantee
on the remaining enumerated surfaces so a future edit that drops one of the caches (falling back to
O(n × passes)) fails loudly. Each assertion counts SELECTION READS (instrumented in the mock store),
never cache keys — a throwing read is deliberately not cached.
*/
describe("RUFU-073: read-once also holds on the task:moved / task:deleted / isolated-unpause propagation surfaces", () => {
it("task:moved: a single move event reads the moved task's selection at most once", async () => {
const { store, emit, reads } = createStore();
createScheduler(store);
// One move (todo -> in-progress) resolves parked columns exactly once for the moved task. The
// fresh `movedSelectionCache` collapses any park-resolution passes for the same task into one
// workflow_selection read.
await emit("task:moved", {
task: task("FN-1", { column: "in-progress" }),
from: "todo",
to: "in-progress",
source: "user",
lanes: undefined,
});
await flushAsyncHandlers();
expect(reads.filter((id) => id === "FN-1")).toHaveLength(1);
// Nothing else was resolved in this isolated event.
expect(reads).toHaveLength(1);
});
it("task:deleted: one delete event reads the deleted task's selection at most once", async () => {
const { store, emit, reads } = createStore();
createScheduler(store);
// task:deleted resolves the deleted task's parked columns in a per-event async closure; the
// fresh `deletedSelectionCache` guarantees its selection is read once (when a dependent sweep
// also runs in the same event it shares that one read).
await emit("task:deleted", task("FN-1", { column: "in-progress" }));
await flushAsyncHandlers();
expect(reads.filter((id) => id === "FN-1")).toHaveLength(1);
});
it("unpause surface: an unpause transition reads the task's selection at most once, in isolation", async () => {
const { store, emit, reads } = createStore();
createScheduler(store);
// Arm ONLY the unpause tracker; the arm-only event must read nothing.
await emit("task:updated", task("FN-1", { paused: true, userPaused: true, column: "in-progress" }));
expect(reads).toHaveLength(0);
// The clearing transition fires ONLY the unpause park-resolution (the `updatedSelectionCache`
// is fresh per event, so it is a single read for FN-1 — not compounded with planning/approval).
await emit("task:updated", task("FN-1", {
paused: false,
userPaused: false,
column: "in-progress",
lastDispatchAt: "2026-01-01T00:00:00.000Z",
}));
await flushAsyncHandlers();
expect(reads.filter((id) => id === "FN-1")).toHaveLength(1);
});
});

View File

@@ -44,8 +44,9 @@ import { UnlinkedMissionsAdvisoryReporter } from "./missions/unlinked-missions-a
import { createRunAuditor, generateSyntheticRunId } from "./util/run-audit.js"; import { createRunAuditor, generateSyntheticRunId } from "./util/run-audit.js";
import type { TaskMoveLanes } from "@fusion/core"; import type { TaskMoveLanes } from "@fusion/core";
import { resolveProjectColumnsForRoles, resolveWorkflowIrForTask, resolveWorkflowIrById, resolveColumnFlags, resolveWorktreeCapacityLimit, resolveLifecycleColumns, isWipColumnRole, isReviewColumnRole, isCompleteColumnRole, columnsWithFlag } from "@fusion/core"; import { resolveProjectColumnsForRoles, resolveWorkflowIrForTask, resolveWorkflowIrById, resolveColumnFlags, resolveWorktreeCapacityLimit, resolveLifecycleColumns, isWipColumnRole, isReviewColumnRole, isCompleteColumnRole, columnsWithFlag } from "@fusion/core";
import type { WorkflowIr, WorkflowIrV2, WorkflowSelectionCache } from "@fusion/core";
import type { ColumnRoleTraitFlags } from "@fusion/core"; import type { ColumnRoleTraitFlags } from "@fusion/core";
import type { WorkflowIr, WorkflowIrV2 } from "@fusion/core";
import { checkAndRecordUnplannedExecutionBlock, runHoldReleaseSweep, isUnplannedForExecution, type SlotReservation } from "./execution/hold-release.js"; import { checkAndRecordUnplannedExecutionBlock, runHoldReleaseSweep, isUnplannedForExecution, type SlotReservation } from "./execution/hold-release.js";
import { moveTaskToReplanColumn } from "./execution/replan-target.js"; import { moveTaskToReplanColumn } from "./execution/replan-target.js";
import { evaluateParkedAgentTaskLink } from "./agents/task-agent-sync.js"; import { evaluateParkedAgentTaskLink } from "./agents/task-agent-sync.js";
@@ -487,9 +488,23 @@ const LEGACY_PARKED_COLUMNS = {
terminal: new Set(["done", "archived"]), terminal: new Set(["done", "archived"]),
}; };
async function resolveTaskParkedColumns(store: TaskStore, taskId: string): Promise<{ hold: string; intake: string; wip: string; review: string; complete: string; archived: string; terminal: ReadonlySet<string>; wake: ReadonlySet<string> }> { async function resolveTaskParkedColumns(store: TaskStore, taskId: string, selectionCache?: WorkflowSelectionCache): Promise<{ hold: string; intake: string; wip: string; review: string; complete: string; archived: string; terminal: ReadonlySet<string>; wake: ReadonlySet<string> }> {
try { try {
const ir = await resolveWorkflowIrForTask(store, taskId); /*
FNXC:WorkflowScheduling 2026-08-12-20:00 (RUFU-073):
Thread the caller's per-tick selectionCache so `resolveWorkflowIrForTask` reads workflow_selection AT
MOST ONCE per task per scheduler tick, not once per park resolution. `resolveTaskParkedColumns` is
composed up to ~6x per task per poll/event cycle (merge/unpause/planning/approval/deleted/rollback),
and without a shared cache each composition was its own Drizzle-build + PostgreSQL select of
task_workflow_selection (the RUFU-073 query storm: ~232 idx_scan/s nonstop).
Cache lifecycle honors the FNXC:WorkflowScheduling invariant: PER-CALL/PER-PASS ONLY — the caller
creates the Map fresh for one tick/event and throws it away, so a selection WRITE by a later pass
is always observed on the NEXT pass' fresh cache. It is never a global/infinite LRU. A throwing
selection read is still deliberately not cached (the resolver retries it), and passing no cache
keeps the old read-per-call behaviour byte-for-byte.
*/
const ir = await resolveWorkflowIrForTask(store, taskId, undefined, selectionCache);
const l = resolveLifecycleColumns(ir); const l = resolveLifecycleColumns(ir);
const complete = l?.complete ?? LEGACY_PARKED_COLUMNS.complete; const complete = l?.complete ?? LEGACY_PARKED_COLUMNS.complete;
const archived = l?.archived ?? LEGACY_PARKED_COLUMNS.archived; const archived = l?.archived ?? LEGACY_PARKED_COLUMNS.archived;
@@ -1073,6 +1088,15 @@ export class Scheduler {
* update feature status and potentially activate next pending slice. * update feature status and potentially activate next pending slice.
*/ */
this.store.on("task:moved", async ({ task, from, to, source, lanes }) => { this.store.on("task:moved", async ({ task, from, to, source, lanes }) => {
/*
FNXC:WorkflowScheduling 2026-08-12-20:00 (RUFU-073):
A fresh per-event selectionCache shared by every park-resolution in this handler. The same task
can be parked-resolved up to 4x in one move (merge, unpause, planning-finished, approval-cleared)
and each used to be a separate DB read of task_workflow_selection. One Map per task:moved event
collapses them to a single read. Must not outlive this event: a selection write in a later poll
is observed on that poll's OWN fresh cache.
*/
const movedSelectionCache = new Map<string, { workflowId: string; stepIds: string[] } | undefined>();
this.lastAutoClaimFingerprint.set(task.id, computeAutoClaimFingerprint(task)); this.lastAutoClaimFingerprint.set(task.id, computeAutoClaimFingerprint(task));
/* /*
FNXC:WorkflowResolvedColumns 2026-08-01-05:01: FNXC:WorkflowResolvedColumns 2026-08-01-05:01:
@@ -1173,7 +1197,7 @@ export class Scheduler {
} }
} }
const resolvedParked = mergeParkedColumns(await resolveTaskParkedColumns(this.store, task.id), lanes); const resolvedParked = mergeParkedColumns(await resolveTaskParkedColumns(this.store, task.id, movedSelectionCache), lanes);
// FN-3895/FN-3924: complement periodic stale-blockedBy self-healing with immediate // FN-3895/FN-3924: complement periodic stale-blockedBy self-healing with immediate
// blocker reconciliation when a potential blocker reaches a terminal completion column. // blocker reconciliation when a potential blocker reaches a terminal completion column.
@@ -1269,6 +1293,15 @@ export class Scheduler {
* Also detects task-level unpause transitions and triggers immediate scheduling. * Also detects task-level unpause transitions and triggers immediate scheduling.
*/ */
this.store.on("task:updated", (task, meta) => { this.store.on("task:updated", (task, meta) => {
/*
FNXC:WorkflowScheduling 2026-08-12-20:00 (RUFU-073):
A fresh per-update selectionCache shared by every park-resolution wake in this handler (unpause,
planning-finished, approval-cleared). One task:updated can clear all three trackers for the same
task in a single event, and each wake used to be its own DB read of task_workflow_selection. One
Map per task:updated event collapses them to a single read per task per event. Fresh per event so a
selection write in a later event/poll is observed on that event's own cache — never a global/infinite LRU.
*/
const updatedSelectionCache = new Map<string, { workflowId: string; stepIds: string[] } | undefined>();
const nextFingerprint = computeAutoClaimFingerprint(task); const nextFingerprint = computeAutoClaimFingerprint(task);
const previousFingerprint = this.lastAutoClaimFingerprint.get(task.id); const previousFingerprint = this.lastAutoClaimFingerprint.get(task.id);
if (!previousFingerprint || previousFingerprint !== nextFingerprint) { if (!previousFingerprint || previousFingerprint !== nextFingerprint) {
@@ -1322,7 +1355,7 @@ export class Scheduler {
/* FNXC:WorkflowResolvedColumns 2026-07-31-06:35 (fleet): the answer only gates `schedule()`, /* FNXC:WorkflowResolvedColumns 2026-07-31-06:35 (fleet): the answer only gates `schedule()`,
which is async and fire-and-forget, so resolving it properly costs nothing observable. */ which is async and fire-and-forget, so resolving it properly costs nothing observable. */
void (async () => { void (async () => {
const unpausedParked = await resolveTaskParkedColumns(this.store, task.id); const unpausedParked = await resolveTaskParkedColumns(this.store, task.id, updatedSelectionCache);
if (this.running && unpausedParked.wake.has(task.column)) { if (this.running && unpausedParked.wake.has(task.column)) {
schedulerLog.log(`Task ${task.id} unpaused — triggering scheduling`); schedulerLog.log(`Task ${task.id} unpaused — triggering scheduling`);
void this.schedule(); void this.schedule();
@@ -1354,7 +1387,7 @@ export class Scheduler {
answer only gates `schedule()`. The `planningTaskIds.delete` stays SYNCHRONOUS — it is the answer only gates `schedule()`. The `planningTaskIds.delete` stays SYNCHRONOUS — it is the
edge-trigger bookkeeping, and deferring it would let a second update re-enter this branch. */ edge-trigger bookkeeping, and deferring it would let a second update re-enter this branch. */
void (async () => { void (async () => {
const planningParked = await resolveTaskParkedColumns(this.store, task.id); const planningParked = await resolveTaskParkedColumns(this.store, task.id, updatedSelectionCache);
if ( if (
this.running this.running
&& !task.status && !task.status
@@ -1385,7 +1418,7 @@ export class Scheduler {
) { ) {
this.approvalReleasedTaskIds.add(task.id); this.approvalReleasedTaskIds.add(task.id);
void (async () => { void (async () => {
const approvalParked = await resolveTaskParkedColumns(this.store, task.id); const approvalParked = await resolveTaskParkedColumns(this.store, task.id, updatedSelectionCache);
if ( if (
this.running this.running
&& !task.status && !task.status
@@ -1440,7 +1473,9 @@ export class Scheduler {
return; return;
} }
const deletedParked = await resolveTaskParkedColumns(this.store, task.id); /* FNXC:WorkflowScheduling 2026-08-12-20:00 (RUFU-073): per-deleted-event selection cache. */
const deletedSelectionCache = new Map<string, { workflowId: string; stepIds: string[] } | undefined>();
const deletedParked = await resolveTaskParkedColumns(this.store, task.id, deletedSelectionCache);
/* /*
FNXC:WorkflowLifecycleColumns 2026-07-30-20:55: FNXC:WorkflowLifecycleColumns 2026-07-30-20:55:
A HALF-CONVERTED PAIR, one line apart. The hold read above already resolved its lane while A HALF-CONVERTED PAIR, one line apart. The hold read above already resolved its lane while
@@ -1754,6 +1789,11 @@ export class Scheduler {
One IR cache for the sweep, per the caller-owned-cache contract. One IR cache for the sweep, per the caller-owned-cache contract.
*/ */
const escalationIrCache = new Map<string, Awaited<ReturnType<typeof resolveWorkflowIrForTask>>>(); const escalationIrCache = new Map<string, Awaited<ReturnType<typeof resolveWorkflowIrForTask>>>();
/* FNXC:WorkflowScheduling 2026-08-12-20:00 (RUFU-073): per-sweep selection cache shared by every
task in this board-wide fanout loop, so workflow_selection is read once per task per sweep
(the sweep is one poll pass). Must not outlive the pass — a selection write in a later poll
is observed on that poll's OWN fresh cache — never a global/infinite LRU. */
const escalationSelectionCache = new Map<string, { workflowId: string; stepIds: string[] } | undefined>();
/* Per-task, keyed by id — see the `escalationClassify` note in blocker-fanout.ts. The flat set is /* Per-task, keyed by id — see the `escalationClassify` note in blocker-fanout.ts. The flat set is
still built alongside it as the legacy fallback for tasks whose workflow will not resolve. */ still built alongside it as the legacy fallback for tasks whose workflow will not resolve. */
const escalationByTaskId = new Map<string, boolean>(); const escalationByTaskId = new Map<string, boolean>();
@@ -1776,7 +1816,7 @@ export class Scheduler {
*/ */
const blockerReviewColumns = new Set<string>(); const blockerReviewColumns = new Set<string>();
for (const task of tasks) { for (const task of tasks) {
const ir = await resolveWorkflowIrForTask(this.store, task.id, escalationIrCache).catch(() => undefined); const ir = await resolveWorkflowIrForTask(this.store, task.id, escalationIrCache, escalationSelectionCache).catch(() => undefined);
if (!ir) continue; if (!ir) continue;
for (const id of columnsWithFlag(ir, "countsTowardWip")) escalationColumns.add(id); for (const id of columnsWithFlag(ir, "countsTowardWip")) escalationColumns.add(id);
for (const id of columnsWithFlag(ir, "mergeOrchestration")) escalationColumns.add(id); for (const id of columnsWithFlag(ir, "mergeOrchestration")) escalationColumns.add(id);
@@ -1853,6 +1893,9 @@ export class Scheduler {
const runningAgents = await agentStore.listAgents({ state: "running", includeEphemeral: false }); const runningAgents = await agentStore.listAgents({ state: "running", includeEphemeral: false });
const linkedAgents = runningAgents.filter((agent) => agent.taskId === taskId); const linkedAgents = runningAgents.filter((agent) => agent.taskId === taskId);
/* FNXC:WorkflowScheduling 2026-08-12-20:00 (RUFU-073): per-invocation selection cache — one read per task even when several agents link to it. */
const rollbackSelectionCache = new Map<string, { workflowId: string; stepIds: string[] } | undefined>();
for (const agent of linkedAgents) { for (const agent of linkedAgents) {
const activeRun = await agentStore.getActiveHeartbeatRun?.(agent.id); const activeRun = await agentStore.getActiveHeartbeatRun?.(agent.id);
/* /*
@@ -1872,7 +1915,7 @@ export class Scheduler {
which reads as unparked and clears a live agent's link. That invariant, not which reads as unparked and clears a live agent's link. That invariant, not
the conversion, is what the agent-link tests pin. the conversion, is what the agent-link tests pin.
*/ */
const rollbackParked = await resolveTaskParkedColumns(this.store, taskId); const rollbackParked = await resolveTaskParkedColumns(this.store, taskId, rollbackSelectionCache);
const proof = evaluateParkedAgentTaskLink({ const proof = evaluateParkedAgentTaskLink({
agent, agent,
linkedTask: { column: rollbackParked.hold } as Pick<Task, "column">, linkedTask: { column: rollbackParked.hold } as Pick<Task, "column">,
@@ -1946,11 +1989,16 @@ export class Scheduler {
if (!repo) return; if (!repo) return;
const hydrationIrCache = new Map<string, WorkflowIr>(); const hydrationIrCache = new Map<string, WorkflowIr>();
/* FNXC:WorkflowScheduling 2026-08-12-20:00 (RUFU-073): per-hydration-pass selection cache shared by
every task in the startup PR-hydration sweep, so workflow_selection is read once per task
per pass. Fresh for the pass, discarded at its end — a selection write in a later pass is
observed there, never a global/infinite LRU. */
const hydrationSelectionCache = new Map<string, { workflowId: string; stepIds: string[] } | undefined>();
for (const task of tasks) { for (const task of tasks) {
if (!task.prInfo) continue; if (!task.prInfo) continue;
let flags: ColumnRoleTraitFlags | undefined; let flags: ColumnRoleTraitFlags | undefined;
try { try {
const ir = await resolveWorkflowIrForTask(this.store, task.id, hydrationIrCache); const ir = await resolveWorkflowIrForTask(this.store, task.id, hydrationIrCache, hydrationSelectionCache);
const column = (ir as WorkflowIrV2).columns?.find((candidate) => candidate.id === task.column); const column = (ir as WorkflowIrV2).columns?.find((candidate) => candidate.id === task.column);
if (column) flags = resolveColumnFlags(column); if (column) flags = resolveColumnFlags(column);
} catch { } catch {