mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-09 06:15:21 +02:00
test: observe assignment lock waits with a two-connection pool
Reuse the lock holder connection for inspection, bound the fixture to two connections, and propagate early transaction failures. Co-Authored-By: Paperclip <noreply@paperclip.ing>
This commit is contained in:
1 parent
70dd278f56
commit
e30a15cc05
1 file changed
+10
-5
@@ -18,7 +18,7 @@ const support = await getEmbeddedPostgresTestSupport();
|
||||
let db: ReturnType<typeof createDb>;
|
||||
beforeAll(async () => {
|
||||
database = await startEmbeddedPostgresTestDatabase("paperclip-human-work-");
|
||||
db = createDb(database.connectionString);
|
||||
db = createDb(database.connectionString, { maxConnections: 2 });
|
||||
}, 30_000);
|
||||
afterAll(async () => { await database?.cleanup(); }, 60_000);
|
||||
|
||||
@@ -121,14 +121,19 @@ const support = await getEmbeddedPostgresTestSupport();
|
||||
const f = await seed(); const issue = await task(f);
|
||||
let release!: () => void;
|
||||
let locked!: () => void;
|
||||
let observeLock!: (pid: number) => Promise<boolean>;
|
||||
const held = new Promise<void>((resolve) => { release = resolve; });
|
||||
const ready = new Promise<void>((resolve) => { locked = resolve; });
|
||||
const commentWrite = db.transaction(async (tx) => {
|
||||
await issueService(db).addComment(issue.id, "Please proceed", { userId: "owner" }, { humanDirectedByUserId: "owner" }, tx);
|
||||
observeLock = async (pid) => {
|
||||
const [state] = await tx.execute(sql`select cardinality(pg_blocking_pids(${pid}::int)) > 0 as blocked`);
|
||||
return state.blocked === true;
|
||||
};
|
||||
locked();
|
||||
await held;
|
||||
});
|
||||
await ready;
|
||||
await Promise.race([ready, commentWrite]);
|
||||
let started!: (pid: number) => void;
|
||||
const waiter = new Promise<number>((resolve) => { started = resolve; });
|
||||
const reassignment = db.transaction(async (tx) => {
|
||||
@@ -137,13 +142,13 @@ const support = await getEmbeddedPostgresTestSupport();
|
||||
await tx.update(issues).set({ assigneeAgentId: f.other.id }).where(eq(issues.id, issue.id));
|
||||
});
|
||||
try {
|
||||
const pid = await waiter;
|
||||
const pid = await Promise.race([waiter, reassignment.then(() => { throw new Error("Reassignment ended before publishing its backend"); })]);
|
||||
const deadline = Date.now() + 5_000;
|
||||
let blocked = false;
|
||||
// Observe an actual PostgreSQL lock wait, not merely two queued promises.
|
||||
while (Date.now() < deadline) {
|
||||
const [state] = await db.execute(sql`select cardinality(pg_blocking_pids(${pid}::int)) > 0 as blocked`);
|
||||
if (state.blocked === true) { blocked = true; break; }
|
||||
// Reuse the lock holder's connection; a third pool slot is unnecessary.
|
||||
if (await observeLock(pid)) { blocked = true; break; }
|
||||
await new Promise((resolve) => setTimeout(resolve, 10));
|
||||
}
|
||||
expect(blocked).toBe(true);
|
||||
|
||||
Reference in new issue
Block a user