FN-8469: prevent workflow definition ID reuse

Prevent workflow definition ID reuse across projects and insertion races.

- Scan global workflow IDs before advancing SQLite or PostgreSQL allocators.
- Retry only confirmed workflow primary-key collisions while preserving unrelated constraint errors.
- Cover stale counters and concurrent allocation behavior in SQLite and PostgreSQL tests.

Files changed:
 .../fn-8469-workflow-definition-id-allocator.md    |   7 ++
 .../__tests__/postgres/workflow-create.pg.test.ts  |  68 +++++++++++
 .../workflow-definition-id-allocator-sync.test.ts  |  60 +++++++++
 .../workflow-definition-id-allocator.test.ts       |  26 ++++
 packages/core/src/task-store/remaining-ops-1.ts    | 134 +++++++++++----------
 .../core/src/task-store/workflow-definitions.ts    | 132 +++++++++++++-------
 6 files changed, 323 insertions(+), 104 deletions(-)

Fusion-Task-Id: FN-8469

Fusion-Task-Lineage: ba5bde88-67ca-47dd-ba86-00aef5258d13

Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
gsxdsm
2026-07-21 20:05:39 -07:00
parent c94920d885
commit 824762cdd8
6 changed files with 323 additions and 104 deletions

View File

@@ -0,0 +1,7 @@
---
"@runfusion/fusion": patch
---
summary: Stop workflow-definition creates from failing when a WF-id is already taken.
category: fix
dev: createWorkflowDefinition allocates past occupied global workflows.id values and retries id-PK unique conflicts instead of leaking Postgres 23505 to plugins/API callers (multi-project / stale next_workflow_definition_id).

View File

@@ -12,6 +12,11 @@
*/
import { describe, it, expect, beforeAll, beforeEach, afterEach, afterAll } from "vitest";
import type { AsyncDataLayer } from "../../postgres/data-layer.js";
import type { TaskStore } from "../../store.js";
import * as schema from "../../postgres/schema/index.js";
import { writeProjectConfig } from "../../task-store/async-settings.js";
import { __setWorkflowDefinitionBeforeInsertForTesting } from "../../task-store/remaining-ops-1.js";
import {
pgDescribe,
@@ -32,6 +37,15 @@ pgTest("workflow definition create (PostgreSQL backend mode)", () => {
afterEach(h.afterEach);
afterAll(h.afterAll);
function boundLayer(projectId: string): AsyncDataLayer {
return { ...h.layer(), projectId };
}
async function boundStore(projectId: string): Promise<TaskStore> {
const { TaskStore: TaskStoreCtor } = await import("../../store.js");
return new TaskStoreCtor(h.rootDir(), undefined, { asyncLayer: boundLayer(projectId) });
}
it("creates, updates, and deletes a workflow definition through the async layer", async () => {
const store = h.store();
expect(store.backendMode).toBe(true);
@@ -88,4 +102,58 @@ pgTest("workflow definition create (PostgreSQL backend mode)", () => {
// The sibling survives the delete (independent rows).
expect(await store.getWorkflowDefinition(second.id)).toBeDefined();
});
it("skips a globally occupied id when another project's counter is stale", async () => {
const projectA = "proj_workflow_allocator_a";
const projectB = "proj_workflow_allocator_b";
const now = new Date().toISOString();
await h.adminDb().insert(schema.project.workflows).values({
id: "WF-002",
name: "project A occupied workflow",
description: "",
icon: null,
ir: BUILTIN_CODING_WORKFLOW_IR as unknown as object,
layout: {},
kind: "workflow",
createdAt: now,
updatedAt: now,
});
await writeProjectConfig(boundLayer(projectB), {}, { nextWorkflowDefinitionId: 2 });
const storeB = await boundStore(projectB);
const created = await storeB.createWorkflowDefinition({ name: "project B workflow", ir: BUILTIN_CODING_WORKFLOW_IR });
expect(created.id).toBe("WF-003");
expect(await storeB.getWorkflowDefinition(created.id)).toMatchObject({ id: created.id });
expect(await (await boundStore(projectA)).getWorkflowDefinition("WF-002")).toMatchObject({ name: "project A occupied workflow" });
});
it("retries an id primary-key collision injected after allocation", async () => {
const now = new Date().toISOString();
let injected = false;
__setWorkflowDefinitionBeforeInsertForTesting(async (id, backendMode) => {
if (injected || !backendMode) return;
injected = true;
await h.adminDb().insert(schema.project.workflows).values({
id,
name: "race winner",
description: "",
icon: null,
ir: BUILTIN_CODING_WORKFLOW_IR as unknown as object,
layout: {},
kind: "workflow",
createdAt: now,
updatedAt: now,
});
});
try {
const created = await h.store().createWorkflowDefinition({ name: "race retry", ir: BUILTIN_CODING_WORKFLOW_IR });
expect(injected).toBe(true);
expect(created.id).toBe("WF-002");
expect(await h.store().getWorkflowDefinition("WF-001")).toMatchObject({ name: "race winner" });
expect(await h.store().getWorkflowDefinition(created.id)).toMatchObject({ name: "race retry" });
} finally {
__setWorkflowDefinitionBeforeInsertForTesting(undefined);
}
});
});

View File

@@ -0,0 +1,60 @@
import { describe, expect, it } from "vitest";
import { BUILTIN_CODING_WORKFLOW_IR } from "../builtin-coding-workflow-ir.js";
import { insertWorkflowDefinitionSyncImpl, nextWorkflowDefinitionIdImpl } from "../task-store/workflow-definitions.js";
import type { TaskStore } from "../store.js";
/**
* SQLite was removed at runtime (VAL-REMOVAL-005), but lifecycle materialization
* retains this synchronous compatibility branch. This narrow SQLite-shaped fake
* executes the real allocator and INSERT functions without reviving SQLite.
*/
function createSyncStoreWithStaleWorkflowCounter(): TaskStore {
const workflowIds = new Set(["WF-002"]);
const meta = new Map([["nextWorkflowDefinitionId", "2"]]);
const db = {
transactionImmediate<T>(operation: () => T): T { return operation(); },
prepare(query: string) {
if (query.includes("SELECT value FROM __meta")) {
return { get: () => meta.has("nextWorkflowDefinitionId") ? { value: meta.get("nextWorkflowDefinitionId") } : undefined };
}
if (query.includes("SELECT id FROM workflows")) {
return { all: () => [...workflowIds].map((id) => ({ id })) };
}
if (query.includes("INSERT INTO __meta")) {
return { run: (value: string) => { meta.set("nextWorkflowDefinitionId", value); } };
}
if (query.includes("INSERT INTO workflows")) {
return {
run: (id: string) => {
if (workflowIds.has(id)) throw new Error("UNIQUE constraint failed: workflows.id");
workflowIds.add(id);
},
};
}
throw new Error(`Unexpected SQLite-shaped query: ${query}`);
},
};
const store = {
db,
nextWorkflowDefinitionId() { return nextWorkflowDefinitionIdImpl(store as TaskStore); },
assertWorkflowIrTraitsValid() {},
workflowDefinitionsCache: null,
};
return store as unknown as TaskStore;
}
describe("workflow definition id allocator (sync materialization path)", () => {
it("allocates beyond a stale __meta counter and occupied workflow row", () => {
const store = createSyncStoreWithStaleWorkflowCounter();
const created = insertWorkflowDefinitionSyncImpl(store, {
name: "fresh workflow",
ir: BUILTIN_CODING_WORKFLOW_IR,
}, true);
expect(created.id).toBe("WF-003");
const second = insertWorkflowDefinitionSyncImpl(store, { name: "second workflow", ir: BUILTIN_CODING_WORKFLOW_IR }, true);
expect(second.id).toBe("WF-004");
});
});

View File

@@ -0,0 +1,26 @@
import { describe, expect, it } from "vitest";
import {
isWorkflowDefinitionIdPrimaryKeyCollision,
maxWorkflowDefinitionSequence,
} from "../task-store/workflow-definitions.js";
describe("workflow definition id allocator helpers", () => {
it("finds only numeric WF ids in the global occupancy set", () => {
expect(maxWorkflowDefinitionSequence([])).toBe(0);
expect(maxWorkflowDefinitionSequence(["WF-001", "WF-010", "custom-flow", "WF-abc"])).toBe(10);
});
it("accepts only workflow-id primary-key unique errors for retry", () => {
expect(isWorkflowDefinitionIdPrimaryKeyCollision({
code: "23505",
constraint: "workflows_pkey",
})).toBe(true);
expect(isWorkflowDefinitionIdPrimaryKeyCollision(new Error("UNIQUE constraint failed: workflows.id"))).toBe(true);
expect(isWorkflowDefinitionIdPrimaryKeyCollision({
code: "23505",
constraint: "some_other_unique",
})).toBe(false);
expect(isWorkflowDefinitionIdPrimaryKeyCollision(new Error("network unavailable"))).toBe(false);
});
});

View File

@@ -34,7 +34,7 @@ import {generateTaskLineageId} from "../task-lineage.js";
import {sanitizeFileScopeInPromptContent} from "../task-store/file-scope.js";
import {type TaskRow} from "../task-store/persistence.js";
import {__setTaskActivityLogLimitsForTesting} from "../task-store/comments.js";
import {nextWorkflowDefinitionIdAsyncImpl} from "../task-store/workflow-definitions.js";
import {isWorkflowDefinitionIdPrimaryKeyCollision, nextWorkflowDefinitionIdAsyncImpl} from "../task-store/workflow-definitions.js";
import {upsertTaskRowInTransaction, buildTaskInsertValues} from "../task-store/async-persistence.js";
import {readTaskRowInTransaction} from "../task-store/async-persistence.js";
import {recordActivityLogEntry as recordActivityLogEntryAsync} from "../task-store/async-audit.js";
@@ -929,6 +929,15 @@ export async function getWorkflowStepImpl(store: TaskStore, id: string): Promise
return template ? store.toBuiltInWorkflowStep(template) : undefined;
}
/** Test-only seam for proving the narrow retry handles a post-allocation race. */
export let workflowDefinitionBeforeInsertForTesting: ((id: string, backendMode: boolean) => void | Promise<void>) | undefined;
export function __setWorkflowDefinitionBeforeInsertForTesting(
hook: typeof workflowDefinitionBeforeInsertForTesting,
): void {
workflowDefinitionBeforeInsertForTesting = hook;
}
export async function createWorkflowDefinitionImpl(store: TaskStore, input: WorkflowDefinitionInput,): Promise<WorkflowDefinition> {
// Rollback compat (#1405): with the flag OFF, persist a pure-v1-equivalent
// graph in the v1 shape so a binary downgrade can still load the row.
@@ -943,71 +952,72 @@ export async function createWorkflowDefinitionImpl(store: TaskStore, input: Work
store.assertWorkflowIrTraitsValid(ir);
const layout = input.layout ?? {};
const now = new Date().toISOString();
// FNXC:SqliteFinalRemoval 2026-06-28:
// Backend mode (PG) allocates the WF-id from project.config via the async
// counter; the sync store.nextWorkflowDefinitionId() reads a SQLite __meta
// row that does not exist in PG. The id is computed up front so the
// definition object is identical across both branches.
const id = store.backendMode
? await nextWorkflowDefinitionIdAsyncImpl(store)
: store.nextWorkflowDefinitionId();
const definition: WorkflowDefinition = {
id,
name,
description: input.description ?? "",
icon: normalizeWorkflowIcon(input.icon),
// KTD-1: fragments are pure-v1 IRs and pass through downgradeIrToV1IfPure
// unchanged; default to "workflow" when the caller omits the kind.
kind: input.kind === "fragment" ? "fragment" : "workflow",
ir,
layout,
createdAt: now,
updatedAt: now,
};
/*
FNXC:WorkflowDefinitionIdAllocator 2026-07-21-12:00:
The global occupancy scan prevents stale per-project counters, but cannot
close a multi-process race after allocation. Retry only a confirmed
`workflows.id` PK conflict: withConfigLock serializes this process, not
another process, and retrying every unique error would hide unrelated
constraints from plugin and API callers.
*/
for (let attempt = 0; attempt < 8; attempt += 1) {
const id = store.backendMode
? await nextWorkflowDefinitionIdAsyncImpl(store)
: store.nextWorkflowDefinitionId();
const definition: WorkflowDefinition = {
id,
name,
description: input.description ?? "",
icon: normalizeWorkflowIcon(input.icon),
kind: input.kind === "fragment" ? "fragment" : "workflow",
ir,
layout,
createdAt: now,
updatedAt: now,
};
try {
await workflowDefinitionBeforeInsertForTesting?.(id, store.backendMode);
if (store.backendMode) {
await store.asyncLayer!.db.insert(schema.project.workflows).values({
id: definition.id,
name: definition.name,
description: definition.description,
icon: definition.icon ?? null,
ir: (flagOnForCreate ? definition.ir : downgradeIrToV1IfPure(definition.ir)) as unknown as object,
layout: definition.layout as unknown as object,
kind: definition.kind,
createdAt: definition.createdAt,
updatedAt: definition.updatedAt,
});
} else {
store.db
.prepare(
`INSERT INTO workflows (id, name, description, icon, ir, layout, kind, createdAt, updatedAt)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)`,
)
.run(
definition.id,
definition.name,
definition.description,
definition.icon ?? null,
serializeWorkflowIr(flagOnForCreate ? definition.ir : downgradeIrToV1IfPure(definition.ir)),
JSON.stringify(definition.layout),
definition.kind,
definition.createdAt,
definition.updatedAt,
);
}
} catch (error) {
if (!isWorkflowDefinitionIdPrimaryKeyCollision(error)) throw error;
continue;
}
if (store.backendMode) {
// FNXC:SqliteFinalRemoval 2026-06-28:
// PG INSERT via Drizzle. ir/layout are jsonb columns, so the OBJECT is
// passed directly (no serializeWorkflowIr/JSON.stringify — that is the
// SQLite TEXT path). Mirrors updateWorkflowDefinitionImpl's backend
// branch; bumpLastModified is skipped in backend mode.
await store.asyncLayer!.db.insert(schema.project.workflows).values({
id: definition.id,
name: definition.name,
description: definition.description,
icon: definition.icon ?? null,
ir: (flagOnForCreate ? definition.ir : downgradeIrToV1IfPure(definition.ir)) as unknown as object,
layout: definition.layout as unknown as object,
kind: definition.kind,
createdAt: definition.createdAt,
updatedAt: definition.updatedAt,
});
store.workflowDefinitionsCache = null;
if (!store.backendMode) store.db.bumpLastModified();
return definition;
}
store.db
.prepare(
`INSERT INTO workflows (id, name, description, icon, ir, layout, kind, createdAt, updatedAt)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)`,
)
.run(
definition.id,
definition.name,
definition.description,
definition.icon ?? null,
serializeWorkflowIr(
flagOnForCreate ? definition.ir : downgradeIrToV1IfPure(definition.ir),
),
JSON.stringify(definition.layout),
definition.kind,
definition.createdAt,
definition.updatedAt,
);
store.workflowDefinitionsCache = null;
store.db.bumpLastModified();
return definition;
throw new Error("Unable to allocate a free workflow definition id after repeated id collisions");
});
}

View File

@@ -26,6 +26,7 @@ import { PluginStore } from "../plugin-store.js";
import { SecretsStore } from "../secrets-store.js";
import { createAsyncDistributedTaskIdAllocator } from "./async-allocator.js";
import { getWorkflowRow, listWorkflowRows } from "../async-workflow-store.js";
import { isPostgresUniqueError } from "../postgres-errors.js";
import { projectOwnershipPartition, projectScopeFor, taskProjectScope } from "../postgres/data-layer.js";
import { getInReviewDurationEvents as getInReviewDurationEventsAsync, getTaskMergedTaskIds as getTaskMergedTaskIdsAsync } from "./async-audit.js";
import { readProjectConfig, writeProjectConfig } from "./async-settings.js";
@@ -180,21 +181,59 @@ export function resolvePluginWorkflowStepImpl(store: TaskStore, id: string): imp
};
}
/** Return the greatest numeric WF sequence while ignoring unrelated custom ids. */
export function maxWorkflowDefinitionSequence(ids: readonly string[]): number {
let max = 0;
for (const id of ids) {
const match = /^WF-(\d+)$/.exec(id);
if (!match) continue;
const value = Number.parseInt(match[1], 10);
if (Number.isSafeInteger(value)) max = Math.max(max, value);
}
return max;
}
function errorMessages(error: unknown): string {
const messages: string[] = [];
let current: unknown = error;
for (let depth = 0; current && depth < 5; depth += 1) {
const value = current as { message?: unknown; constraint?: unknown; detail?: unknown; cause?: unknown };
for (const candidate of [value.message, value.constraint, value.detail]) {
if (typeof candidate === "string") messages.push(candidate);
}
current = value.cause;
}
return messages.join(" ");
}
/**
* FNXC:SqliteFinalRemoval 2026-06-28:
* Backend-mode (PG) sibling of nextWorkflowDefinitionIdImpl. SQLite stored the
* WF-id counter in a __meta row read+incremented inside a transactionImmediate;
* PG has no __meta table so the counter lives in project.config
* (next_workflow_definition_id). The read+increment is serialized by the
* caller's withConfigLock (mirrors createWorkflowStepImpl's WS-id counter port),
* preserving the WF-### format and the never-reuse-across-deletes intent. The
* existing settings are passed back through writeProjectConfig so bumping the
* counter never clobbers the project settings object.
* Return true only for the global `workflows.id` primary-key target. A bare
* unique violation is deliberately insufficient because future workflow-table
* unique constraints must still reach callers unchanged.
*/
export function isWorkflowDefinitionIdPrimaryKeyCollision(error: unknown): boolean {
const message = errorMessages(error);
const targetsWorkflowId = /(?:project\.)?workflows_pkey\b|(?:project\.)?workflows\.id\b|(?:UNIQUE|PRIMARY KEY) constraint failed:\s*(?:project\.)?workflows\.id\b/i.test(message);
return targetsWorkflowId && (isPostgresUniqueError(error) || /SQLITE_CONSTRAINT|UNIQUE constraint failed|PRIMARY KEY constraint failed/i.test(message));
}
/**
* FNXC:WorkflowDefinitionIdAllocator 2026-07-21-12:00:
* `project.workflows.id` is global while `config.next_workflow_definition_id`
* belongs to one project. Allocate above the full unscoped workflows table as
* well as the monotonic local counter, otherwise a stale second-project counter
* can reissue an id owned by another project. The create path separately retries
* an id-PK race because this scan and withConfigLock are process-local.
*/
export async function nextWorkflowDefinitionIdAsyncImpl(store: TaskStore): Promise<string> {
const layer = store.asyncLayer!;
const configRow = await readProjectConfig(layer);
const next = configRow.nextWorkflowDefinitionId ?? 1;
const [configRow, workflows] = await Promise.all([
readProjectConfig(layer),
// listWorkflowRows deliberately has no project_id filter: workflow ids are global PKs.
listWorkflowRows(layer),
]);
const counter = configRow.nextWorkflowDefinitionId ?? 1;
const next = Math.max(counter, maxWorkflowDefinitionSequence(workflows.map(({ id }) => id)) + 1);
await writeProjectConfig(layer, configRow.settings ?? {}, {
nextWorkflowDefinitionId: next + 1,
});
@@ -209,7 +248,9 @@ export function nextWorkflowDefinitionIdImpl(store: TaskStore): string {
const row = store.db.prepare("SELECT value FROM __meta WHERE key = 'nextWorkflowDefinitionId'").get() as
| { value: string }
| undefined;
const next = row ? parseInt(row.value, 10) || 1 : 1;
const counter = row ? parseInt(row.value, 10) || 1 : 1;
const workflowRows = store.db.prepare("SELECT id FROM workflows").all() as Array<{ id: string }>;
const next = Math.max(counter, maxWorkflowDefinitionSequence(workflowRows.map(({ id }) => id)) + 1);
store.db
.prepare(
"INSERT INTO __meta (key, value) VALUES ('nextWorkflowDefinitionId', ?) ON CONFLICT(key) DO UPDATE SET value = excluded.value",
@@ -355,36 +396,43 @@ export function insertWorkflowDefinitionSyncImpl(store: TaskStore,
store.assertWorkflowIrTraitsValid(ir);
const layout = input.layout ?? {};
const now = new Date().toISOString();
const id = store.nextWorkflowDefinitionId();
const definition: WorkflowDefinition = {
id,
name,
description: input.description ?? "",
icon: normalizeWorkflowIcon(input.icon),
kind: input.kind === "fragment" ? "fragment" : "workflow",
ir,
layout,
createdAt: now,
updatedAt: now,
};
store.db
.prepare(
`INSERT INTO workflows (id, name, description, icon, ir, layout, kind, createdAt, updatedAt)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)`,
)
.run(
definition.id,
definition.name,
definition.description,
definition.icon ?? null,
serializeWorkflowIr(flagOn ? definition.ir : downgradeIrToV1IfPure(definition.ir)),
JSON.stringify(definition.layout),
definition.kind,
definition.createdAt,
definition.updatedAt,
);
store.workflowDefinitionsCache = null;
return definition;
for (let attempt = 0; attempt < 8; attempt += 1) {
const definition: WorkflowDefinition = {
id: store.nextWorkflowDefinitionId(),
name,
description: input.description ?? "",
icon: normalizeWorkflowIcon(input.icon),
kind: input.kind === "fragment" ? "fragment" : "workflow",
ir,
layout,
createdAt: now,
updatedAt: now,
};
try {
store.db
.prepare(
`INSERT INTO workflows (id, name, description, icon, ir, layout, kind, createdAt, updatedAt)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)`,
)
.run(
definition.id,
definition.name,
definition.description,
definition.icon ?? null,
serializeWorkflowIr(flagOn ? definition.ir : downgradeIrToV1IfPure(definition.ir)),
JSON.stringify(definition.layout),
definition.kind,
definition.createdAt,
definition.updatedAt,
);
} catch (error) {
if (!isWorkflowDefinitionIdPrimaryKeyCollision(error)) throw error;
continue;
}
store.workflowDefinitionsCache = null;
return definition;
}
throw new Error("Unable to allocate a free workflow definition id after repeated id collisions");
}
export async function isWorkflowCliCommandApprovedImpl(store: TaskStore, command: string): Promise<boolean> {