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. */
|
/** 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" };
|
||||||
|
|||||||
@@ -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 { 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 {
|
||||||
|
|||||||
Reference in New Issue
Block a user