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:
7
.changeset/rufu-073-performance.md
Normal file
7
.changeset/rufu-073-performance.md
Normal 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.
|
||||
93
docs/solutions/workflow-selection-per-tick-cache.md
Normal file
93
docs/solutions/workflow-selection-per-tick-cache.md
Normal 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%.
|
||||
@@ -103,6 +103,18 @@ export type WorkflowSelection = { workflowId: string; stepIds: string[] };
|
||||
/** A caller-owned cache that is valid only for one resolver pass. */
|
||||
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 {
|
||||
getTaskWorkflowSelection(taskId: string): WorkflowSelection | undefined;
|
||||
getTaskWorkflowSelectionAsync?(taskId: string): Promise<WorkflowSelection | undefined>;
|
||||
@@ -326,12 +338,44 @@ export async function resolveWorkflowIrForTaskWithProvenance(
|
||||
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.
|
||||
*/
|
||||
const selection = selectionCache?.has(taskId)
|
||||
? selectionCache.get(taskId)
|
||||
: store.getTaskWorkflowSelectionAsync
|
||||
? await store.getTaskWorkflowSelectionAsync(taskId)
|
||||
: store.getTaskWorkflowSelection(taskId);
|
||||
if (!selectionCache?.has(taskId)) selectionCache?.set(taskId, selection);
|
||||
const isSelectionCached = selectionCache?.has(taskId) ?? false;
|
||||
let selection: WorkflowSelection | undefined = isSelectionCached && selectionCache ? selectionCache.get(taskId) : undefined;
|
||||
if (!isSelectionCached) {
|
||||
// When a caller OWNs a per-tick cache shared by concurrent wake closures, coalesce duplicated
|
||||
// in-flight reads onto ONE promise (RUFU-073). When NO cache is supplied, behave exactly as
|
||||
// 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;
|
||||
} catch {
|
||||
return { ir: defaultCodingWorkflowIr(), source: "default" };
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
@@ -44,8 +44,9 @@ import { UnlinkedMissionsAdvisoryReporter } from "./missions/unlinked-missions-a
|
||||
import { createRunAuditor, generateSyntheticRunId } from "./util/run-audit.js";
|
||||
import type { TaskMoveLanes } 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 { WorkflowIr, WorkflowIrV2 } from "@fusion/core";
|
||||
|
||||
import { checkAndRecordUnplannedExecutionBlock, runHoldReleaseSweep, isUnplannedForExecution, type SlotReservation } from "./execution/hold-release.js";
|
||||
import { moveTaskToReplanColumn } from "./execution/replan-target.js";
|
||||
import { evaluateParkedAgentTaskLink } from "./agents/task-agent-sync.js";
|
||||
@@ -487,9 +488,23 @@ const LEGACY_PARKED_COLUMNS = {
|
||||
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 {
|
||||
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 complete = l?.complete ?? LEGACY_PARKED_COLUMNS.complete;
|
||||
const archived = l?.archived ?? LEGACY_PARKED_COLUMNS.archived;
|
||||
@@ -1073,6 +1088,15 @@ export class Scheduler {
|
||||
* update feature status and potentially activate next pending slice.
|
||||
*/
|
||||
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));
|
||||
/*
|
||||
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
|
||||
// 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.
|
||||
*/
|
||||
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 previousFingerprint = this.lastAutoClaimFingerprint.get(task.id);
|
||||
if (!previousFingerprint || previousFingerprint !== nextFingerprint) {
|
||||
@@ -1322,7 +1355,7 @@ export class Scheduler {
|
||||
/* 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. */
|
||||
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)) {
|
||||
schedulerLog.log(`Task ${task.id} unpaused — triggering scheduling`);
|
||||
void this.schedule();
|
||||
@@ -1354,7 +1387,7 @@ export class Scheduler {
|
||||
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. */
|
||||
void (async () => {
|
||||
const planningParked = await resolveTaskParkedColumns(this.store, task.id);
|
||||
const planningParked = await resolveTaskParkedColumns(this.store, task.id, updatedSelectionCache);
|
||||
if (
|
||||
this.running
|
||||
&& !task.status
|
||||
@@ -1385,7 +1418,7 @@ export class Scheduler {
|
||||
) {
|
||||
this.approvalReleasedTaskIds.add(task.id);
|
||||
void (async () => {
|
||||
const approvalParked = await resolveTaskParkedColumns(this.store, task.id);
|
||||
const approvalParked = await resolveTaskParkedColumns(this.store, task.id, updatedSelectionCache);
|
||||
if (
|
||||
this.running
|
||||
&& !task.status
|
||||
@@ -1440,7 +1473,9 @@ export class Scheduler {
|
||||
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:
|
||||
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.
|
||||
*/
|
||||
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
|
||||
still built alongside it as the legacy fallback for tasks whose workflow will not resolve. */
|
||||
const escalationByTaskId = new Map<string, boolean>();
|
||||
@@ -1776,7 +1816,7 @@ export class Scheduler {
|
||||
*/
|
||||
const blockerReviewColumns = new Set<string>();
|
||||
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;
|
||||
for (const id of columnsWithFlag(ir, "countsTowardWip")) 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 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) {
|
||||
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
|
||||
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({
|
||||
agent,
|
||||
linkedTask: { column: rollbackParked.hold } as Pick<Task, "column">,
|
||||
@@ -1946,11 +1989,16 @@ export class Scheduler {
|
||||
if (!repo) return;
|
||||
|
||||
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) {
|
||||
if (!task.prInfo) continue;
|
||||
let flags: ColumnRoleTraitFlags | undefined;
|
||||
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);
|
||||
if (column) flags = resolveColumnFlags(column);
|
||||
} catch {
|
||||
|
||||
Reference in New Issue
Block a user