FN-8998: scope workflow definitions by project

Keep custom workflows isolated within their owning project.

- Model workflows with a project-scoped composite key and scope definition mutations.
- Preserve globally safe ID occupancy while isolating analytics name lookups.
- Add PostgreSQL coverage and document the project partition contract.

Files changed:
 .changeset/fn-8998-workflows-project-partition.md  |   7 +
 docs/storage.md                                    |   1 +
 .../workflows-project-isolation.pg.test.ts         | 146 +++++++++++++++++++++
 .../core/src/async-stores/async-workflow-store.ts  |  19 ++-
 packages/core/src/board/workflow-analytics.ts      |   7 +-
 packages/core/src/postgres/schema/project.ts       |  17 ++-
 .../core/src/task-store/workflow-definitions.ts    |  25 ++--
 packages/core/src/task-store/workflow-ops.ts       |  27 +++-
 8 files changed, 224 insertions(+), 25 deletions(-)

Fusion-Task-Id: FN-8998

Fusion-Task-Lineage: b6e8a100-de4d-4ff9-98df-4504d9050740

Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
gsxdsm
2026-08-11 20:25:59 -07:00
parent 9c5176f55e
commit 11334b1249
8 changed files with 224 additions and 25 deletions

View File

@@ -0,0 +1,7 @@
---
"@runfusion/fusion": patch
---
summary: Keep custom workflows private to their project on shared databases.
category: fix
dev: Models project.workflows as (project_id, id), scopes predicates with projectScopeFor, and preserves global ID occupancy allocation.

View File

@@ -37,6 +37,7 @@ See the [2026-07-14 PostgreSQL runtime cutover review](./postgres-migration-revi
- The `0000_initial.sql` baseline defines the table and indexes only. The later `0025_symbol_locks.sql` migration enables and forces RLS, creates `fusion_project_isolation`, and attaches `fusion_assign_project_id` after `0006_project_ownership.sql` creates that function/policy machinery. Both fresh full-applier and upgrade paths therefore end with the same project-isolation contract. - The `0000_initial.sql` baseline defines the table and indexes only. The later `0025_symbol_locks.sql` migration enables and forces RLS, creates `fusion_project_isolation`, and attaches `fusion_assign_project_id` after `0006_project_ownership.sql` creates that function/policy machinery. Both fresh full-applier and upgrade paths therefore end with the same project-isolation contract.
- `project.agent_ratings` is project-owned with composite `(project_id, id)` identity, allowing the same rating id in separate projects without cross-project reads or deletes. The dynamic `0006_project_ownership.sql` migration reconciles the physical table; `0055_fn_8988_agent_ratings_project_partition.sql` repeats that guarantee idempotently for historical drift. Bound `addRating`, `getRatings`, and `deleteRating` apply the project ownership partition, while unbound compatibility layers retain trigger-stamped writes and unscoped reads/deletes. - `project.agent_ratings` is project-owned with composite `(project_id, id)` identity, allowing the same rating id in separate projects without cross-project reads or deletes. The dynamic `0006_project_ownership.sql` migration reconciles the physical table; `0055_fn_8988_agent_ratings_project_partition.sql` repeats that guarantee idempotently for historical drift. Bound `addRating`, `getRatings`, and `deleteRating` apply the project ownership partition, while unbound compatibility layers retain trigger-stamped writes and unscoped reads/deletes.
- FN-8997 audited `workflow_steps`, `chat_room_members`, and `chat_room_messages`: their Drizzle declarations now model 0006's `project_id` and composite keys, and bound workflow/chat helpers scope reads and mutations on that partition. Chat isolation requires both membership/message predicates **and** the parent `chat_rooms.project_id` predicate; either leg alone can resolve a foreign row when room IDs collide. `plugins` remains an intentionally unmodeled compatibility table because it has no runtime Drizzle path. Migration `0056_fn_8997_project_ownership_declaration_drift.sql` is idempotent and adds only partition-prefixed predicate indexes; it does not rewrite healthy 0006 ownership columns or keys. - FN-8997 audited `workflow_steps`, `chat_room_members`, and `chat_room_messages`: their Drizzle declarations now model 0006's `project_id` and composite keys, and bound workflow/chat helpers scope reads and mutations on that partition. Chat isolation requires both membership/message predicates **and** the parent `chat_rooms.project_id` predicate; either leg alone can resolve a foreign row when room IDs collide. `plugins` remains an intentionally unmodeled compatibility table because it has no runtime Drizzle path. Migration `0056_fn_8997_project_ownership_declaration_drift.sql` is idempotent and adds only partition-prefixed predicate indexes; it does not rewrite healthy 0006 ownership columns or keys.
- `project.workflows` has project-local `(project_id, id)` identity. Bound definition reads, updates, deletes, companion workflow settings/prompt-override deletes, and analytics name prefetches use `projectScopeFor`; blank/unbound layers deliberately retain cross-project compatibility reads. The per-project workflow-id counter intentionally scans occupancy across every partition before allocation, because burning a colliding ID is safer than reusing a legacy or stale-counter ID held elsewhere.
- Startup and Batch 1 self-healing expire locks when their lease elapsed or the owner task is terminal/missing. They never move a task or alter scheduler, worktree, semaphore, or verification state. Run-audit events are `symbol-lock:acquired`, `symbol-lock:acquire-conflict`, `symbol-lock:renewed`, `symbol-lock:released`, `symbol-lock:reconcile-stale`, and deduplicated `symbol-lock:reconcile-stale-no-action`; metadata uses only counts/outcomes and normalized opaque keys. - Startup and Batch 1 self-healing expire locks when their lease elapsed or the owner task is terminal/missing. They never move a task or alter scheduler, worktree, semaphore, or verification state. Run-audit events are `symbol-lock:acquired`, `symbol-lock:acquire-conflict`, `symbol-lock:renewed`, `symbol-lock:released`, `symbol-lock:reconcile-stale`, and deduplicated `symbol-lock:reconcile-stale-no-action`; metadata uses only counts/outcomes and normalized opaque keys.
- FN-8405 adds `Task.declaredSymbols` as the durable, normalized task declaration source. `## Declared Symbols` in PROMPT.md is parsed only on create/update writes: an absent key may hydrate from the prompt, while a present `undefined`, `null` (update), or `[]` clears and suppresses hydration; a non-empty explicit array wins. Store resolution (`resolveTaskSymbols` and `resolveTaskSymbolsForWorkItem({ taskId })`) reads only the durable field, and slim projections plus archive/restore retain it. Scheduler admission remains a separate FN-8306 consumer; File Scope is never treated as a symbol source. - FN-8405 adds `Task.declaredSymbols` as the durable, normalized task declaration source. `## Declared Symbols` in PROMPT.md is parsed only on create/update writes: an absent key may hydrate from the prompt, while a present `undefined`, `null` (update), or `[]` clears and suppresses hydration; a non-empty explicit array wins. Store resolution (`resolveTaskSymbols` and `resolveTaskSymbolsForWorkItem({ taskId })`) reads only the durable field, and slim projections plus archive/restore retain it. Scheduler admission remains a separate FN-8306 consumer; File Scope is never treated as a symbol source.

View File

@@ -0,0 +1,146 @@
import { afterAll, afterEach, beforeAll, beforeEach, expect, it } from "vitest";
import { and, eq } from "drizzle-orm";
import {
createSharedPgTaskStoreTestHarness,
pgDescribe,
type SharedPgTaskStoreHarness,
} from "../../__test-utils__/pg-test-harness.js";
import { getWorkflowRow, listWorkflowRows } from "../../async-stores/async-workflow-store.js";
import { aggregateWorkflowAnalytics } from "../../board/workflow-analytics.js";
import type { AsyncDataLayer } from "../../postgres/data-layer.js";
import * as schema from "../../postgres/schema/index.js";
import type { TaskStore } from "../../store.js";
import { writeProjectConfig } from "../../task-store/async/async-settings.js";
import { nextWorkflowDefinitionIdAsyncImpl } from "../../task-store/workflow-definitions.js";
import { BUILTIN_CODING_WORKFLOW_IR } from "../../workflows/builtin-coding-workflow-ir.js";
pgDescribe("workflows project isolation", () => {
const h: SharedPgTaskStoreHarness = createSharedPgTaskStoreTestHarness({
prefix: "fusion_workflows_isolation",
});
beforeAll(h.beforeAll);
afterAll(h.afterAll);
beforeEach(h.beforeEach);
afterEach(h.afterEach);
const now = "2026-08-12T03:02:00.000Z";
const ir = BUILTIN_CODING_WORKFLOW_IR as unknown as object;
const bind = (projectId: string): AsyncDataLayer => ({ ...h.layer(), projectId });
const boundStore = async (projectId: string): Promise<TaskStore> => {
const { TaskStore: TaskStoreCtor } = await import("../../store.js");
return new TaskStoreCtor(h.rootDir(), undefined, { asyncLayer: bind(projectId) });
};
async function insertWorkflow(
layer: AsyncDataLayer,
id: string,
name: string,
workflowIr: object = ir,
): Promise<void> {
await layer.db.insert(schema.project.workflows).values({
...(layer.projectId?.trim() ? { projectId: layer.projectId } : {}),
id, name, description: "", icon: null, ir: workflowIr, layout: {}, kind: "workflow", createdAt: now, updatedAt: now,
});
}
it("models the live 0006 ownership shape exactly", async () => {
const columns = await h.adminSql()<Array<{ column_name: string; data_type: string; is_nullable: string; column_default: string | null }>>`
SELECT column_name, data_type, is_nullable, column_default
FROM information_schema.columns
WHERE table_schema = 'project' AND table_name = 'workflows' AND column_name = 'project_id'
`;
expect(columns).toEqual([expect.objectContaining({
column_name: "project_id", data_type: "text", is_nullable: "NO",
column_default: expect.stringContaining("current_setting"),
})]);
const keys = await h.adminSql()<Array<{ conname: string; columns: string[] }>>`
SELECT c.conname, array_agg(a.attname ORDER BY k.ordinality) AS columns
FROM pg_constraint c
CROSS JOIN LATERAL unnest(c.conkey) WITH ORDINALITY k(attnum, ordinality)
JOIN pg_attribute a ON a.attrelid = c.conrelid AND a.attnum = k.attnum
WHERE c.conrelid = 'project.workflows'::regclass AND c.contype = 'p'
GROUP BY c.conname
`;
expect(keys).toEqual([{ conname: "workflows_pkey", columns: ["project_id", "id"] }]);
expect(await h.adminSql()<Array<{ polname: string }>>`
SELECT polname FROM pg_policy WHERE polrelid = 'project.workflows'::regclass
`).toEqual([{ polname: "fusion_project_isolation" }]);
expect(await h.adminSql()<Array<{ tgname: string }>>`
SELECT tgname FROM pg_trigger WHERE tgrelid = 'project.workflows'::regclass
AND NOT tgisinternal
`).toEqual([{ tgname: "fusion_assign_project_id" }]);
});
it("isolates colliding workflow reads, writes, deletes, analytics, and preserves unbound compatibility", async () => {
/*
FNXC:WorkflowDefinitionProjectIsolation 2026-08-12-03:02:
Duplicate WF IDs are expected because counters are per project. These assertions exercise
owner-connected layers, where RLS is bypassed and each application predicate must preserve
the foreign workflow plus its companion settings and prompt overrides.
*/
const projectA = bind("workflows-project-a");
const projectB = bind("workflows-project-b");
const foreignIr = { foreign: "project-b-workflow-ir" };
await insertWorkflow(projectA, "WF-003", "Project A flow");
await insertWorkflow(projectA, "WF-005", "Project A only");
await insertWorkflow(projectB, "WF-003", "Project B flow", foreignIr);
await insertWorkflow(projectB, "WF-004", "Project B only");
for (const layer of [projectA, projectB]) {
await layer.db.insert(schema.project.workflowSettings).values({ workflowId: "WF-003", projectId: layer.projectId!, values: {}, updatedAt: now });
await layer.db.insert(schema.project.workflowPromptOverrides).values({ workflowId: "WF-003", projectId: layer.projectId!, overrides: {}, updatedAt: now });
}
expect((await listWorkflowRows(projectA)).map((row) => row.id)).toEqual(["WF-003", "WF-005"]);
expect((await listWorkflowRows(projectB)).map((row) => row.id)).toEqual(["WF-003", "WF-004"]);
expect(await getWorkflowRow(projectA, "WF-003")).toMatchObject({ id: "WF-003", name: "Project A flow" });
expect(await getWorkflowRow(projectB, "WF-003")).toMatchObject({ id: "WF-003", name: "Project B flow" });
expect("projectId" in (await getWorkflowRow(projectA, "WF-003"))!).toBe(false);
expect(await listWorkflowRows(bind("workflows-empty"))).toEqual([]);
const storeA = await boundStore(projectA.projectId!);
await storeA.createTaskWithReservedId(
{ description: "Analytics workflow name fixture", column: "done", workflowId: "WF-003" },
{ taskId: "FN-WORKFLOW-ANALYTICS", createdAt: now, updatedAt: now, applyDefaultWorkflowSteps: false },
);
await projectA.db.update(schema.project.tasks).set({ columnMovedAt: now })
.where(and(eq(schema.project.tasks.id, "FN-WORKFLOW-ANALYTICS"), eq(schema.project.tasks.projectId, projectA.projectId!)));
const analytics = await aggregateWorkflowAnalytics(projectA, {
from: "2026-08-12T00:00:00.000Z",
to: "2026-08-12T23:59:59.999Z",
});
expect(analytics.workflows).toContainEqual(expect.objectContaining({
workflowId: "WF-003", workflowName: "Project A flow", tasksCompleted: 1,
}));
await storeA.updateWorkflowDefinition("WF-003", { description: "A updated" });
expect(await getWorkflowRow(projectA, "WF-003")).toMatchObject({ description: "A updated" });
expect(await getWorkflowRow(projectB, "WF-003")).toMatchObject({
name: "Project B flow", description: "", ir: JSON.stringify(foreignIr), updatedAt: now,
});
expect(await projectB.db.select({ ir: schema.project.workflows.ir }).from(schema.project.workflows)
.where(and(eq(schema.project.workflows.id, "WF-003"), eq(schema.project.workflows.projectId, projectB.projectId!))))
.toEqual([{ ir: foreignIr }]);
await storeA.deleteWorkflowDefinition("WF-003");
expect(await getWorkflowRow(projectA, "WF-003")).toBeUndefined();
expect(await projectA.db.select().from(schema.project.workflowSettings).where(and(eq(schema.project.workflowSettings.workflowId, "WF-003"), eq(schema.project.workflowSettings.projectId, projectA.projectId!)))).toEqual([]);
expect(await projectA.db.select().from(schema.project.workflowPromptOverrides).where(and(eq(schema.project.workflowPromptOverrides.workflowId, "WF-003"), eq(schema.project.workflowPromptOverrides.projectId, projectA.projectId!)))).toEqual([]);
expect(await getWorkflowRow(projectB, "WF-003")).toMatchObject({
name: "Project B flow", ir: JSON.stringify(foreignIr),
});
expect(await projectB.db.select().from(schema.project.workflowSettings).where(and(eq(schema.project.workflowSettings.workflowId, "WF-003"), eq(schema.project.workflowSettings.projectId, projectB.projectId!)))).toHaveLength(1);
expect(await projectB.db.select().from(schema.project.workflowPromptOverrides).where(and(eq(schema.project.workflowPromptOverrides.workflowId, "WF-003"), eq(schema.project.workflowPromptOverrides.projectId, projectB.projectId!)))).toHaveLength(1);
await expect(storeA.deleteWorkflowDefinition("WF-004")).rejects.toThrow("Workflow 'WF-004' not found");
const unbound = { ...h.layer(), projectId: "" } satisfies AsyncDataLayer;
expect((await listWorkflowRows(unbound)).map((row) => row.name)).toEqual(["Project A only", "Project B flow", "Project B only"]);
await insertWorkflow(unbound, "WF-006", "Unbound trigger stamped");
const unboundStored = await unbound.db.select({ projectId: schema.project.workflows.projectId })
.from(schema.project.workflows).where(eq(schema.project.workflows.id, "WF-006"));
expect(unboundStored[0]?.projectId).toBe("__legacy_unscoped__");
await writeProjectConfig(projectA, {}, { nextWorkflowDefinitionId: 2 });
expect(await nextWorkflowDefinitionIdAsyncImpl(storeA)).toBe("WF-007");
});
});

View File

@@ -14,9 +14,9 @@
* shared `toWorkflowDefinition` mapper expects JSON strings (it parseWorkflowIr's * shared `toWorkflowDefinition` mapper expects JSON strings (it parseWorkflowIr's
* them), so we re-stringify here to keep the mapper backend-agnostic. * them), so we re-stringify here to keep the mapper backend-agnostic.
*/ */
import { asc, eq } from "drizzle-orm"; import { and, asc, eq } from "drizzle-orm";
import * as schema from "../postgres/schema/index.js"; import * as schema from "../postgres/schema/index.js";
import type { AsyncDataLayer } from "../postgres/data-layer.js"; import { projectScopeFor, type AsyncDataLayer } from "../postgres/data-layer.js";
import type { StoredWorkflowRow } from "../workflows/workflow-definition-types.js"; import type { StoredWorkflowRow } from "../workflows/workflow-definition-types.js";
/** SQLite-shaped workflow row (ir/layout as JSON strings) consumed by toWorkflowDefinition. */ /** SQLite-shaped workflow row (ir/layout as JSON strings) consumed by toWorkflowDefinition. */
@@ -41,6 +41,7 @@ export async function listWorkflowRows(layer: AsyncDataLayer): Promise<WorkflowR
const rows = await layer.db const rows = await layer.db
.select() .select()
.from(schema.project.workflows) .from(schema.project.workflows)
.where(projectScopeFor(schema.project.workflows.projectId, layer.projectId))
.orderBy(asc(schema.project.workflows.createdAt)); .orderBy(asc(schema.project.workflows.createdAt));
return rows.map(rowToWorkflowRow); return rows.map(rowToWorkflowRow);
} }
@@ -50,7 +51,19 @@ export async function getWorkflowRow(layer: AsyncDataLayer, id: string): Promise
const rows = await layer.db const rows = await layer.db
.select() .select()
.from(schema.project.workflows) .from(schema.project.workflows)
.where(eq(schema.project.workflows.id, id)) .where(and(
eq(schema.project.workflows.id, id),
projectScopeFor(schema.project.workflows.projectId, layer.projectId),
))
.limit(1); .limit(1);
return rows[0] ? rowToWorkflowRow(rows[0]) : undefined; return rows[0] ? rowToWorkflowRow(rows[0]) : undefined;
} }
/**
* Intentionally unscoped workflow-id occupancy scan used only by the allocator.
* A project-local allocator must burn a colliding ID rather than reuse an ID in
* another partition when a stale counter or legacy row is present.
*/
export async function listWorkflowIdsAcrossProjects(layer: AsyncDataLayer): Promise<Array<{ id: string }>> {
return layer.db.select({ id: schema.project.workflows.id }).from(schema.project.workflows);
}

View File

@@ -2,7 +2,7 @@ import { isReviewColumnRole, isWipColumnRole, type ColumnRoleTraitFlags } from "
import { resolveProjectColumnsForRoles, type ProjectLaneVocabularyStore } from "../project-lane-vocabulary.js"; import { resolveProjectColumnsForRoles, type ProjectLaneVocabularyStore } from "../project-lane-vocabulary.js";
import { sql } from "drizzle-orm"; import { sql } from "drizzle-orm";
import type { Database } from "../db/db.js"; import type { Database } from "../db/db.js";
import type { AsyncDataLayer } from "../postgres/data-layer.js"; import { projectScopeFor, type AsyncDataLayer } from "../postgres/data-layer.js";
import { BUILTIN_WORKFLOWS, getBuiltinWorkflow, isBuiltinWorkflowId } from "../workflows/builtin-workflows.js"; import { BUILTIN_WORKFLOWS, getBuiltinWorkflow, isBuiltinWorkflowId } from "../workflows/builtin-workflows.js";
import { costFor, type CostResult, type ModelPricingOverrides } from "../ai/model-pricing.js"; import { costFor, type CostResult, type ModelPricingOverrides } from "../ai/model-pricing.js";
import type { TokenTotals } from "./token-analytics.js"; import type { TokenTotals } from "./token-analytics.js";
@@ -592,9 +592,10 @@ async function aggregateWorkflowAnalyticsAsync(
modifiedFiles: r.modifiedFiles == null ? null : JSON.stringify(r.modifiedFiles), modifiedFiles: r.modifiedFiles == null ? null : JSON.stringify(r.modifiedFiles),
})); }));
// Prefetch all custom workflow names once; builtins resolve in-memory. // Prefetch custom names once; an unbound analytics layer intentionally sees all partitions.
const workflowScope = projectScopeFor(sql`project_id`, layer.projectId);
const workflowNameRows = (await layer.db.execute( const workflowNameRows = (await layer.db.execute(
sql`SELECT id, name FROM project.workflows`, sql`SELECT id, name FROM project.workflows ${workflowScope ? sql`WHERE ${workflowScope}` : sql``}`,
)) as Array<{ id: string; name: string | null }>; )) as Array<{ id: string; name: string | null }>;
const names = new Map<string, string>(); const names = new Map<string, string>();
for (const row of workflowNameRows) { for (const row of workflowNameRows) {

View File

@@ -706,8 +706,18 @@ export const workflowSteps = projectSchema.table("workflow_steps", {
index("idxWorkflowStepsProjectCreatedAt").on(t.projectId, t.createdAt), index("idxWorkflowStepsProjectCreatedAt").on(t.projectId, t.createdAt),
]); ]);
/*
FNXC:WorkflowDefinitionProjectPartition 2026-08-12-03:02:
Migration 0006 physically partitions workflows by project_id. WF-<n> identifiers use a
per-project counter, so collisions are normal; owner connections set fusion.project_bypass,
so RLS cannot protect unscoped application queries. The post-0006 harness confirms the physical
column, RLS policy, trigger, and ordered key are already exact, so no reconciliation migration is
needed. Model that composite identity here while keeping the database default so trigger-stamped
inserts can omit projectId.
*/
export const workflows = projectSchema.table("workflows", { export const workflows = projectSchema.table("workflows", {
id: text("id").primaryKey(), projectId: text("project_id").notNull().default(sql`current_setting('fusion.project_id', true)`),
id: text("id").notNull(),
name: text("name").notNull(), name: text("name").notNull(),
description: text("description").notNull().default(""), description: text("description").notNull().default(""),
icon: text("icon"), icon: text("icon"),
@@ -716,7 +726,10 @@ export const workflows = projectSchema.table("workflows", {
kind: text("kind").notNull().default("workflow"), kind: text("kind").notNull().default("workflow"),
createdAt: text("created_at").notNull(), createdAt: text("created_at").notNull(),
updatedAt: text("updated_at").notNull(), updatedAt: text("updated_at").notNull(),
}, (t) => [index("idxWorkflowsCreatedAt").on(t.createdAt)]); }, (t) => [
primaryKey({ columns: [t.projectId, t.id] }),
index("idxWorkflowsCreatedAt").on(t.createdAt),
]);
export const taskWorkflowSelection = projectSchema.table("task_workflow_selection", { export const taskWorkflowSelection = projectSchema.table("task_workflow_selection", {
projectId: text("project_id").notNull().default(sql`current_setting('fusion.project_id', true)`), projectId: text("project_id").notNull().default(sql`current_setting('fusion.project_id', true)`),

View File

@@ -24,7 +24,7 @@ import { type PluginGateVerdict } from "../plugins/plugin-gate-verdict.js";
import { PluginStore } from "../stores/plugin-store.js"; import { PluginStore } from "../stores/plugin-store.js";
import { SecretsStore } from "../secrets/secrets-store.js"; import { SecretsStore } from "../secrets/secrets-store.js";
import { createAsyncDistributedTaskIdAllocator } from "./async/async-allocator.js"; import { createAsyncDistributedTaskIdAllocator } from "./async/async-allocator.js";
import { getWorkflowRow, listWorkflowRows } from "../async-stores/async-workflow-store.js"; import { getWorkflowRow, listWorkflowIdsAcrossProjects, listWorkflowRows } from "../async-stores/async-workflow-store.js";
import { isPostgresUniqueError } from "../db/postgres-errors.js"; import { isPostgresUniqueError } from "../db/postgres-errors.js";
import {resolveColumnCapacity, resolveCapacityPoolId} from "../workflows/workflow-capacity.js"; import {resolveColumnCapacity, resolveCapacityPoolId} from "../workflows/workflow-capacity.js";
import {readTaskRow as readTaskRowAsync} from "./async/async-persistence.js"; import {readTaskRow as readTaskRowAsync} from "./async/async-persistence.js";
@@ -212,9 +212,9 @@ function errorMessages(error: unknown): string {
} }
/** /**
* Return true only for the global `workflows.id` primary-key target. A bare * Return true only for the `workflows` primary-key target. A bare unique
* unique violation is deliberately insufficient because future workflow-table * violation is deliberately insufficient because future workflow-table unique
* unique constraints must still reach callers unchanged. * constraints must still reach callers unchanged.
*/ */
export function isWorkflowDefinitionIdPrimaryKeyCollision(error: unknown): boolean { export function isWorkflowDefinitionIdPrimaryKeyCollision(error: unknown): boolean {
const message = errorMessages(error); const message = errorMessages(error);
@@ -223,19 +223,20 @@ export function isWorkflowDefinitionIdPrimaryKeyCollision(error: unknown): boole
} }
/** /**
* FNXC:WorkflowDefinitionIdAllocator 2026-07-21-12:00: * FNXC:WorkflowDefinitionIdAllocator 2026-08-12-03:02:
* `project.workflows.id` is global while `config.next_workflow_definition_id` * `project.workflows` now has project-local composite identity while
* belongs to one project. Allocate above the full unscoped workflows table as * `config.next_workflow_definition_id` remains per-project. Allocate above the
* well as the monotonic local counter, otherwise a stale second-project counter * full unscoped table as well as the local counter: legacy __legacy_unscoped__
* can reissue an id owned by another project. The create path separately retries * rows and stale counters can still reissue an ID held by another partition.
* an id-PK race because this scan and withConfigLock are process-local. * Burning an ID is harmless; reusing one is not. The create path separately
* retries a same-project PK race because this scan and withConfigLock are
* process-local.
*/ */
export async function nextWorkflowDefinitionIdAsyncImpl(store: TaskStore): Promise<string> { export async function nextWorkflowDefinitionIdAsyncImpl(store: TaskStore): Promise<string> {
const layer = store.asyncLayer!; const layer = store.asyncLayer!;
const [configRow, workflows] = await Promise.all([ const [configRow, workflows] = await Promise.all([
readProjectConfig(layer), readProjectConfig(layer),
// listWorkflowRows deliberately has no project_id filter: workflow ids are global PKs. listWorkflowIdsAcrossProjects(layer),
listWorkflowRows(layer),
]); ]);
const counter = configRow.nextWorkflowDefinitionId ?? 1; const counter = configRow.nextWorkflowDefinitionId ?? 1;
const next = Math.max(counter, maxWorkflowDefinitionSequence(workflows.map(({ id }) => id)) + 1); const next = Math.max(counter, maxWorkflowDefinitionSequence(workflows.map(({ id }) => id)) + 1);

View File

@@ -321,7 +321,10 @@ export async function updateWorkflowDefinitionImpl(store: TaskStore, id: string,
ir: downgradeIrToV1IfPure(next.ir), ir: downgradeIrToV1IfPure(next.ir),
layout: next.layout, layout: next.layout,
updatedAt: next.updatedAt, updatedAt: next.updatedAt,
}).where(eq(schema.project.workflows.id, id)); }).where(and(
eq(schema.project.workflows.id, id),
projectScopeFor(schema.project.workflows.projectId, layer.projectId),
));
store.workflowDefinitionsCache = null; store.workflowDefinitionsCache = null;
return next; return next;
@@ -381,12 +384,26 @@ export async function deleteWorkflowDefinitionImpl(store: TaskStore, id: string)
const occupantTaskIds = await store.listWorkflowOccupantTaskIds(id, false); const occupantTaskIds = await store.listWorkflowOccupantTaskIds(id, false);
// FNXC:PostgresCutover 2026-06-28: async deletes for backend mode /*
const deleted = await layer.db.delete(schema.project.workflows).where(eq(schema.project.workflows.id, id)).returning(); FNXC:WorkflowDefinitionProjectPartition 2026-08-12-03:02:
Before project ownership predicates, deleting a colliding workflow ID removed a foreign
workflow and its settings/prompt overrides under owner connections that bypass RLS. Keep
every delete in this transaction on the same bound partition.
*/
const deleted = await layer.db.delete(schema.project.workflows).where(and(
eq(schema.project.workflows.id, id),
projectScopeFor(schema.project.workflows.projectId, layer.projectId),
)).returning();
if (deleted.length === 0) throw new Error(`Workflow '${id}' not found`); if (deleted.length === 0) throw new Error(`Workflow '${id}' not found`);
store.workflowDefinitionsCache = null; store.workflowDefinitionsCache = null;
await layer.db.delete(schema.project.workflowSettings).where(eq(schema.project.workflowSettings.workflowId, id)); await layer.db.delete(schema.project.workflowSettings).where(and(
await layer.db.delete(schema.project.workflowPromptOverrides).where(eq(schema.project.workflowPromptOverrides.workflowId, id)); eq(schema.project.workflowSettings.workflowId, id),
projectScopeFor(schema.project.workflowSettings.projectId, layer.projectId),
));
await layer.db.delete(schema.project.workflowPromptOverrides).where(and(
eq(schema.project.workflowPromptOverrides.workflowId, id),
projectScopeFor(schema.project.workflowPromptOverrides.projectId, layer.projectId),
));
// Cascade: clear the project default when it pointed at this workflow. // Cascade: clear the project default when it pointed at this workflow.