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:
7
.changeset/fn-8998-workflows-project-partition.md
Normal file
7
.changeset/fn-8998-workflows-project-partition.md
Normal 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.
|
||||
@@ -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.
|
||||
|
||||
|
||||
@@ -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");
|
||||
});
|
||||
});
|
||||
@@ -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<WorkflowR
|
||||
const rows = await layer.db
|
||||
.select()
|
||||
.from(schema.project.workflows)
|
||||
.where(projectScopeFor(schema.project.workflows.projectId, layer.projectId))
|
||||
.orderBy(asc(schema.project.workflows.createdAt));
|
||||
return rows.map(rowToWorkflowRow);
|
||||
}
|
||||
@@ -50,7 +51,19 @@ export async function getWorkflowRow(layer: AsyncDataLayer, id: string): Promise
|
||||
const rows = await layer.db
|
||||
.select()
|
||||
.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);
|
||||
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);
|
||||
}
|
||||
|
||||
@@ -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<string, string>();
|
||||
for (const row of workflowNameRows) {
|
||||
|
||||
@@ -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-<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", {
|
||||
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)`),
|
||||
|
||||
@@ -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<string> {
|
||||
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);
|
||||
|
||||
@@ -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.
|
||||
|
||||
Reference in New Issue
Block a user