Skip to content

Commit 552ce68

Browse files
committed
fix(workflows): guard restore against concurrent workspace archive
1 parent 30bc9d8 commit 552ce68

2 files changed

Lines changed: 44 additions & 0 deletions

File tree

‎apps/sim/lib/projects/__integration__/foundation.integration.ts‎

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -45,6 +45,7 @@ import {
4545
lockWorkspaceProject,
4646
transferWorkspaceProjects,
4747
} from '@/lib/projects/membership'
48+
import { restoreWorkflow } from '@/lib/workflows/lifecycle'
4849
import { buildNewWorkflowRow } from '@/lib/workflows/persistence/new-workflow-row'
4950
import { createWorkspaceInTransaction } from '@/lib/workspaces/create'
5051
import { archiveWorkspace } from '@/lib/workspaces/lifecycle'
@@ -838,6 +839,46 @@ describe('Project foundation at the database and application boundary', () => {
838839
}
839840
)
840841

842+
check('workflow restore cannot overtake a concurrent Project archive', async () => {
843+
const f = await fixture()
844+
const workflowId = await addWorkflow(f.ids[0], f.ownerId)
845+
await db.update(workflow).set({ archivedAt: new Date() }).where(eq(workflow.id, workflowId))
846+
const held = createDeferred<number>()
847+
const release = createDeferred<void>()
848+
const archive = db.transaction(async (tx) => {
849+
const [connection] = await tx.execute<{ pid: number }>(sql`SELECT pg_backend_pid() AS pid`)
850+
await tx.select().from(workflow).where(eq(workflow.id, workflowId)).for('update')
851+
await archiveProjectInTransaction(tx, f.projectId)
852+
held.resolve(connection.pid)
853+
await release.promise
854+
})
855+
const blocker = await held.promise
856+
const restore = restoreWorkflow(workflowId, request).then(
857+
(result) => result,
858+
(error: unknown) => error
859+
)
860+
try {
861+
let waiting = false
862+
for (let attempt = 0; attempt < 100; attempt++) {
863+
const rows = await db.execute(sql`
864+
SELECT 1 FROM pg_stat_activity WHERE ${blocker} = ANY(pg_blocking_pids(pid))
865+
`)
866+
if (rows.length) {
867+
waiting = true
868+
break
869+
}
870+
await sleep(10)
871+
}
872+
expect(waiting).toBe(true)
873+
} finally {
874+
release.resolve()
875+
await archive
876+
}
877+
expect(await restore).toMatchObject({ code: 'not_found' })
878+
const [row] = await db.select().from(workflow).where(eq(workflow.id, workflowId))
879+
expect(row.archivedAt).not.toBeNull()
880+
})
881+
841882
check(
842883
'Project archive rolls back as a unit and retries without losing environments',
843884
async () => {

‎apps/sim/lib/workflows/lifecycle.ts‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ import { mcpPubSub } from '@/lib/mcp/pubsub'
2222
import { releaseWebhookPathClaims } from '@/lib/webhooks/path-claims'
2323
import { supersedeInFlightDeploymentOperations } from '@/lib/workflows/persistence/deployment-operations'
2424
import { getWorkflowById } from '@/lib/workflows/utils'
25+
import { lockActiveWorkspace } from '@/lib/workspaces/active-workspace'
2526

2627
const logger = createLogger('WorkflowLifecycle')
2728

@@ -178,6 +179,8 @@ export async function restoreWorkflow(
178179
const archivedAt = existingWorkflow.archivedAt
179180

180181
await db.transaction(async (tx) => {
182+
if (existingWorkflow.workspaceId) await lockActiveWorkspace(tx, existingWorkflow.workspaceId)
183+
181184
await tx
182185
.update(workflow)
183186
.set({

0 commit comments

Comments
 (0)