Skip to content

Commit 523338a

Browse files
committed
fix(projects): retry writes after concurrent membership changes
1 parent 5121347 commit 523338a

2 files changed

Lines changed: 55 additions & 1 deletion

File tree

‎packages/db/migrations/0394_project_membership_enforcement.sql‎

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -142,8 +142,15 @@ BEGIN
142142
IF TG_OP <> 'INSERT' THEN previous_id := OLD.workspace_id; END IF;
143143
IF TG_OP <> 'DELETE' THEN next_id := NEW.workspace_id; END IF;
144144
END IF;
145-
SELECT array_agg(project_id) INTO owners FROM project_workspace WHERE workspace_id IN (previous_id, next_id);
145+
SELECT array_agg(project_id ORDER BY workspace_id) INTO owners FROM project_workspace WHERE workspace_id IN (previous_id, next_id);
146146
PERFORM project_contract_lock_projects(owners);
147+
IF owners IS DISTINCT FROM (
148+
SELECT array_agg(project_id ORDER BY workspace_id) FROM project_workspace
149+
WHERE workspace_id IN (previous_id, next_id)
150+
) THEN
151+
RAISE EXCEPTION 'Environment changed Projects while acquiring its lifecycle lock; retry the transaction'
152+
USING ERRCODE = '40001';
153+
END IF;
147154
END IF;
148155
IF TG_OP = 'DELETE' THEN RETURN OLD; END IF;
149156
RETURN NEW;

‎packages/db/scripts/project-contract.integration.ts‎

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ import { backfillProjects } from '@sim/db/project-backfill'
88
import { readTestDatabaseUrl } from '@sim/db/testing/test-infrastructure'
99
import { createDeferred } from '@sim/testing/helpers/deferred'
1010
import { getPostgresErrorCode } from '@sim/utils/errors'
11+
import { sleep } from '@sim/utils/helpers'
1112
import { generateId } from '@sim/utils/id'
1213
import { drizzle } from 'drizzle-orm/postgres-js'
1314
import { migrate } from 'drizzle-orm/postgres-js/migrator'
@@ -220,6 +221,52 @@ describe('Project expand/backfill/contract against PostgreSQL', () => {
220221
})
221222
})
222223

224+
it('retries a workflow write when its environment moves Projects while the writer waits', async () => {
225+
await database(async (sql) => {
226+
await seed(sql)
227+
await enforce(sql)
228+
const moved = createDeferred<void>()
229+
const release = createDeferred<void>()
230+
const transfer = sql.begin(async (tx) => {
231+
await tx`INSERT INTO project (id, name, owner_id) VALUES ('destination', 'Destination', 'owner')`
232+
await tx`UPDATE workspace SET forked_from_workspace_id = NULL WHERE id = 'child'`
233+
await tx`UPDATE project_workspace SET project_id = 'destination' WHERE workspace_id = 'child'`
234+
moved.resolve()
235+
await release.promise
236+
})
237+
await moved.promise
238+
const writer = sql
239+
.begin(async (tx) => {
240+
await tx`SET LOCAL application_name = 'project_contract_stale_writer'`
241+
await tx`INSERT INTO workflow VALUES ('racing', 'child', NULL)`
242+
})
243+
.then(
244+
() => null,
245+
(error: unknown) => error
246+
)
247+
try {
248+
let waiting = false
249+
for (let attempt = 0; attempt < 100; attempt++) {
250+
const rows =
251+
await sql`SELECT 1 FROM pg_stat_activity WHERE application_name = 'project_contract_stale_writer' AND wait_event_type = 'Lock'`
252+
if (rows.length) {
253+
waiting = true
254+
break
255+
}
256+
await sleep(10)
257+
}
258+
expect(waiting).toBe(true)
259+
} finally {
260+
release.resolve()
261+
}
262+
await transfer
263+
expect(getPostgresErrorCode(await writer)).toBe('40001')
264+
expect(await sql`SELECT 1 FROM workflow WHERE id = 'racing'`).toHaveLength(0)
265+
await sql`INSERT INTO workflow VALUES ('retried', 'child', NULL)`
266+
expect(await sql`SELECT 1 FROM workflow WHERE id = 'retried'`).toHaveLength(1)
267+
})
268+
})
269+
223270
it.each(['read committed', 'repeatable read'] as const)(
224271
'prevents concurrent last-environment removal under %s',
225272
async (isolation) => {

0 commit comments

Comments
 (0)