Skip to content

Commit 85b96fb

Browse files
authored
fix(forks): route new workflows through one row builder and harden lineage locking (#8406)
* fix(forks): route new workflows through one row builder and harden lineage locking * fix(workflows): read the fork-sync policy inside each create transaction
1 parent c47ed6c commit 85b96fb

40 files changed

Lines changed: 583 additions & 980 deletions

File tree

‎apps/docs/content/docs/platform/enterprise/forks.mdx‎

Lines changed: 8 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -29,7 +29,7 @@ On Sim Cloud, your organization may also need the feature turned on for your acc
2929

3030
### 1. Open Forks
3131

32-
Go to **Settings → Organization → Workspace forks** in the workspace you want to fork from (or manage).
32+
Go to **Settings → Workspace → Workspace forks** in the workspace you want to fork from (or manage).
3333

3434
<Image src="/static/enterprise/forks-list.png" alt="Workspace Forks settings page showing Parent and Forks sections with Docs, See activity, and Create fork actions" width={900} height={369} />
3535

@@ -62,7 +62,7 @@ Click **Fork**. The child workspace is created immediately. Deployed workflows l
6262

6363
### 3. Open the parent edge (from the child)
6464

65-
Open the **child** workspace → **Settings → Organization → Workspace forks**. On the **Parent** row, open the menu and choose **Edit mappings**.
65+
Open the **child** workspace → **Settings → Workspace → Workspace forks**. On the **Parent** row, open the menu and choose **Edit mappings**.
6666

6767
Child rows (when you are on the parent) only offer **Open workspace** and **Disconnect** — mapping and sync are owned by the child configuring how it relates to its parent.
6868

@@ -136,14 +136,14 @@ Above the list, **Sync new workflows by default** decides where a **newly create
136136

137137
| Setting | A new workflow… |
138138
|---------|-----------------|
139-
| **On** (default) | joins fork sync — it arrives checked and syncs as soon as you deploy it |
140-
| **Off** | starts outside fork sync — it arrives unchecked and only syncs after you check it |
139+
| **Sync** (default) | joins fork sync — it arrives checked and syncs as soon as you deploy it |
140+
| **Don't sync** | starts outside fork sync — it arrives unchecked and only syncs after you check it |
141141

142142
Three things to know:
143143

144-
- **It applies to the whole fork lineage.** The toggle writes every workspace in the lineage — the root, every ancestor, every descendant — so a parent and its forks can never disagree about what "new" means. Any workspace admin in the lineage can change it, and each member gets its own audit entry naming the workspace the change came from. A new fork inherits the value at creation.
144+
- **It applies to the whole fork lineage.** The toggle writes every workspace in the lineage — the root, every ancestor, every descendant — so a parent and its forks can never disagree about what "new" means. Any workspace admin in the lineage can change it, and each workspace whose value changes gets its own audit entry naming the workspace the change came from. A new fork inherits the value at creation.
145145
- **It is forward-only.** Flipping it never moves an existing workflow in or out of sync. The checkbox list above stays the record of what syncs.
146-
- **"New" means genuinely new.** Creating, duplicating, or importing a workflow takes this setting, as does the blank starter workflow a fork gets when there is nothing to copy. A workflow that arrives as a **copy** — from a fork, or from a push or pull — inherits its source's own checkbox instead, so a workflow you deliberately synced never lands unsynced in the child.
146+
- **"New" means genuinely new.** Creating, duplicating, or importing a workflow takes this setting, as does the blank starter workflow a fork gets when there is nothing to copy. A workflow that arrives as a **copy** — from a fork, or from a push or pull — ignores this setting. Only synced workflows are copied, and they always arrive synced, so a workflow you deliberately synced never lands unsynced on the other side.
147147

148148
**Example:** a template workspace turns this off so every scratch workflow the team creates stays local, then checks only the handful meant to reach the forks.
149149

@@ -399,7 +399,7 @@ Schedules, webhooks, and triggers are not live in the child until you **deploy**
399399
{ question: "Why is Sync greyed out?", answer: "Usually a blocking reference, an unmapped credential or secret, or a required dependent field (label, channel, document, …) still empty. Open Blocking sync and the mapping sections — each row explains what to fix. Sync also stays disabled while details are loading or if loading failed (reload the page)." },
400400
{ question: "Is sync a merge?", answer: "No. Deploy is like a commit; sync is a force push or force pull of deployed workflows onto the target. Use Rollback only for the last sync into a workspace, and remember copied resources may remain." },
401401
{ question: "Who can disconnect a fork I cannot open?", answer: "Any admin on your side of the edge. Disconnect does not require access to the other workspace — so you are not stuck if the other side lost membership." },
402-
{ question: "I deployed a new workflow and sync ignored it. Why?", answer: "Sync new workflows by default is off for this fork lineage, so the workflow was created outside fork sync. Open Settings → Organization → Workspace forks and check it under Synced workflows. Turning the toggle back on only affects workflows created after that — it never moves an existing one." },
402+
{ question: "I deployed a new workflow and sync ignored it. Why?", answer: "Sync new workflows by default is off for this fork lineage, so the workflow was created outside fork sync. Open Settings → Workspace → Workspace forks and check it under Synced workflows. Turning the toggle back on only affects workflows created after that — it never moves an existing one." },
403403
{ question: "Does turning Sync new workflows by default off stop my current syncs?", answer: "No. It is forward-only and never rewrites an existing workflow's checkbox, so everything already synced keeps syncing. It also applies to every workspace in the fork lineage, not just the one you changed it from." }
404404
]} />
405405

@@ -413,4 +413,4 @@ Self-hosted deployments turn Forks on with an environment variable instead of th
413413
|----------|-------------|
414414
| `FORKING_ENABLED`, `NEXT_PUBLIC_FORKING_ENABLED` | Enables workspace forking when billing is not used as the entitlement gate |
415415

416-
Once enabled, use the same **Settings → Organization → Workspace forks** UI as Sim Cloud. Only workspace admins can manage forks.
416+
Once enabled, use the same **Settings → Workspace → Workspace forks** UI as Sim Cloud. Only workspace admins can manage forks.

‎apps/sim/app/api/superuser/import-workflow/route.ts‎

Lines changed: 12 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -12,13 +12,13 @@ import { loadCopilotChatMessages } from '@/lib/mothership/chat/lifecycle'
1212
import { appendCopilotChatMessages } from '@/lib/mothership/chat/messages-store'
1313
import { verifyEffectiveSuperUser } from '@/lib/permissions/super-user'
1414
import { parseWorkflowJson } from '@/lib/workflows/operations/import-export'
15+
import { buildNewWorkflowRow } from '@/lib/workflows/persistence/new-workflow-row'
1516
import {
1617
loadWorkflowFromNormalizedTables,
1718
saveWorkflowToNormalizedTables,
1819
} from '@/lib/workflows/persistence/utils'
1920
import { sanitizeForExport } from '@/lib/workflows/sanitization/json-sanitizer'
2021
import { deduplicateWorkflowName } from '@/lib/workflows/utils'
21-
import { resolveForkSyncExclusionForNewWorkflow } from '@/ee/workspace-forking/lib/sync-default'
2222

2323
const logger = createLogger('SuperUserImportWorkflow')
2424

@@ -130,30 +130,23 @@ export const POST = withRouteHandler(async (request: NextRequest) => {
130130

131131
// Create new workflow record
132132
const newWorkflowId = generateId()
133-
const now = new Date()
134133
const dedupedName = await deduplicateWorkflowName(
135134
`[Debug Import] ${sourceWorkflow.name}`,
136135
targetWorkspaceId,
137136
null
138137
)
139138

140-
await db.insert(workflow).values({
141-
id: newWorkflowId,
142-
userId: session.user.id,
143-
workspaceId: targetWorkspaceId,
144-
folderId: null,
145-
name: dedupedName,
146-
description: sourceWorkflow.description,
147-
lastSynced: now,
148-
createdAt: now,
149-
updatedAt: now,
150-
isDeployed: false, // Never copy deployment status
151-
runCount: 0,
152-
variables: sourceWorkflow.variables || {},
153-
// An imported workflow is a NEW workflow in the target workspace, so it takes that
154-
// workspace's fork-sync policy rather than the column default.
155-
forkSyncExcluded: await resolveForkSyncExclusionForNewWorkflow(db, targetWorkspaceId),
156-
})
139+
await db.insert(workflow).values(
140+
await buildNewWorkflowRow(db, {
141+
id: newWorkflowId,
142+
userId: session.user.id,
143+
workspaceId: targetWorkspaceId,
144+
folderId: null,
145+
name: dedupedName,
146+
description: sourceWorkflow.description,
147+
variables: sourceWorkflow.variables || {},
148+
})
149+
)
157150

158151
// Save using existing persistence logic
159152
const saveResult = await saveWorkflowToNormalizedTables(newWorkflowId, importedData, {

‎apps/sim/app/api/v1/admin/workflows/import/route.ts‎

Lines changed: 11 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@ import { adminV1ImportWorkflowContract } from '@/lib/api/contracts/v1/admin'
3030
import { parseRequest } from '@/lib/api/server'
3131
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
3232
import { parseWorkflowJson } from '@/lib/workflows/operations/import-export'
33+
import { buildNewWorkflowRow } from '@/lib/workflows/persistence/new-workflow-row'
3334
import { prepareWorkflowStateForPersistence } from '@/lib/workflows/persistence/prepare-state'
3435
import { saveWorkflowToNormalizedTables } from '@/lib/workflows/persistence/utils'
3536
import { deduplicateWorkflowName } from '@/lib/workflows/utils'
@@ -41,7 +42,6 @@ import {
4142
notFoundResponse,
4243
} from '@/app/api/v1/admin/responses'
4344
import { extractWorkflowMetadata, type WorkflowImportRequest } from '@/app/api/v1/admin/types'
44-
import { resolveForkSyncExclusionForNewWorkflow } from '@/ee/workspace-forking/lib/sync-default'
4545

4646
const logger = createLogger('AdminWorkflowImportAPI')
4747

@@ -113,27 +113,18 @@ export const POST = withRouteHandler(
113113
)
114114

115115
const workflowId = generateId()
116-
const now = new Date()
117116
const dedupedName = await deduplicateWorkflowName(workflowName, workspaceId, folderId || null)
118117

119-
await db.insert(workflow).values({
120-
id: workflowId,
121-
userId: workspaceData.ownerId,
122-
workspaceId,
123-
folderId: folderId || null,
124-
name: dedupedName,
125-
description: workflowDescription,
126-
lastSynced: now,
127-
createdAt: now,
128-
updatedAt: now,
129-
isDeployed: false,
130-
runCount: 0,
131-
variables: {},
132-
// An imported workflow is a NEW workflow in this workspace, so it takes the
133-
// workspace's fork-sync policy. Without this it lands on the column default and
134-
// silently joins fork sync in a workspace that opted out.
135-
forkSyncExcluded: await resolveForkSyncExclusionForNewWorkflow(db, workspaceId),
136-
})
118+
await db.insert(workflow).values(
119+
await buildNewWorkflowRow(db, {
120+
id: workflowId,
121+
userId: workspaceData.ownerId,
122+
workspaceId,
123+
folderId: folderId || null,
124+
name: dedupedName,
125+
description: workflowDescription,
126+
})
127+
)
137128

138129
/**
139130
* Same normalization the editor and the v1 import API run, via the one

‎apps/sim/app/api/v1/admin/workspaces/[id]/import/route.ts‎

Lines changed: 11 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -45,6 +45,7 @@ import {
4545
extractWorkflowsFromZip,
4646
parseWorkflowJson,
4747
} from '@/lib/workflows/operations/import-export'
48+
import { buildNewWorkflowRow } from '@/lib/workflows/persistence/new-workflow-row'
4849
import { prepareWorkflowStateForPersistence } from '@/lib/workflows/persistence/prepare-state'
4950
import { saveWorkflowToNormalizedTables } from '@/lib/workflows/persistence/utils'
5051
import { deduplicateWorkflowName } from '@/lib/workflows/utils'
@@ -62,7 +63,6 @@ import type {
6263
WorkspaceImportRequest,
6364
WorkspaceImportResponse,
6465
} from '@/app/api/v1/admin/types'
65-
import { resolveForkSyncExclusionForNewWorkflow } from '@/ee/workspace-forking/lib/sync-default'
6666

6767
const logger = createLogger('AdminWorkspaceImportAPI')
6868

@@ -348,27 +348,18 @@ async function importSingleWorkflow(
348348
}
349349

350350
const workflowId = generateId()
351-
const now = new Date()
352351
const dedupedName = await deduplicateWorkflowName(workflowName, workspaceId, targetFolderId)
353352

354-
await db.insert(workflow).values({
355-
id: workflowId,
356-
userId: ownerId,
357-
workspaceId,
358-
folderId: targetFolderId,
359-
name: dedupedName,
360-
description: workflowData.metadata?.description || 'Imported via Admin API',
361-
lastSynced: now,
362-
createdAt: now,
363-
updatedAt: now,
364-
isDeployed: false,
365-
runCount: 0,
366-
variables: {},
367-
// An imported workflow is a NEW workflow in this workspace, so it takes the
368-
// workspace's fork-sync policy. Without this it lands on the column default and
369-
// silently joins fork sync in a workspace that opted out.
370-
forkSyncExcluded: await resolveForkSyncExclusionForNewWorkflow(db, workspaceId),
371-
})
353+
await db.insert(workflow).values(
354+
await buildNewWorkflowRow(db, {
355+
id: workflowId,
356+
userId: ownerId,
357+
workspaceId,
358+
folderId: targetFolderId,
359+
name: dedupedName,
360+
description: workflowData.metadata?.description || 'Imported via Admin API',
361+
})
362+
)
372363

373364
/**
374365
* Same normalization the editor, the v1 import API and the single-workflow

‎apps/sim/app/api/workspaces/[id]/fork/sync-default/route.ts‎

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -13,10 +13,8 @@ export const PUT = defineInternalJsonRoute({
1313
auth: internalSessionAuth,
1414
operation: forkOperations.syncDefault,
1515
/**
16-
* Rated, unlike its sibling fork routes. This is the one that writes workspaces the
17-
* caller may not administer, under the feature's coarsest advisory lock, so an admin of
18-
* any single lineage member could otherwise loop it and starve fork creation across the
19-
* whole lineage.
16+
* Rate-limited, unlike sibling fork routes: it writes the whole lineage under the coarsest
17+
* fork lock, so looping it could starve fork creation lineage-wide.
2018
*/
2119
rateLimit: internalRateLimits.user({ bucketName: 'workspace-fork-sync-default' }),
2220
errorPolicy: internalForkErrorPolicy,

‎apps/sim/ee/workspace-forking/application/lineage-details.ts‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,10 @@
11
import { db } from '@sim/db'
22
import { workspace } from '@sim/db/schema'
33
import { eq } from 'drizzle-orm'
4+
import { readForkSyncNewWorkflowsExcluded } from '@/lib/workflows/persistence/new-workflow-row'
45
import { getEffectiveWorkspacePermission } from '@/lib/workspaces/permissions/utils'
56
import { getForkChildren, getForkParent } from '@/ee/workspace-forking/lib/lineage/lineage'
67
import { getUndoableRunForTarget } from '@/ee/workspace-forking/lib/promote/promote-run-store'
7-
import { resolveForkSyncExclusionForNewWorkflow } from '@/ee/workspace-forking/lib/sync-default'
88

99
/**
1010
* Annotates a lineage node with whether the viewer holds any access to it (explicit
@@ -40,7 +40,7 @@ export const getWorkspaceForkLineageDetails = defineForkUseCase({
4040
getForkChildren(workspaceId),
4141
getUndoableRunForTarget(db, workspaceId),
4242
// Lineage-uniform, so this workspace's own value is the lineage's value.
43-
resolveForkSyncExclusionForNewWorkflow(db, workspaceId),
43+
readForkSyncNewWorkflowsExcluded(db, workspaceId),
4444
])
4545

4646
const [parent, children] = await Promise.all([

‎apps/sim/ee/workspace-forking/application/operations.ts‎

Lines changed: 2 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -99,12 +99,8 @@ export const forkOperations = {
9999
oauthScope: 'api:write',
100100
}),
101101
/**
102-
* Admin on the CALLING workspace is sufficient, and the write then fans out to every
103-
* ancestor and descendant, because the default is meaningless unless it is uniform
104-
* across a lineage. Flipping it to "sync new workflows" restores the historical
105-
* behaviour rather than granting anything new, and it never moves an existing workflow
106-
* in or out of sync - so each member records its own audit entry rather than the write
107-
* being restricted to one workspace.
102+
* Admin on the calling workspace is sufficient; the write fans out to the whole lineage
103+
* because the default must be uniform, and it never moves an existing workflow.
108104
*
109105
* permission-group-exempt: the new-workflow fork-sync default is workspace configuration governed by the admin role.
110106
*/

‎apps/sim/ee/workspace-forking/application/recovery-and-mappings.ts‎

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -52,10 +52,8 @@ export const updateWorkspaceForkMappings = defineForkUseCase<
5252
input.direction === 'push' ? input.otherWorkspaceId : input.workspaceId
5353
return db.transaction(async (tx) => {
5454
await setForkLockTimeout(tx)
55-
// Rank 4 - see the rank table on `acquireForkLineageLock`. Unlike promote and
56-
// rollback this takes no rank-3 target lock: it rewrites only this edge's mapping
57-
// rows, never the target's workflows, so nothing contends with a sync into the
58-
// target. Skipping a higher rank is not an ordering violation.
55+
// Rank 4 - see the rank table on `acquireForkLineageLock`. No target lock: this
56+
// rewrites only the edge's mapping rows.
5957
await acquireForkEdgeLock(tx, edge.childWorkspaceId)
6058
const [currentEdge] = await tx
6159
.select({ parentId: workspace.forkedFromWorkspaceId })

‎apps/sim/ee/workspace-forking/application/revision.ts‎

Lines changed: 3 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -143,11 +143,8 @@ export async function loadForkPreviewRevision(
143143
/**
144144
* Locks normalized graph rows as well as workflow metadata, including realtime-only writes.
145145
*
146-
* Rank 5 - the heaviest acquirer in the fork module, and the one the rank table on
147-
* `acquireForkLineageLock` exists for. It takes `FOR UPDATE` on the `workspace` rows, so
148-
* any caller that also needs the rank-2 lineage lock must take that one FIRST; doing it
149-
* the other way round deadlocks against `unlinkForkEdge`, which holds the lineage key and
150-
* then updates the same `workspace` row.
146+
* Rank 5 - see the rank table on `acquireForkLineageLock`. Takes `FOR UPDATE` on `workspace`
147+
* rows, so a caller needing the rank-2 lineage lock must take it first.
151148
*/
152149
export async function lockForkRevision(tx: DbTransaction, scope: ForkRevisionScope): Promise<void> {
153150
const workspaceIds = [
@@ -207,21 +204,11 @@ export async function assertForkSourceVersions(
207204
sourceWorkspaceId: string,
208205
expected: ReadonlyMap<string, { id: string; digest: string }>
209206
): Promise<void> {
210-
// Verify exactly the workflows that were ADMITTED, rather than re-deriving the source
211-
// predicate here. Re-deriving it duplicated `listDeployedWorkflows`'s filter, so the day
212-
// a caller admitted a different set - "Copy unsynced workflows" admits sync-excluded
213-
// workflows - this query returned fewer rows and every such fork failed on a phantom
214-
// size mismatch. Keying off `expected` cannot drift from the admitted set by construction.
215-
if (expected.size === 0) return
216-
const admittedIds = sql.join(
217-
[...expected.keys()].map((id) => sql`${id}`),
218-
sql`, `
219-
)
220207
const rows = await tx.execute<{ workflowId: string; id: string; digest: string }>(sql`
221208
SELECT w.id AS "workflowId", d.id, md5(d.state::text) AS digest FROM ${workflow} w
222209
JOIN ${workflowDeploymentVersion} d ON d.workflow_id = w.id AND d.is_active = true
223210
WHERE w.workspace_id = ${sourceWorkspaceId} AND w.is_deployed = true
224-
AND w.archived_at IS NULL AND w.id IN (${admittedIds})
211+
AND w.archived_at IS NULL AND w.fork_sync_excluded = false
225212
`)
226213
if (
227214
rows.length !== expected.size ||

0 commit comments

Comments
 (0)