diff --git a/.changeset/rufu-073-performance.md b/.changeset/rufu-073-performance.md new file mode 100644 index 0000000000..41068143db --- /dev/null +++ b/.changeset/rufu-073-performance.md @@ -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. \ No newline at end of file diff --git a/docs/solutions/workflow-selection-per-tick-cache.md b/docs/solutions/workflow-selection-per-tick-cache.md new file mode 100644 index 0000000000..a86abd07d5 --- /dev/null +++ b/docs/solutions/workflow-selection-per-tick-cache.md @@ -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`. 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>`: 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%. \ No newline at end of file diff --git a/packages/core/src/workflows/workflow-ir-resolver.ts b/packages/core/src/workflows/workflow-ir-resolver.ts index ad8b935c41..2f34295540 100644 --- a/packages/core/src/workflows/workflow-ir-resolver.ts +++ b/packages/core/src/workflows/workflow-ir-resolver.ts @@ -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; +/* +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>>(); + export interface WorkflowIrResolverStore { getTaskWorkflowSelection(taskId: string): WorkflowSelection | undefined; getTaskWorkflowSelectionAsync?(taskId: string): Promise; @@ -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> | undefined; + let coalesced = false; + if (selectionCache) { + inflight = inflightSelectionReads.get(selectionCache); + if (!inflight) { + inflight = new Map>(); + 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" }; diff --git a/packages/engine/src/__tests__/scheduler-workflow-selection-read-once-per-tick.test.ts b/packages/engine/src/__tests__/scheduler-workflow-selection-read-once-per-tick.test.ts new file mode 100644 index 0000000000..17580099f0 --- /dev/null +++ b/packages/engine/src/__tests__/scheduler-workflow-selection-read-once-per-tick.test.ts @@ -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) 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[] } = {}) { + const useAsync = opts.async !== false; + const tasks = opts.tasks ?? []; + const listeners = new Map 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 = {}) { + 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(); + 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); + }); +}); diff --git a/packages/engine/src/scheduler.ts b/packages/engine/src/scheduler.ts index c3cfe920ea..364f46586f 100644 --- a/packages/engine/src/scheduler.ts +++ b/packages/engine/src/scheduler.ts @@ -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; wake: ReadonlySet }> { +async function resolveTaskParkedColumns(store: TaskStore, taskId: string, selectionCache?: WorkflowSelectionCache): Promise<{ hold: string; intake: string; wip: string; review: string; complete: string; archived: string; terminal: ReadonlySet; wake: ReadonlySet }> { 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(); 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(); 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(); + 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>>(); + /* 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(); /* 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(); @@ -1776,7 +1816,7 @@ export class Scheduler { */ const blockerReviewColumns = new Set(); 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(); + 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, @@ -1946,11 +1989,16 @@ export class Scheduler { if (!repo) return; const hydrationIrCache = new Map(); + /* 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(); 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 {