diff --git a/.changeset/fn-8998-workflows-project-partition.md b/.changeset/fn-8998-workflows-project-partition.md new file mode 100644 index 0000000000..fd7cd4439a --- /dev/null +++ b/.changeset/fn-8998-workflows-project-partition.md @@ -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. diff --git a/docs/storage.md b/docs/storage.md index 8bc2e3f663..f8e7fa6e43 100644 --- a/docs/storage.md +++ b/docs/storage.md @@ -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. - `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. +- `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. - 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. diff --git a/packages/core/src/__tests__/postgres/workflows-project-isolation.pg.test.ts b/packages/core/src/__tests__/postgres/workflows-project-isolation.pg.test.ts new file mode 100644 index 0000000000..9db38b85cd --- /dev/null +++ b/packages/core/src/__tests__/postgres/workflows-project-isolation.pg.test.ts @@ -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 => { + 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 { + 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()>` + 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()>` + 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()>` + SELECT polname FROM pg_policy WHERE polrelid = 'project.workflows'::regclass + `).toEqual([{ polname: "fusion_project_isolation" }]); + expect(await h.adminSql()>` + 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"); + }); +}); diff --git a/packages/core/src/async-stores/async-workflow-store.ts b/packages/core/src/async-stores/async-workflow-store.ts index 7c39ed00f5..8a71dda118 100644 --- a/packages/core/src/async-stores/async-workflow-store.ts +++ b/packages/core/src/async-stores/async-workflow-store.ts @@ -14,9 +14,9 @@ * shared `toWorkflowDefinition` mapper expects JSON strings (it parseWorkflowIr's * 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 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"; /** SQLite-shaped workflow row (ir/layout as JSON strings) consumed by toWorkflowDefinition. */ @@ -41,6 +41,7 @@ export async function listWorkflowRows(layer: AsyncDataLayer): Promise> { + return layer.db.select({ id: schema.project.workflows.id }).from(schema.project.workflows); +} diff --git a/packages/core/src/board/workflow-analytics.ts b/packages/core/src/board/workflow-analytics.ts index 90b640f101..888dc87417 100644 --- a/packages/core/src/board/workflow-analytics.ts +++ b/packages/core/src/board/workflow-analytics.ts @@ -2,7 +2,7 @@ import { isReviewColumnRole, isWipColumnRole, type ColumnRoleTraitFlags } from " import { resolveProjectColumnsForRoles, type ProjectLaneVocabularyStore } from "../project-lane-vocabulary.js"; import { sql } from "drizzle-orm"; 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 { costFor, type CostResult, type ModelPricingOverrides } from "../ai/model-pricing.js"; import type { TokenTotals } from "./token-analytics.js"; @@ -592,9 +592,10 @@ async function aggregateWorkflowAnalyticsAsync( 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( - 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 }>; const names = new Map(); for (const row of workflowNameRows) { diff --git a/packages/core/src/postgres/schema/project.ts b/packages/core/src/postgres/schema/project.ts index 7908532c9c..91083c2a1f 100644 --- a/packages/core/src/postgres/schema/project.ts +++ b/packages/core/src/postgres/schema/project.ts @@ -706,8 +706,18 @@ export const workflowSteps = projectSchema.table("workflow_steps", { index("idxWorkflowStepsProjectCreatedAt").on(t.projectId, t.createdAt), ]); +/* +FNXC:WorkflowDefinitionProjectPartition 2026-08-12-03:02: +Migration 0006 physically partitions workflows by project_id. WF- 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", { - 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(), description: text("description").notNull().default(""), icon: text("icon"), @@ -716,7 +726,10 @@ export const workflows = projectSchema.table("workflows", { kind: text("kind").notNull().default("workflow"), createdAt: text("created_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", { projectId: text("project_id").notNull().default(sql`current_setting('fusion.project_id', true)`), diff --git a/packages/core/src/task-store/workflow-definitions.ts b/packages/core/src/task-store/workflow-definitions.ts index e11a56a9df..2edc23e4aa 100644 --- a/packages/core/src/task-store/workflow-definitions.ts +++ b/packages/core/src/task-store/workflow-definitions.ts @@ -24,7 +24,7 @@ import { type PluginGateVerdict } from "../plugins/plugin-gate-verdict.js"; import { PluginStore } from "../stores/plugin-store.js"; import { SecretsStore } from "../secrets/secrets-store.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 {resolveColumnCapacity, resolveCapacityPoolId} from "../workflows/workflow-capacity.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 - * unique violation is deliberately insufficient because future workflow-table - * unique constraints must still reach callers unchanged. + * Return true only for the `workflows` 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); @@ -223,19 +223,20 @@ export function isWorkflowDefinitionIdPrimaryKeyCollision(error: unknown): boole } /** - * 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. + * FNXC:WorkflowDefinitionIdAllocator 2026-08-12-03:02: + * `project.workflows` now has project-local composite identity while + * `config.next_workflow_definition_id` remains per-project. Allocate above the + * full unscoped table as well as the local counter: legacy __legacy_unscoped__ + * rows and stale counters can still reissue an ID held by another partition. + * 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 { const layer = store.asyncLayer!; const [configRow, workflows] = await Promise.all([ readProjectConfig(layer), - // listWorkflowRows deliberately has no project_id filter: workflow ids are global PKs. - listWorkflowRows(layer), + listWorkflowIdsAcrossProjects(layer), ]); const counter = configRow.nextWorkflowDefinitionId ?? 1; const next = Math.max(counter, maxWorkflowDefinitionSequence(workflows.map(({ id }) => id)) + 1); diff --git a/packages/core/src/task-store/workflow-ops.ts b/packages/core/src/task-store/workflow-ops.ts index fea4875318..bb1469ec44 100644 --- a/packages/core/src/task-store/workflow-ops.ts +++ b/packages/core/src/task-store/workflow-ops.ts @@ -321,7 +321,10 @@ export async function updateWorkflowDefinitionImpl(store: TaskStore, id: string, ir: downgradeIrToV1IfPure(next.ir), layout: next.layout, 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; return next; @@ -381,12 +384,26 @@ export async function deleteWorkflowDefinitionImpl(store: TaskStore, id: string) 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`); store.workflowDefinitionsCache = null; - await layer.db.delete(schema.project.workflowSettings).where(eq(schema.project.workflowSettings.workflowId, id)); - await layer.db.delete(schema.project.workflowPromptOverrides).where(eq(schema.project.workflowPromptOverrides.workflowId, id)); + await layer.db.delete(schema.project.workflowSettings).where(and( + 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.