From 6de6a8103e61aaadefd8b836a4acc9ff9cb64ba8 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 06:16:36 -0700 Subject: [PATCH 1/2] Add concurrency regression tests for upsertPolicy Reconstructs the real-Postgres finding that two concurrent upsertPolicy calls patching different fields of the same tenant's policy can silently revert each other, for both an existing row and a brand-new tenant's first write. A plain Promise.all of two real calls does not reliably force the interleaving on a fast local connection, so the existing-row case uses a third connection that takes the row's lock first and holds it open while both real calls start and queue up behind it, guaranteeing the overlap. Mirrors the existing consumePendingInvite concurrency test in the same file. --- .../access-policy/test/store.drizzle.test.ts | 98 +++++++++++++++++++ 1 file changed, 98 insertions(+) diff --git a/packages/access-policy/test/store.drizzle.test.ts b/packages/access-policy/test/store.drizzle.test.ts index 54d25d962..3a24ee877 100644 --- a/packages/access-policy/test/store.drizzle.test.ts +++ b/packages/access-policy/test/store.drizzle.test.ts @@ -180,6 +180,104 @@ describeIfDb("createDrizzleAccessPolicyStore", () => { } }); + test("upsertPolicy is atomic: two concurrent patches to different fields on an existing row both land, neither reverts the other", async () => { + const setupSql = postgres(scratchUrl, { max: 1 }); + try { + const store = createDrizzleAccessPolicyStore(drizzle(setupSql)); + await store.upsertPolicy("tnt_race_existing", { + selfSignup: "off", + allowedDomains: [], + tenancyCreation: "owners", + }); + } finally { + await setupSql.end(); + } + + // A plain `Promise.all` of two real calls does not reliably force + // the worst-case interleaving on a fast local connection: one + // call's whole read-modify-write often finishes before the other's + // read even starts, so the two never actually overlap. Instead, a + // third connection takes the row's lock first and holds it open + // while both real `upsertPolicy` calls start and queue up behind + // it — releasing it then guarantees both calls' reads had to + // happen without seeing the other's write yet, exactly the + // interleaving that silently reverted one admin's change. + const holderSql = postgres(scratchUrl, { max: 1 }); + const sqlA = postgres(scratchUrl, { max: 1 }); + const sqlB = postgres(scratchUrl, { max: 1 }); + try { + const storeA = createDrizzleAccessPolicyStore(drizzle(sqlA)); + const storeB = createDrizzleAccessPolicyStore(drizzle(sqlB)); + + let holderReady: () => void = () => undefined; + const holderHasLock = new Promise((resolve) => { + holderReady = resolve; + }); + let releaseHolder: () => void = () => undefined; + const releaseSignal = new Promise((resolve) => { + releaseHolder = resolve; + }); + const holderTx = holderSql.begin(async (tx) => { + await tx`select * from access_policy.policy where tenant_id = 'tnt_race_existing' for update`; + holderReady(); + await releaseSignal; + }); + + await holderHasLock; + + const racers = Promise.all([ + storeA.upsertPolicy("tnt_race_existing", { selfSignup: "open" }), + storeB.upsertPolicy("tnt_race_existing", { + allowedDomains: ["acme.example"], + }), + ]); + // Give both calls time to actually issue their row-locking read + // and start queuing behind the holder before it releases. + await new Promise((resolve) => setTimeout(resolve, 100)); + + releaseHolder(); + await holderTx; + await racers; + + const policy = await storeA.getPolicy("tnt_race_existing"); + expect(policy.selfSignup).toBe("open"); + expect(policy.allowedDomains).toEqual(["acme.example"]); + expect(policy.tenancyCreation).toBe("owners"); + } finally { + await holderSql.end(); + await sqlA.end(); + await sqlB.end(); + } + }); + + test("upsertPolicy is atomic: two concurrent first-writes for a brand-new tenant both land, neither reverts the other", async () => { + const sql = postgres(scratchUrl, { max: 5 }); + try { + const store = createDrizzleAccessPolicyStore(drizzle(sql)); + + // No row exists yet for this tenant, so both calls race the + // create path too — the ensure-row-then-lock step inside + // `upsertPolicy` has to serialize this case as well, not only + // the existing-row case above. + await Promise.all([ + store.upsertPolicy("tnt_race_new", { selfSignup: "open" }), + store.upsertPolicy("tnt_race_new", { + allowedDomains: ["acme.example"], + }), + ]); + + const policy = await store.getPolicy("tnt_race_new"); + expect(policy.selfSignup).toBe("open"); + expect(policy.allowedDomains).toEqual(["acme.example"]); + + const rows = + await sql`select count(*)::int as count from access_policy.policy where tenant_id = 'tnt_race_new'`; + expect(rows[0]?.["count"]).toBe(1); + } finally { + await sql.end(); + } + }); + test("pending invites: a domain match is found for any email on that domain", async () => { const sql = postgres(scratchUrl, { max: 1 }); try { From 1b61def3dfa92f1c92de400d208ed15384f076c7 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 06:16:43 -0700 Subject: [PATCH 2/2] Fix lost-update race in access-policy upsertPolicy upsertPolicy read the current policy row and wrote a merged result back as two separate round trips, with no transaction, row lock, or optimistic-concurrency guard. Two concurrent PATCHes touching different fields of the same tenant's policy could both read the same baseline, so whichever write landed second silently reverted the other's field. Wraps upsertPolicy in a transaction with SELECT ... FOR UPDATE rather than an optimistic version-stamp check: two transactions starting in the same wall-clock tick could still both pass a version check, which is exactly the bug this closes. Because the policy row may not exist yet, an INSERT ... ON CONFLICT DO NOTHING guarantees a row exists before the row lock is taken, serializing concurrent first-writes for a brand-new tenant on the insert's unique index instead of racing them. --- packages/access-policy/src/store.ts | 84 ++++++++++++++++++----------- 1 file changed, 53 insertions(+), 31 deletions(-) diff --git a/packages/access-policy/src/store.ts b/packages/access-policy/src/store.ts index 8f35bf09b..926534ee2 100644 --- a/packages/access-policy/src/store.ts +++ b/packages/access-policy/src/store.ts @@ -110,41 +110,63 @@ export function createDrizzleAccessPolicyStore< return row !== undefined; }, + // Read-modify-write under a transaction with a row lock rather than + // an optimistic version-stamp check: two transactions starting in + // the same wall-clock tick could still both pass a version check, + // which is exactly the lost-update bug this closes. + // + // Because this method's row may not exist yet (it is an upsert), + // `SELECT ... FOR UPDATE` alone has nothing to lock for a brand-new + // tenant. The `INSERT ... ON CONFLICT DO NOTHING` below guarantees + // a row first: two concurrent first-writes for the same tenant + // serialize on that insert's unique-index conflict — the loser + // blocks until the winner's transaction commits, then no-ops and + // its own subsequent `SELECT ... FOR UPDATE` sees the winner's + // committed row rather than racing it. async upsertPolicy(tenantId, patch) { - const current = await db - .select() - .from(policy) - .where(eq(policy.tenantId, tenantId)); - const existing = - current[0] === undefined - ? DEFAULT_ACCESS_POLICY - : resolveAccessPolicy(current[0]); - const next: AccessPolicy = { - selfSignup: patch.selfSignup ?? existing.selfSignup, - allowedDomains: patch.allowedDomains ?? existing.allowedDomains, - tenancyCreation: patch.tenancyCreation ?? existing.tenancyCreation, - }; - const now = new Date(); - await db - .insert(policy) - .values({ - tenantId, - selfSignup: next.selfSignup, - allowedDomains: serializeAllowedDomains(next.allowedDomains), - tenancyCreation: next.tenancyCreation, - createdAt: now, - updatedAt: now, - }) - .onConflictDoUpdate({ - target: policy.tenantId, - set: { + return db.transaction(async (tx) => { + const now = new Date(); + await tx + .insert(policy) + .values({ + tenantId, + selfSignup: DEFAULT_ACCESS_POLICY.selfSignup, + allowedDomains: serializeAllowedDomains( + DEFAULT_ACCESS_POLICY.allowedDomains, + ), + tenancyCreation: DEFAULT_ACCESS_POLICY.tenancyCreation, + createdAt: now, + updatedAt: now, + }) + .onConflictDoNothing({ target: policy.tenantId }); + + const [current] = await tx + .select() + .from(policy) + .where(eq(policy.tenantId, tenantId)) + .for("update"); + if (current === undefined) { + throw new Error( + `upsertPolicy: no access_policy.policy row for tenant ${tenantId} after ensuring one exists`, + ); + } + const existing = resolveAccessPolicy(current); + const next: AccessPolicy = { + selfSignup: patch.selfSignup ?? existing.selfSignup, + allowedDomains: patch.allowedDomains ?? existing.allowedDomains, + tenancyCreation: patch.tenancyCreation ?? existing.tenancyCreation, + }; + await tx + .update(policy) + .set({ selfSignup: next.selfSignup, allowedDomains: serializeAllowedDomains(next.allowedDomains), tenancyCreation: next.tenancyCreation, - updatedAt: now, - }, - }); - return next; + updatedAt: new Date(), + }) + .where(eq(policy.tenantId, tenantId)); + return next; + }); }, async createPendingInvite(tenantId, input) {