Files
fusion/packages/engine/src/workflow-column-boundary.ts
gsxdsm 83209e64dc fix(workflows): align stages with board columns (#2378)
## Summary

The Coding (Ideas) workflow now behaves like the board it presents:
Ideas stays inert, Todo owns planning and plan review, In progress owns
implementation, and In review owns code review and merge. The restored
preset is intentionally limited to that five-stage path, while the
existing Coding workflow remains unchanged.

Workflow execution now suspends at Todo→In progress instead of running
the implementation node early. A durable, single-owner continuation
records the exact resume node and survives process restarts; the
scheduler remains the only component allowed to admit the task into WIP.
Disabled optional review groups traverse the same boundary without
invoking a reviewer, avoiding the prior stuck-task behavior.

Workflow validation also rejects capacity holds with no reachable WIP
destination, so deterministic lifecycle deadlocks fail at authoring time
rather than after a task is running.

Session-settled decisions carried from planning: columns are execution
invariants, scheduler-owned WIP admission is preserved, the existing
Coding (Ideas) preset is restored and simplified, and invalid release
topology is rejected (user-approved).

## Validation

- `pnpm lint`
- `pnpm verify:fast`
- `pnpm test:gate` (296 engine, 128 PostgreSQL core, and 63 CI-shape
tests)
- Focused workflow lifecycle tests (106 assertions)
- PostgreSQL regression coverage proves atomic continuation replacement
and database rejection of a second active owner


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

* **New Features**
* Added durable, resumable workflow execution across capacity boundaries
(including explicit suspend/resume at the correct node).
* Introduced Todo “plan review” workflow continuations and automated
planning/capacity draining.
* Restored Coding (Ideas) as a selectable built-in and updated its lane
placement; improved optional-step group enablement support.
* **Bug Fixes**
  * User moves back to Todo now cancels active workflow continuations.
* Rejected workflow boundary transitions now surface as errors (instead
of silently continuing).
* Workflows with undriveable capacity-hold configurations are now
rejected.
* **Tests / Data**
* Expanded coverage for workflow suspension, continuations, and
continuation replacement; updated database schema to persist
continuation metadata and enforce single active continuation.
<!-- end of auto-generated comment: release notes by coderabbit.ai -->
2026-07-21 12:17:47 -07:00

355 lines
14 KiB
TypeScript

/*
FNXC:WorkflowColumnBoundary 2026-07-18-20:35:
U1 — the graph is the sole lifecycle driver. As traversal crosses from a node in
column A into a node in column B, the card MOVES to column B through the store's
trait-hook `moveTask` path. This controller is the seam the graph executor calls
on every real-node entry; it owns the trait decision (move vs. park), the audit
emission (KTD-12), and the durable IR pin / drift guard (KTD-3).
Rules encoded here:
- KTD-1: `end` and failure-terminal arrivals never move — the executor never
calls `onNodeEntry` for a node without a column, and failure edges in the
benchmark terminate at `end` (columnless), so the card stays put and parks in
place. Only column-bearing node entries move.
- KTD-2: a hold→wip boundary is NEVER graph-moved. The graph parks the card at
the ready-for-release seam and the scheduler's capacity sweep (U4) is the sole
actor that performs the hold→wip move. This controller records the seam and
performs no move on that boundary.
- KTD-3: each node entry pins the resolved IR (version/content hash) durably on
the run state; `detectDrift` compares the stored pin against the current IR at
run start and parks with `task:reconcile-workflow-drift` on mismatch instead
of traversing a mutated graph.
- KTD-12: `task:column-transition` and `task:reconcile-workflow-drift` carry
ids/counts/outcomes-only metadata — no prose, no model ids.
Idempotency (U1 scenario 4 — kill/restart mid-transition settles exactly once):
the controller tracks the card's current column and only moves when the entered
node's column differs. A re-entered node (rework loop) or a re-run after restart
in the same column is a no-op, so exactly one move + one audit event is produced
per real boundary crossing. The crash-safe `transitionPending` marker written
in-txn by `moveTaskInternal` is the store-level half of that guarantee.
*/
import {
type WorkflowIr,
type WorkflowIrNode,
type WorkflowIrPin,
computeWorkflowIrPin,
detectWorkflowDrift,
findWorkflowColumn,
isHoldToWipBoundary,
resolveColumnFlags,
} from "@fusion/core";
/** Run-audit event emitted by the boundary controller (KTD-12, ids/counts only). */
export type WorkflowColumnBoundaryAuditEvent =
| {
type: "task:column-transition";
taskId: string;
workflowId: string;
fromColumn: string;
toColumn: string;
irHash: string;
nodeId: string;
}
| {
type: "task:reconcile-workflow-drift";
taskId: string;
workflowId: string;
pinnedNodeId: string;
reason: string;
};
/** The single store move seam. Wired to `store.moveTask(id, toColumn, {
* moveSource: "engine", workflowMoveSource: "workflow-graph" })`. */
export type WorkflowColumnMove = (
toColumn: string,
ctx: { fromColumn: string; nodeId: string },
) => Promise<void>;
export interface WorkflowColumnBoundaryDeps {
taskId: string;
workflowId: string;
ir: WorkflowIr;
/** The card's lifecycle column at run start (task.column). */
initialColumn: string;
/** The store move seam. Absent → the controller records intent but performs no
* move (byte-inert; used by unit tests that assert decisions only). */
moveTask?: WorkflowColumnMove;
/** Emit an ids/counts/outcomes-only run-audit event (KTD-12). */
emitAudit?: (event: WorkflowColumnBoundaryAuditEvent) => void | Promise<void>;
/** KTD-3: persist the per-node-entry IR pin (wired to the task row's
* `workflowIrPin*` fields via `createStoreIrPinPersistence` in production). */
pinNodeEntry?: (pin: WorkflowIrPin) => void | Promise<void>;
/*
FNXC:WorkflowIrPin 2026-07-19-21:10 (KTD-3 drift-park loop fix, PR #2342):
Clear the durable pin the moment drift is DETECTED. The pin that fired the drift
guard is by definition stale (its node/column/hash no longer exists in the current
IR); leaving it on the row made the park permanent — every requeue re-loaded the
same pin, re-fired detectDrift, and re-failed until an operator manually nulled
the fields. Clearing at drift-park time makes the park self-correcting: the next
ordinary requeue loads no prior pin, re-resolves the CURRENT IR fresh, and
proceeds — adopting the changed workflow is exactly the desired outcome.
Wired to `createStoreIrPinPersistence().clearPin` in production; absent → the
pre-fix posture (pin left in place) for minimal test harnesses.
*/
clearPin?: () => void | Promise<void>;
/** KTD-3: the pin recorded by a prior (possibly crashed) run, for drift check. */
priorPin?: WorkflowIrPin;
/** Optional diagnostics sink; never throws into the run. */
onWarn?: (message: string, detail: Record<string, unknown>) => void;
/** Persist a durable continuation before control returns to the scheduler. */
onSuspend?: (suspension: Extract<WorkflowColumnBoundaryEntryResult, { kind: "suspended" }>) => void | Promise<void>;
}
/** The seam the graph executor consumes. */
export interface WorkflowColumnBoundary {
/** The card's current lifecycle column (updated after each successful move). */
currentColumn(): string;
/** Cross into `node.column` when it differs from the current column. */
onNodeEntry(node: WorkflowIrNode): Promise<WorkflowColumnBoundaryEntryResult | void>;
/** KTD-3 drift guard — run once at graph start. Returns true when the pinned
* node/column is gone from the current IR (run must park, not traverse). */
detectDrift(): Promise<boolean>;
}
export type WorkflowColumnBoundaryEntryResult =
| { kind: "entered" }
| {
kind: "suspended";
reason: "capacity";
nodeId: string;
fromColumn: string;
toColumn: string;
irHash: string;
};
/*
FNXC:WorkflowIrPin 2026-07-19-18:30 (KTD-3 / U9b):
Store-backed KTD-3 pin persistence. The U9b schema landed the durable pin as three
task-row fields (`workflowIrPin` = content hash, `workflowIrPinNodeId`,
`workflowIrPinColumnId` — migration 0026), threaded through updateTask /
serialization / slim projections in core. This factory binds that row surface to
the boundary's `pinNodeEntry` / `loadPriorPin` seams:
- pinNodeEntry writes via `updateTask` ONLY when the pin actually changed
(same node re-entry / rework loops and same-hash chains are free — one row
write per real node entry, not per traversal step);
- loadPriorPin reads the fields back off the task row at run start so a
restart/re-entry compares the crashed run's pin against the CURRENT IR and
takes the drift-park path (`task:reconcile-workflow-drift`) on mismatch.
Degradation: a store lacking `updateTask`/`getTask`, or whose rows do not carry
the pin fields (in-memory fakes, pre-U9b DBs), yields no prior pin and no-op
writes — exactly the pre-wiring inert posture, so legacy harnesses are untouched.
*/
/** The minimal task-row surface the pin persistence needs. Structural on purpose
* so real stores and test fakes both fit without importing the full TaskStore. */
export interface WorkflowIrPinStoreSurface {
updateTask?: (
id: string,
updates: {
workflowIrPin?: string | null;
workflowIrPinNodeId?: string | null;
workflowIrPinColumnId?: string | null;
},
) => unknown;
getTask?: (id: string) => Promise<{
workflowIrPin?: string;
workflowIrPinNodeId?: string;
workflowIrPinColumnId?: string;
}>;
}
/** Build the store-backed `pinNodeEntry`/`loadPriorPin` pair for one task. */
export function createStoreIrPinPersistence(
store: WorkflowIrPinStoreSurface,
taskId: string,
): {
pinNodeEntry: (pin: WorkflowIrPin) => Promise<void>;
loadPriorPin: () => Promise<WorkflowIrPin | undefined>;
clearPin: () => Promise<void>;
} {
// Last pin known to be on the row (loaded or written) — the change-only gate.
let lastPin: WorkflowIrPin | undefined;
const samePin = (a: WorkflowIrPin, b: WorkflowIrPin): boolean =>
a.nodeId === b.nodeId && a.irHash === b.irHash && a.columnId === b.columnId;
return {
pinNodeEntry: async (pin) => {
if (typeof store.updateTask !== "function") return; // no-op degradation
if (lastPin && samePin(lastPin, pin)) return; // unchanged → no row write
await store.updateTask(taskId, {
workflowIrPin: pin.irHash,
workflowIrPinNodeId: pin.nodeId,
workflowIrPinColumnId: pin.columnId ?? null,
});
lastPin = pin;
},
loadPriorPin: async () => {
if (typeof store.getTask !== "function") return undefined;
try {
const row = await store.getTask(taskId);
if (!row?.workflowIrPin || !row.workflowIrPinNodeId) return undefined;
const pin: WorkflowIrPin = {
nodeId: row.workflowIrPinNodeId,
irHash: row.workflowIrPin,
columnId: row.workflowIrPinColumnId ?? undefined,
};
lastPin = pin;
return pin;
} catch {
// A store that cannot read the row degrades to "no prior pin" — the
// drift guard stays inert rather than failing the run on bookkeeping.
return undefined;
}
},
// FNXC:WorkflowIrPin 2026-07-19-21:10: null all three row fields (production
// task-update.ts treats null as clear-to-undefined) so a requeued run loads
// NO prior pin and re-resolves the current IR instead of re-firing drift.
clearPin: async () => {
if (typeof store.updateTask !== "function") return; // no-op degradation
await store.updateTask(taskId, {
workflowIrPin: null,
workflowIrPinNodeId: null,
workflowIrPinColumnId: null,
});
lastPin = undefined;
},
};
}
/**
* Build the boundary controller for one graph run. Additive: when the executor
* has no `columnBoundary` dep wired it performs no lifecycle moves at all, so
* every legacy runner/executor test stays byte-identical.
*/
export function createWorkflowColumnBoundary(
deps: WorkflowColumnBoundaryDeps,
): WorkflowColumnBoundary {
let column = deps.initialColumn;
const warn = (message: string, detail: Record<string, unknown>): void => {
try {
deps.onWarn?.(message, detail);
} catch {
/* diagnostics must never affect the run */
}
};
const flagsFor = (columnId: string) => {
const col = findWorkflowColumn(deps.ir, columnId);
return col ? resolveColumnFlags(col) : {};
};
return {
currentColumn: () => column,
async detectDrift(): Promise<boolean> {
const pin = deps.priorPin;
if (!pin) return false;
const reason = detectWorkflowDrift(deps.ir, pin);
if (!reason) return false;
try {
await deps.emitAudit?.({
type: "task:reconcile-workflow-drift",
taskId: deps.taskId,
workflowId: deps.workflowId,
pinnedNodeId: pin.nodeId,
reason,
});
} catch (err) {
warn("drift audit emit failed", { error: err instanceof Error ? err.message : String(err) });
}
// FNXC:WorkflowIrPin 2026-07-19-21:10 (drift-park loop fix): the pin that
// fired the guard is stale — clear it NOW so the park self-corrects on the
// next requeue (fresh IR resolution) instead of looping forever. Fail-soft:
// a failed clear merely re-fires drift next run (the pre-fix behavior),
// never fails this run's bookkeeping.
try {
await deps.clearPin?.();
} catch (err) {
warn("stale drift pin clear failed", { error: err instanceof Error ? err.message : String(err) });
}
return true;
},
async onNodeEntry(node: WorkflowIrNode): Promise<WorkflowColumnBoundaryEntryResult> {
const toColumn = node.column;
// KTD-1: a columnless node (e.g. `end`) never moves the card.
if (!toColumn) return { kind: "entered" };
// KTD-3: pin the resolved IR for this node-entry (durable seam).
try {
await deps.pinNodeEntry?.(computeWorkflowIrPin(deps.ir, node.id));
} catch (err) {
warn("ir pin write failed", { nodeId: node.id, error: err instanceof Error ? err.message : String(err) });
}
// Idempotent: a re-entered/rework node or a same-column node chain no-ops.
if (toColumn === column) return { kind: "entered" };
const fromColumn = column;
// KTD-2: never graph-move a hold→wip boundary — the scheduler is the sole
// mover there; the card parks at the ready-for-release seam (U4 performs the
// actual release). Record the seam but perform no move here.
if (isHoldToWipBoundary(flagsFor(fromColumn), flagsFor(toColumn))) {
warn("hold→wip boundary parked at ready-for-release seam (scheduler-owned)", {
fromColumn,
toColumn,
irHash: computeWorkflowIrPin(deps.ir, node.id).irHash,
nodeId: node.id,
});
const suspension = {
kind: "suspended",
reason: "capacity",
nodeId: node.id,
fromColumn,
toColumn,
irHash: computeWorkflowIrPin(deps.ir, node.id).irHash,
} as const;
await deps.onSuspend?.(suspension);
return suspension;
}
// The single mover: the store's trait-hook moveTask path.
if (deps.moveTask) {
try {
await deps.moveTask(toColumn, { fromColumn, nodeId: node.id });
} catch (err) {
// A rejected move (capacity, invariant) leaves the card in its current
// column; routing/parking is U4/U5. Do not advance `column` and do not
// emit a transition audit for a move that did not happen.
warn("graph column move rejected", {
fromColumn,
toColumn,
nodeId: node.id,
error: err instanceof Error ? err.message : String(err),
});
throw err;
}
}
column = toColumn;
try {
await deps.emitAudit?.({
type: "task:column-transition",
taskId: deps.taskId,
workflowId: deps.workflowId,
fromColumn,
toColumn,
irHash: computeWorkflowIrPin(deps.ir, node.id).irHash,
nodeId: node.id,
});
} catch (err) {
warn("column-transition audit emit failed", {
fromColumn,
toColumn,
error: err instanceof Error ? err.message : String(err),
});
}
return { kind: "entered" };
},
};
}