|
1 | 1 | import { |
2 | 2 | backfillLegacyKnowledgeBaseWorkspaces, |
3 | 3 | createPostgresLegacyKnowledgeBaseWorkspaceStore, |
| 4 | + type LegacyKnowledgeBaseMoveOutcome, |
4 | 5 | selectLegacyKnowledgeBaseWorkspace, |
5 | 6 | } from '@sim/db/script-migrations/0013_backfill_legacy_knowledge_base_workspaces' |
| 7 | +import { sleep } from '@sim/utils/helpers' |
6 | 8 | import { generateId } from '@sim/utils/id' |
7 | 9 | import postgres, { type Sql } from 'postgres' |
8 | 10 | import { afterAll, beforeAll, beforeEach, describe, expect, it } from 'vitest' |
@@ -286,6 +288,71 @@ describe.runIf(Boolean(databaseUrl))('legacy KB workspace backfill in PostgreSQL |
286 | 288 | expect(await userBytes('owner')).toBe(300) |
287 | 289 | }) |
288 | 290 |
|
| 291 | + it('waits for an application-held workspace lock instead of aborting the migration', async () => { |
| 292 | + await workspace('destination') |
| 293 | + await kb() |
| 294 | + await document('kb', 100) |
| 295 | + let markWriterReady!: (pid: number) => void |
| 296 | + const writerReady = new Promise<number>((resolve) => { |
| 297 | + markWriterReady = resolve |
| 298 | + }) |
| 299 | + let releaseWriter!: () => void |
| 300 | + const writerReleased = new Promise<void>((resolve) => { |
| 301 | + releaseWriter = resolve |
| 302 | + }) |
| 303 | + const writer = sql.begin(async (tx) => { |
| 304 | + await tx`SELECT id FROM workspace WHERE id = 'destination' FOR NO KEY UPDATE` |
| 305 | + const [connection] = await tx<{ pid: number }[]>`SELECT pg_backend_pid() AS pid` |
| 306 | + markWriterReady(connection.pid) |
| 307 | + await writerReleased |
| 308 | + }) |
| 309 | + let move: Promise<LegacyKnowledgeBaseMoveOutcome> | undefined |
| 310 | + try { |
| 311 | + const writerPid = await Promise.race([ |
| 312 | + writerReady, |
| 313 | + writer.then(() => { |
| 314 | + throw new Error('Workspace writer finished before the concurrency check') |
| 315 | + }), |
| 316 | + ]) |
| 317 | + let settled = false |
| 318 | + move = subject.moveCandidate('kb') |
| 319 | + void move.then( |
| 320 | + () => { |
| 321 | + settled = true |
| 322 | + }, |
| 323 | + () => { |
| 324 | + settled = true |
| 325 | + } |
| 326 | + ) |
| 327 | + let waiting = false |
| 328 | + const deadline = Date.now() + 2_000 |
| 329 | + while (!settled && Date.now() < deadline) { |
| 330 | + const [waiter] = await admin<{ waiting: boolean }[]>` |
| 331 | + SELECT EXISTS (SELECT 1 FROM pg_stat_activity |
| 332 | + WHERE ${writerPid} = ANY(pg_blocking_pids(pid)) AND wait_event_type = 'Lock') AS waiting |
| 333 | + ` |
| 334 | + if (waiter.waiting) { |
| 335 | + waiting = true |
| 336 | + break |
| 337 | + } |
| 338 | + await sleep(1) |
| 339 | + } |
| 340 | + expect(waiting, 'The migration must wait for the workspace writer').toBe(true) |
| 341 | + /** Real lock contention must outlast the former five-second timeout. */ |
| 342 | + await sleep(5_100) |
| 343 | + expect(settled).toBe(false) |
| 344 | + releaseWriter() |
| 345 | + await writer |
| 346 | + expect(await move).toBe('moved') |
| 347 | + expect((await scopedKb()).workspace_id).toBe('destination') |
| 348 | + expect(await userBytes('owner')).toBe(100) |
| 349 | + } finally { |
| 350 | + releaseWriter() |
| 351 | + await writer |
| 352 | + await move?.catch(() => undefined) |
| 353 | + } |
| 354 | + }, 15_000) |
| 355 | + |
289 | 356 | it('rolls back scope, names, and accounting together on failure and resumes safely', async () => { |
290 | 357 | await workspace('destination') |
291 | 358 | await kb() |
|
0 commit comments