From 8a8e5d7e2ebd77ead95d25de0bea0d01012fcee6 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 12:02:18 -0700 Subject: [PATCH 1/3] Close the GitHub connect grant/webhook-trigger duplication race startReviewingRepos's hasRepoGrant/hasWebhookTrigger checks are a fast path for a sequential retry (CL-7134), not a lock: two concurrent "start reviewing" calls for the same repo (a double-click, or a client retrying an in-flight request) can both read "not set up yet" before either write lands, minting two repo: grants and two live webhook triggers with different secrets for one repo. A duplicate grant is a permissions defect, not an untidy row -- the room's access to that repo now depends on which of two indistinguishable rows a later revoke happens to touch, so "disconnect this repo" can silently fail to actually revoke access. A duplicate trigger risks a code-review run firing twice, or the wrong secret being shown for the one GitHub actually has configured. Interchange's own `grant` table is never modified to fix this: no index, no constraint, not even a workbench-owned migration issuing DDL against it, and `mintRepoGrant` stays a plain insert -- the same shape Interchange's own POST /grants route writes. The concurrency guard lives entirely in a new workbench-owned table instead: `webhook_triggers.repo_review_lease` (own schema, no relationship to any Interchange table -- every workbench package already stores a platform id as a plain text column, never a cross-schema FK, and "repo" is a GitHub identity Interchange has no row for at all). `startReviewingRepos` acquires this lease first, per repo; only the winner runs the rest of that repo's body (hasRepoGrant/mintRepoGrant/hasWebhookTrigger/createWebhookTrigger, otherwise unchanged from CL-7134), so two concurrent callers can never both reach it for the same repo. The lease is released in a `finally` as soon as that body finishes, success or failure, so a legitimate retry is never blocked by its own prior attempt; a database-side compare-and-swap (`ON CONFLICT ... DO UPDATE ... WHERE`) also lets a lease older than two minutes be stolen, so a crash mid-work self-heals instead of permanently stranding that repo's setup (the actual claim a stranded lease makes is only ever "someone claimed this as of time T", never "the grant/trigger exist" -- CL-7213's own precedent). The paired webhook_trigger fix stays in our own schema too: a unique index on webhook_trigger(tenant_id, workflow_definition_id, name), reconciling any pre-existing duplicates first (keeps the oldest row per group). The store gains ensure() -- insert-first with onConflictDoNothing and a re-select on conflict, returning the real persisted secret -- as a method distinct from create(), which keeps its plain-insert contract so the generic management API never hands a caller someone else's trigger under a different name; a duplicate create() now 409s. apps/hub's createWebhookTrigger port binds to ensure() as defense-in-depth even though the lease already makes this call-site single-flight per repo. Regression tests, confirmed red before each fix existed: - packages/webhook-triggers/test/repo-review-lease.drizzle.test.ts: two concurrent acquire() calls settle on exactly one winner, a stale lease can be stolen, a fresh one cannot, and release() lets an immediate reacquire succeed. - packages/webhook-triggers/test/store.drizzle.test.ts: two concurrent ensure() calls settle on exactly one trigger, same id and secret. - packages/webhook-triggers/test/management-routes.test.ts: a second create() for an existing name/definition 409s. - packages/workflow-catalog/src/connect-github-setup.test.ts: two concurrent startReviewingRepos calls sharing lease-backed fakes mint exactly one grant and one trigger for the same repo -- the pure reconstruction of the audit's own (deleted) reproduction. --- apps/hub/src/index.ts | 29 +++- .../chat-ui/test/connect-github-flow.test.tsx | 10 ++ packages/webhook-triggers/src/index.ts | 12 +- .../webhook-triggers/src/management-routes.ts | 48 +++++- packages/webhook-triggers/src/migrations.ts | 49 ++++++ .../webhook-triggers/src/repo-review-lease.ts | 104 +++++++++++ packages/webhook-triggers/src/schema.ts | 52 +++++- packages/webhook-triggers/src/store.ts | 77 +++++++++ .../test/management-routes.test.ts | 22 +++ .../webhook-triggers/test/migrations.test.ts | 25 ++- .../test/repo-review-lease.drizzle.test.ts | 162 ++++++++++++++++++ .../test/store.drizzle.test.ts | 44 +++++ .../webhook-triggers/test/test-support.ts | 63 +++++-- .../connect-github-credential-link.test.ts | 4 + .../src/connect-github-routes.test.ts | 2 + .../src/connect-github-routes.ts | 21 +++ .../src/connect-github-setup.test.ts | 46 +++++ .../src/connect-github-setup.ts | 101 +++++++++-- scripts/checks/no-product-tenancy.ts | 13 +- 19 files changed, 822 insertions(+), 62 deletions(-) create mode 100644 packages/webhook-triggers/src/repo-review-lease.ts create mode 100644 packages/webhook-triggers/test/repo-review-lease.drizzle.test.ts diff --git a/apps/hub/src/index.ts b/apps/hub/src/index.ts index 9e31e6374..fe6c3ff89 100644 --- a/apps/hub/src/index.ts +++ b/apps/hub/src/index.ts @@ -163,6 +163,7 @@ import { createWorkflowCommandPlugin, } from "@corbits/commands"; import { + createDrizzleRepoReviewLeaseStore, createDrizzleWebhookTriggerStore, createWebhookIngressRoutes, createWebhookTriggerRoutes, @@ -1978,6 +1979,12 @@ export async function createHub(config: HubConfig) { db, credentialCipher, ); + // CL-7242: the sole concurrency backstop for the GitHub connect + // card's start-reviewing step -- see + // packages/webhook-triggers/src/repo-review-lease.ts for why this + // lives in our own schema rather than as any change to Interchange's + // `grant` table. + const repoReviewLeaseStore = createDrizzleRepoReviewLeaseStore(db); // Shared by every folded-run first-turn mail send below (webhook // triggers and routines alike) — a `CryptoProviderCache` is keyed by // instance id, which is globally unique across this hub regardless of @@ -2248,6 +2255,19 @@ export async function createHub(config: HubConfig) { }); return row?.id; }, + acquireRepoReviewLease: (tenantId, repo) => + repoReviewLeaseStore.acquire(tenantId, repo.name), + releaseRepoReviewLease: (tenantId, repo) => + repoReviewLeaseStore.release(tenantId, repo.name), + // hasRepoGrant/mintRepoGrant go through Interchange's native + // grants HTTP surface (never a direct `grant` table write -- + // see native-repo-grants.ts). That table carries no unique + // constraint over tenant/resource/action, so a bare read-then- + // POST here would itself be a duplicate-grant race; safe only + // because the caller in connect-github-routes.ts reaches this + // once `acquireRepoReviewLease` has already made this call-site + // single-flight per (tenant, repo) -- see + // packages/webhook-triggers/src/repo-review-lease.ts (CL-7242). hasRepoGrant: (tenantId, repo, cookies) => hasRepoGrantViaHttp(selfApi, tenantId, repo, cookies), mintRepoGrant: (tenantId, repo, cookies) => @@ -2258,7 +2278,14 @@ export async function createHub(config: HubConfig) { codeReviewDefinitionId, repo, ) => { - const row = await webhookTriggerStore.create({ + // `ensure`, not `create`: a concurrent "start reviewing" call + // for the same repo can race this one past `hasWebhookTrigger` + // above, and `webhook_trigger_tenant_definition_name_unique` + // (packages/webhook-triggers migration 0003, CL-7242) is what + // actually resolves that — the loser gets the winner's real + // row back instead of minting a second live trigger with a + // different secret. + const row = await webhookTriggerStore.ensure({ id: generateId("workflowRun"), tenantId, name: webhookTriggerName(repo), diff --git a/packages/chat-ui/test/connect-github-flow.test.tsx b/packages/chat-ui/test/connect-github-flow.test.tsx index e14809ace..d9f6c836e 100644 --- a/packages/chat-ui/test/connect-github-flow.test.tsx +++ b/packages/chat-ui/test/connect-github-flow.test.tsx @@ -67,7 +67,17 @@ function buildHarness() { let connected = false; let subscriber: ((state: ConnectGithubQuery) => void) | undefined; + const heldLeases = new Set(); + const setupPorts: ConnectGithubSetupPorts = { + async acquireRepoReviewLease(repo) { + if (heldLeases.has(repo.name)) return false; + heldLeases.add(repo.name); + return true; + }, + async releaseRepoReviewLease(repo) { + heldLeases.delete(repo.name); + }, async hasRepoGrant() { return false; }, diff --git a/packages/webhook-triggers/src/index.ts b/packages/webhook-triggers/src/index.ts index 88dcf45c3..c17ee0e9c 100644 --- a/packages/webhook-triggers/src/index.ts +++ b/packages/webhook-triggers/src/index.ts @@ -4,7 +4,17 @@ export { type ApplyWebhookTriggersMigrationsReport, type WebhookTriggersMigration, } from "./migrations"; -export { webhookTrigger, type WebhookTriggerRow } from "./schema"; +export { + webhookTrigger, + repoReviewLease, + type WebhookTriggerRow, + type RepoReviewLeaseRow, +} from "./schema"; +export { + createDrizzleRepoReviewLeaseStore, + type RepoReviewLeaseStore, + type RepoReviewLeaseDb, +} from "./repo-review-lease"; export { createDrizzleWebhookTriggerStore, type CreateWebhookTriggerInput, diff --git a/packages/webhook-triggers/src/management-routes.ts b/packages/webhook-triggers/src/management-routes.ts index 43f25f657..3c7a09f56 100644 --- a/packages/webhook-triggers/src/management-routes.ts +++ b/packages/webhook-triggers/src/management-routes.ts @@ -15,6 +15,7 @@ import { type } from "arktype"; import type { RequireGrant, TenantEnv } from "@intx/hub-api"; import { idResource } from "@intx/hub-api"; import { generateId } from "@intx/hub-common"; +import { pgErrorCode, PG_UNIQUE_VIOLATION } from "@intx/db"; import { generateWebhookSecret } from "./signature"; import type { WebhookTriggerRow } from "./schema"; @@ -24,6 +25,21 @@ const ErrorEnvelope = (code: string, message: string) => ({ error: { code, message }, }); +/** + * True for a Postgres unique-violation (`23505`) — the shape a duplicate + * `(tenant, workflow definition, name)` now raises through + * `0003_webhook_trigger_tenant_definition_name_unique`. `pgErrorCode` + * walks Drizzle's wrapped cause chain, since a real insert failure + * arrives as a `DrizzleQueryError` rather than the raw driver error; + * this package's in-memory test fake stamps `.code` directly to match. + * Never silently retried as an `ensure`: this route's `create` promises + * a genuinely new row, so a collision is reported to the caller as a + * conflict rather than handed back someone else's trigger. + */ +function isUniqueViolation(error: unknown): boolean { + return pgErrorCode(error) === PG_UNIQUE_VIOLATION; +} + const CreateTriggerBody = type({ name: "string", workflowDefinitionId: "string", @@ -95,15 +111,29 @@ export function createWebhookTriggerRoutes( const secret = generateWebhookSecret(); - const row = await deps.store.create({ - id: generateId("workflowRun"), - tenantId: tenant.id, - name: body.name, - workflowDefinitionId: body.workflowDefinitionId, - inputTemplate: body.inputTemplate, - secret, - createdBy: principal.id, - }); + let row: WebhookTriggerRow; + try { + row = await deps.store.create({ + id: generateId("workflowRun"), + tenantId: tenant.id, + name: body.name, + workflowDefinitionId: body.workflowDefinitionId, + inputTemplate: body.inputTemplate, + secret, + createdBy: principal.id, + }); + } catch (cause) { + if (isUniqueViolation(cause)) { + return c.json( + ErrorEnvelope( + "conflict", + "a trigger with this name already exists for this workflow definition", + ), + 409, + ); + } + throw cause; + } return c.json({ ...publicView(row), secret }, 201); }); diff --git a/packages/webhook-triggers/src/migrations.ts b/packages/webhook-triggers/src/migrations.ts index d2765e4eb..ec2043c85 100644 --- a/packages/webhook-triggers/src/migrations.ts +++ b/packages/webhook-triggers/src/migrations.ts @@ -49,6 +49,55 @@ export const webhookTriggersMigrations: readonly WebhookTriggersMigration[] = [ ON "webhook_triggers"."webhook_trigger" ("tenant_id"); `, }, + { + // CL-7242: startReviewingRepos's check-then-act (hasWebhookTrigger + // then createWebhookTrigger) reads (tenant_id, workflow_definition_id, + // name) with no atomic backstop, so two concurrent "start reviewing" + // calls for the same repo can both read "no trigger yet" and both + // insert -- two live triggers with the same name but different + // secrets. Reconcile first: a database already carrying the race's + // duplicates would otherwise fail CREATE UNIQUE INDEX. Keep the + // oldest row per tuple and delete the rest. This whole migration + // runs inside one transaction (applyWebhookTriggersMigrations wraps + // each entry in `sql.begin`), so the delete and the index build + // can't be split by a concurrent writer. + name: "0003_webhook_trigger_tenant_definition_name_unique", + sql: ` + DELETE FROM "webhook_triggers"."webhook_trigger" AS t + USING "webhook_triggers"."webhook_trigger" AS older + WHERE t.tenant_id = older.tenant_id + AND t.workflow_definition_id = older.workflow_definition_id + AND t.name = older.name + AND (older.created_at, older.id) < (t.created_at, t.id); + + CREATE UNIQUE INDEX IF NOT EXISTS "webhook_trigger_tenant_definition_name_unique" + ON "webhook_triggers"."webhook_trigger" ("tenant_id", "workflow_definition_id", "name"); + `, + }, + { + // CL-7242: the paired fix to 0003, in our own schema rather than + // Interchange's. `startReviewingRepos` acquires a short-lived + // lease on `(tenant_id, repo)` before its check-then-act body + // (hasRepoGrant/mintRepoGrant, hasWebhookTrigger/createWebhookTrigger) + // runs, so two concurrent calls for the same repo can never both + // enter that body -- only one can hold the lease at a time. The + // unique index must exist before any `ON CONFLICT (tenant_id, repo)` + // is issued against this table at runtime (Postgres requires a + // matching unique constraint/index for that clause, or the insert + // errors), so both land in this one migration, index second. + name: "0004_repo_review_lease", + sql: ` + CREATE TABLE IF NOT EXISTS "webhook_triggers"."repo_review_lease" ( + "id" text PRIMARY KEY, + "tenant_id" text NOT NULL, + "repo" text NOT NULL, + "leased_at" timestamptz NOT NULL DEFAULT now() + ); + + CREATE UNIQUE INDEX IF NOT EXISTS "repo_review_lease_tenant_repo_unique" + ON "webhook_triggers"."repo_review_lease" ("tenant_id", "repo"); + `, + }, ]; // Named distinctly from the platform's setup ledger and from any diff --git a/packages/webhook-triggers/src/repo-review-lease.ts b/packages/webhook-triggers/src/repo-review-lease.ts new file mode 100644 index 000000000..6fa2b18ef --- /dev/null +++ b/packages/webhook-triggers/src/repo-review-lease.ts @@ -0,0 +1,104 @@ +// Closes a check-then-act race in the GitHub connect card's +// start-reviewing step (CL-7242): `startReviewingRepos` +// (`@corbits/workflow-catalog`) loops over selected repos doing +// hasRepoGrant/mintRepoGrant then hasWebhookTrigger/createWebhookTrigger +// per repo, each a plain read followed by a conditional write with no +// atomic backstop. Two concurrent calls for the same repo (a +// double-click, or a client retrying an in-flight request) can both +// read "not set up yet" before either write lands, so both mint a +// grant and both create a trigger. +// +// A lease acquired before that per-repo body runs is the actual +// backstop: only the caller that wins the lease proceeds into +// hasRepoGrant/mintRepoGrant/hasWebhookTrigger/createWebhookTrigger, +// which stay exactly as CL-7134 left them (a fast path for a +// *sequential* retry after a mid-loop failure, safe now that they can +// never run concurrently for the same repo). The lease is released as +// soon as that body finishes (success or failure) so a legitimate +// retry is never blocked by its own prior attempt; a 2-minute +// staleness window is a crash-only backstop, in case a process dies +// before its `finally` can run. +// +// This table carries no schema-level relationship to Interchange's +// own `grant` table (or any other platform table): the workaround +// this closes lives entirely in workbench-owned state, and the actual +// source of truth for whether a repo's grant/trigger exist stays +// `hasRepoGrant`/`hasWebhookTrigger` against the real platform data, +// completely unchanged. The lease only ever asserts "someone claimed +// responsibility for this repo's setup as of `leasedAt`" — never +// "the grant/trigger exist" — so it can never assert something untrue +// the way a stale "done" marker could (CL-7213's own precedent). +import { and, eq, lt } from "drizzle-orm"; +import type { PostgresJsDatabase } from "drizzle-orm/postgres-js"; + +import { repoReviewLease } from "./schema"; + +export type RepoReviewLeaseDb< + TSchema extends Record = Record, +> = PostgresJsDatabase; + +/** Comfortably longer than a single repo's synchronous mint-and-create + * work should ever take, so a live lease is never mistaken for stale; + * short enough that a crashed holder self-heals well within a person + * re-clicking "Start reviewing" a few times. */ +const LEASE_STALE_AFTER_MS = 2 * 60 * 1000; + +export interface RepoReviewLeaseStore { + /** + * True if this call now holds the lease on `(tenantId, repo)` — + * either no lease existed, or the existing one is older than the + * staleness window and was stolen. False means another call + * currently holds (or very recently held) it; the caller must skip + * this repo rather than proceed. + */ + acquire(tenantId: string, repo: string): Promise; + /** Releases a held lease so an immediate legitimate retry (e.g. the + * next repo in a fresh `startReviewingRepos` call) never waits out + * the staleness window. Safe to call even if this caller never held + * the lease (e.g. `acquire` returned false) — a no-op in that case. */ + release(tenantId: string, repo: string): Promise; +} + +export function createDrizzleRepoReviewLeaseStore< + TSchema extends Record, +>(db: RepoReviewLeaseDb): RepoReviewLeaseStore { + return { + async acquire(tenantId, repo) { + const now = new Date(); + const staleBefore = new Date(now.getTime() - LEASE_STALE_AFTER_MS); + // Insert-first with a conditional steal, not select-then-insert: + // the unique index on (tenant_id, repo) makes this one atomic + // compare-and-swap on the DB side. The insert succeeds when no + // row exists yet; the `DO UPDATE ... WHERE` steals an existing + // row only when it's stale, and otherwise leaves it untouched + // and returns nothing — Postgres, not app-level timing, decides + // who wins. + const rows = await db + .insert(repoReviewLease) + .values({ + id: `lease_${crypto.randomUUID()}`, + tenantId, + repo, + leasedAt: now, + }) + .onConflictDoUpdate({ + target: [repoReviewLease.tenantId, repoReviewLease.repo], + set: { leasedAt: now }, + where: lt(repoReviewLease.leasedAt, staleBefore), + }) + .returning({ id: repoReviewLease.id }); + return rows.length > 0; + }, + + async release(tenantId, repo) { + await db + .delete(repoReviewLease) + .where( + and( + eq(repoReviewLease.tenantId, tenantId), + eq(repoReviewLease.repo, repo), + ), + ); + }, + }; +} diff --git a/packages/webhook-triggers/src/schema.ts b/packages/webhook-triggers/src/schema.ts index 82159e362..3e30c24c4 100644 --- a/packages/webhook-triggers/src/schema.ts +++ b/packages/webhook-triggers/src/schema.ts @@ -1,13 +1,26 @@ -// The one product table `@corbits/webhook-triggers` owns: a trigger -// row per external-webhook-to-workflow binding. It lives in its own -// `webhook_triggers` Postgres schema, fully siloed from the platform's -// `public` schema — see docs/package-migrations.md. +// Two product tables `@corbits/webhook-triggers` owns: a trigger row +// per external-webhook-to-workflow binding, and a short-lived lease +// row (`repoReviewLease`) closing a concurrency race in the GitHub +// connect card's start-reviewing step (CL-7242, `./repo-review-lease.ts`). +// Both live in this package's own `webhook_triggers` Postgres schema, +// fully siloed from the platform's `public` schema — see +// docs/package-migrations.md. `repoReviewLease` carries no foreign key +// to any Interchange core table: every workbench package already +// stores a platform id as a plain text column, never a cross-schema +// FK (see e.g. packages/chat/src/schema.ts's own header), and "repo" +// is a GitHub identity Interchange has no row for at all. // // The signing secret is encrypted at rest via Interchange's // `CredentialCipher` seam — see `./store.ts` for the encrypt/decrypt // wiring and `./signature.ts` for the security-model note on what that // does and does not close. -import { boolean, pgSchema, text, timestamp } from "drizzle-orm/pg-core"; +import { + boolean, + pgSchema, + text, + timestamp, + uniqueIndex, +} from "drizzle-orm/pg-core"; export const webhookTriggersSchema = pgSchema("webhook_triggers"); @@ -40,3 +53,32 @@ export const webhookTrigger = webhookTriggersSchema.table("webhook_trigger", { }); export type WebhookTriggerRow = typeof webhookTrigger.$inferSelect; + +/** + * A short-lived lease serializing `startReviewingRepos`' per-repo + * mint-grant-and-create-trigger work (CL-7242): two concurrent calls + * for the same `(tenantId, repo)` can both read "not set up yet" + * before either write lands, so the lease is acquired first and is + * the sole thing preventing both from proceeding. Never a record of + * *completion* — only ever "someone claimed responsibility for this + * repo's setup as of `leasedAt`" — so it can never assert something + * untrue the way a stale "done" marker could (CL-7213). See + * `./repo-review-lease.ts` for the acquire/release/steal-if-stale + * semantics this table backs. + */ +export const repoReviewLease = webhookTriggersSchema.table( + "repo_review_lease", + { + id: text("id").primaryKey(), + tenantId: text("tenant_id").notNull(), + repo: text("repo").notNull(), + leasedAt: timestamp("leased_at", { withTimezone: true }) + .notNull() + .defaultNow(), + }, + (t) => [ + uniqueIndex("repo_review_lease_tenant_repo_unique").on(t.tenantId, t.repo), + ], +); + +export type RepoReviewLeaseRow = typeof repoReviewLease.$inferSelect; diff --git a/packages/webhook-triggers/src/store.ts b/packages/webhook-triggers/src/store.ts index 900a517f3..50557dfd3 100644 --- a/packages/webhook-triggers/src/store.ts +++ b/packages/webhook-triggers/src/store.ts @@ -36,7 +36,27 @@ export interface CreateWebhookTriggerInput { } export interface WebhookTriggerStore { + /** + * Always inserts a new row — a second call with the same + * `(tenantId, workflowDefinitionId, name)` throws a unique-constraint + * violation rather than silently reusing the first row. The generic + * management API (`management-routes.ts`) wants this: reusing an + * existing row here would mean a caller posting a *different* + * `inputTemplate` under a name that already exists gets back 201 and + * someone else's trigger. + */ create(input: CreateWebhookTriggerInput): Promise; + /** + * Idempotent create: a second call with the same + * `(tenantId, workflowDefinitionId, name)` returns the first call's + * row untouched rather than inserting (or throwing) — the id and the + * real, already-persisted secret, never a freshly generated one that + * was never stored. This is what a retry-safe *mint* wants + * (`startReviewingRepos`'s trigger-per-repo convention, CL-7242): two + * concurrent calls for the same repo must settle on exactly one live + * trigger. + */ + ensure(input: CreateWebhookTriggerInput): Promise; get( tenantId: string, triggerId: string, @@ -100,6 +120,63 @@ export function createDrizzleWebhookTriggerStore< return { ...row, secret: input.secret }; }, + async ensure(input) { + const encryptedSecret = await credentialCipher.encrypt( + input.secret, + credentialAad(input.id, "secret"), + ); + // Insert-first, not select-then-insert: two concurrent calls for + // the same (tenant, definition, name) both attempt the insert, + // the unique index serializes them, and the loser's empty + // `returning()` re-selects the winner's row rather than creating + // a duplicate trigger. + const inserted = await db + .insert(webhookTrigger) + .values({ + id: input.id, + tenantId: input.tenantId, + name: input.name, + workflowDefinitionId: input.workflowDefinitionId, + inputTemplate: input.inputTemplate, + secret: encryptedSecret, + enabled: true, + createdBy: input.createdBy, + }) + .onConflictDoNothing({ + target: [ + webhookTrigger.tenantId, + webhookTrigger.workflowDefinitionId, + webhookTrigger.name, + ], + }) + .returning(); + const row = inserted[0]; + if (row) { + // As in `create`, the plaintext is already known — return it + // directly rather than round-tripping through an extra decrypt. + return { ...row, secret: input.secret }; + } + const [existing] = await db + .select() + .from(webhookTrigger) + .where( + and( + eq(webhookTrigger.tenantId, input.tenantId), + eq(webhookTrigger.workflowDefinitionId, input.workflowDefinitionId), + eq(webhookTrigger.name, input.name), + ), + ) + .limit(1); + if (existing === undefined) { + throw new Error( + "expected webhook trigger row after conflicting insert", + ); + } + // The winner's real, already-persisted secret — never the + // freshly generated plaintext this call never actually stored. + return decrypted(existing); + }, + async get(tenantId, triggerId) { const [row] = await db .select() diff --git a/packages/webhook-triggers/test/management-routes.test.ts b/packages/webhook-triggers/test/management-routes.test.ts index 6236f3bfd..508cdba25 100644 --- a/packages/webhook-triggers/test/management-routes.test.ts +++ b/packages/webhook-triggers/test/management-routes.test.ts @@ -57,6 +57,28 @@ describe("POST /", () => { const body = (await response.json()) as { error: { code: string } }; expect(body.error.code).toBe("bad_request"); }); + + test("a second create for the same name/definition 409s rather than duplicating", async () => { + const { app } = buildApp(); + const post = () => + app.request("/", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ + name: "Granola note-taker", + workflowDefinitionId: "def_1", + inputTemplate: "New note: {{note.title}}", + }), + }); + + const first = await post(); + expect(first.status).toBe(201); + + const second = await post(); + expect(second.status).toBe(409); + const body = (await second.json()) as { error: { code: string } }; + expect(body.error.code).toBe("conflict"); + }); }); describe("GET / and GET /:id", () => { diff --git a/packages/webhook-triggers/test/migrations.test.ts b/packages/webhook-triggers/test/migrations.test.ts index addd51948..8040daeeb 100644 --- a/packages/webhook-triggers/test/migrations.test.ts +++ b/packages/webhook-triggers/test/migrations.test.ts @@ -20,6 +20,13 @@ function scratchUrlFor(e2eUrl: string): string { const databaseUrl = e2eDatabaseUrl(); const describeIfDb = databaseUrl === undefined ? describe.skip : describe; +const migrationNames = [ + "0001_webhook_trigger", + "0002_webhook_trigger_tenant_index", + "0003_webhook_trigger_tenant_definition_name_unique", + "0004_repo_review_lease", +]; + describeIfDb("applyWebhookTriggersMigrations", () => { const scratchUrl = scratchUrlFor( databaseUrl ?? "postgres://localhost:5432/unused", @@ -56,11 +63,6 @@ describeIfDb("applyWebhookTriggersMigrations", () => { } }, 20000); - const migrationNames = [ - "0001_webhook_trigger", - "0002_webhook_trigger_tenant_index", - ]; - test("applies the trigger table into its own schema and is idempotent on a second run", async () => { const first = await applyWebhookTriggersMigrations(scratchUrl); expect(first.applied).toEqual(migrationNames); @@ -146,21 +148,16 @@ describeIfDb("applyWebhookTriggersMigrations concurrency", () => { ...first.alreadyApplied, ...second.alreadyApplied, ].sort(), - ).toEqual( - ["0001_webhook_trigger", "0002_webhook_trigger_tenant_index"] - .flatMap((name) => [name, name]) - .sort(), - ); + ).toEqual(migrationNames.flatMap((name) => [name, name]).sort()); const sql = postgres(scratchUrl, { max: 1, onnotice: () => undefined }); try { const ledgerRows = await sql.unsafe( `SELECT name FROM "webhook_triggers"."webhook_triggers_migrations" ORDER BY name`, ); - expect(ledgerRows.map((row) => String(row["name"]))).toEqual([ - "0001_webhook_trigger", - "0002_webhook_trigger_tenant_index", - ]); + expect(ledgerRows.map((row) => String(row["name"]))).toEqual( + migrationNames, + ); } finally { await sql.end(); } diff --git a/packages/webhook-triggers/test/repo-review-lease.drizzle.test.ts b/packages/webhook-triggers/test/repo-review-lease.drizzle.test.ts new file mode 100644 index 000000000..65865a181 --- /dev/null +++ b/packages/webhook-triggers/test/repo-review-lease.drizzle.test.ts @@ -0,0 +1,162 @@ +// DB-gated: skipped when no DATABASE_URL is reachable, mirroring +// `store.drizzle.test.ts`. Runs against its own scratch database. +// +// CL-7242: proves the actual compare-and-swap `acquire` relies on — +// concurrent callers racing the same (tenant, repo), a stale lease +// being stolen, and a released lease being immediately reacquirable — +// against a real Postgres. A mocked port proves nothing about +// `ON CONFLICT ... DO UPDATE ... WHERE`. +import { afterAll, beforeAll, describe, expect, test } from "bun:test"; +import { drizzle } from "drizzle-orm/postgres-js"; +import postgres from "postgres"; + +import { applyWebhookTriggersMigrations } from "../src/migrations"; +import { createDrizzleRepoReviewLeaseStore } from "../src/repo-review-lease"; + +function scratchUrlFor(e2eUrl: string): string { + const url = new URL(e2eUrl); + const database = url.pathname.replace(/^\//, ""); + url.pathname = `/${database}_repo_review_lease_drizzle_test`; + return url.toString(); +} + +const databaseUrl = process.env["DATABASE_URL"] ?? ""; +const describeIfDb = databaseUrl === "" ? describe.skip : describe; + +const TENANT_ID = "tnt_1"; + +describeIfDb("createDrizzleRepoReviewLeaseStore", () => { + const scratchUrl = scratchUrlFor( + databaseUrl || "postgres://localhost:5432/unused", + ); + const scratchTarget = new URL(scratchUrl); + const scratchDatabase = scratchTarget.pathname.replace(/^\//, ""); + + beforeAll(async () => { + const maintenanceUrl = new URL(scratchUrl); + maintenanceUrl.pathname = "/postgres"; + const maintenance = postgres(maintenanceUrl.toString(), { + max: 1, + onnotice: () => undefined, + }); + try { + await maintenance.unsafe(`DROP DATABASE IF EXISTS "${scratchDatabase}"`); + await maintenance.unsafe(`CREATE DATABASE "${scratchDatabase}"`); + } finally { + await maintenance.end(); + } + await applyWebhookTriggersMigrations(scratchUrl); + }); + + afterAll(async () => { + const maintenanceUrl = new URL(scratchUrl); + maintenanceUrl.pathname = "/postgres"; + const maintenance = postgres(maintenanceUrl.toString(), { + max: 1, + onnotice: () => undefined, + }); + try { + await maintenance.unsafe(`DROP DATABASE IF EXISTS "${scratchDatabase}"`); + } finally { + await maintenance.end(); + } + }); + + test("two concurrent acquire() calls for the same tenant/repo settle on exactly one winner", async () => { + const sql = postgres(scratchUrl, { max: 5 }); + try { + const db = drizzle(sql); + const store = createDrizzleRepoReviewLeaseStore(db); + + const results = await Promise.all([ + store.acquire(TENANT_ID, "acme/widgets"), + store.acquire(TENANT_ID, "acme/widgets"), + ]); + + expect(results.filter(Boolean)).toHaveLength(1); + + const rows = await sql` + select id from webhook_triggers.repo_review_lease + where tenant_id = ${TENANT_ID} and repo = 'acme/widgets' + `; + expect(rows).toHaveLength(1); + } finally { + await sql.end(); + } + }); + + test("a stale lease can be stolen by a new acquire()", async () => { + const sql = postgres(scratchUrl, { max: 5 }); + try { + const db = drizzle(sql); + const store = createDrizzleRepoReviewLeaseStore(db); + + expect(await store.acquire(TENANT_ID, "acme/stale-repo")).toBe(true); + + // Backdate the lease well past the staleness window, simulating + // a holder that crashed mid-work and never released it. + await sql` + update webhook_triggers.repo_review_lease + set leased_at = now() - interval '10 minutes' + where tenant_id = ${TENANT_ID} and repo = 'acme/stale-repo' + `; + + expect(await store.acquire(TENANT_ID, "acme/stale-repo")).toBe(true); + + const rows = await sql` + select id from webhook_triggers.repo_review_lease + where tenant_id = ${TENANT_ID} and repo = 'acme/stale-repo' + `; + expect(rows).toHaveLength(1); + } finally { + await sql.end(); + } + }); + + test("a fresh lease cannot be stolen", async () => { + const sql = postgres(scratchUrl, { max: 5 }); + try { + const db = drizzle(sql); + const store = createDrizzleRepoReviewLeaseStore(db); + + expect(await store.acquire(TENANT_ID, "acme/fresh-repo")).toBe(true); + expect(await store.acquire(TENANT_ID, "acme/fresh-repo")).toBe(false); + } finally { + await sql.end(); + } + }); + + test("release() lets an immediate reacquire succeed without waiting out the staleness window", async () => { + const sql = postgres(scratchUrl, { max: 5 }); + try { + const db = drizzle(sql); + const store = createDrizzleRepoReviewLeaseStore(db); + + expect(await store.acquire(TENANT_ID, "acme/released-repo")).toBe(true); + await store.release(TENANT_ID, "acme/released-repo"); + + const rows = await sql` + select id from webhook_triggers.repo_review_lease + where tenant_id = ${TENANT_ID} and repo = 'acme/released-repo' + `; + expect(rows).toHaveLength(0); + + expect(await store.acquire(TENANT_ID, "acme/released-repo")).toBe(true); + } finally { + await sql.end(); + } + }); + + test("release() on a lease this caller never held is a no-op", async () => { + const sql = postgres(scratchUrl, { max: 5 }); + try { + const db = drizzle(sql); + const store = createDrizzleRepoReviewLeaseStore(db); + await expect( + store.release(TENANT_ID, "acme/never-leased"), + ).resolves.toBeUndefined(); + } finally { + await sql.end(); + } + }); +}); diff --git a/packages/webhook-triggers/test/store.drizzle.test.ts b/packages/webhook-triggers/test/store.drizzle.test.ts index 9ef2496d9..b42a66cd2 100644 --- a/packages/webhook-triggers/test/store.drizzle.test.ts +++ b/packages/webhook-triggers/test/store.drizzle.test.ts @@ -175,5 +175,49 @@ describeIfDb( await sql.end(); } }); + + // CL-7242: `startReviewingRepos`'s check-then-act + // (hasWebhookTrigger then createWebhookTrigger) can race two + // concurrent "start reviewing" calls for the same repo past the + // read before either write lands. `ensure` is the actual backstop + // apps/hub's `createWebhookTrigger` port binds to instead of + // `create` — reconstructs the audit's finding (2 live triggers for + // one repo) against the real `webhook_trigger_tenant_definition_name_unique` + // index (migration 0003). + test("two concurrent ensure() calls for the same tenant/definition/name settle on exactly one trigger", async () => { + const sql = postgres(scratchUrl, { max: 5 }); + try { + const db = drizzle(sql); + const cipher = createEnvKeyCredentialCipher(KEY); + const store = createDrizzleWebhookTriggerStore(db, cipher); + + const input = (id: string) => ({ + id, + tenantId: TENANT_ID, + name: "acme/widgets pull-request-opened", + workflowDefinitionId: "def_code_review", + inputTemplate: "Review the pull request at {{pull_request.html_url}}", + secret: `secret-${id}`, + createdBy: "user_1", + }); + + const [first, second] = await Promise.all([ + store.ensure(input("wht_race_1")), + store.ensure(input("wht_race_2")), + ]); + + expect(first.id).toBe(second.id); + expect(first.secret).toBe(second.secret); + + const rows = await sql` + select id from webhook_triggers.webhook_trigger + where tenant_id = ${TENANT_ID} and workflow_definition_id = 'def_code_review' + and name = 'acme/widgets pull-request-opened' + `; + expect(rows).toHaveLength(1); + } finally { + await sql.end(); + } + }); }, ); diff --git a/packages/webhook-triggers/test/test-support.ts b/packages/webhook-triggers/test/test-support.ts index 643630d77..398dce0d8 100644 --- a/packages/webhook-triggers/test/test-support.ts +++ b/packages/webhook-triggers/test/test-support.ts @@ -38,20 +38,59 @@ export function principal(id: string) { export function createInMemoryWebhookTriggerStore(): WebhookTriggerStore { const rows = new Map(); + function findByName( + tenantId: string, + workflowDefinitionId: string, + name: string, + ): WebhookTriggerRow | undefined { + return [...rows.values()].find( + (row) => + row.tenantId === tenantId && + row.workflowDefinitionId === workflowDefinitionId && + row.name === name, + ); + } + + function newRow(input: CreateWebhookTriggerInput): WebhookTriggerRow { + return { + id: input.id, + tenantId: input.tenantId, + name: input.name, + workflowDefinitionId: input.workflowDefinitionId, + inputTemplate: input.inputTemplate, + secret: input.secret, + enabled: true, + createdBy: input.createdBy, + createdAt: new Date(), + lastFiredAt: null, + }; + } + return { async create(input: CreateWebhookTriggerInput) { - const row: WebhookTriggerRow = { - id: input.id, - tenantId: input.tenantId, - name: input.name, - workflowDefinitionId: input.workflowDefinitionId, - inputTemplate: input.inputTemplate, - secret: input.secret, - enabled: true, - createdBy: input.createdBy, - createdAt: new Date(), - lastFiredAt: null, - }; + if (findByName(input.tenantId, input.workflowDefinitionId, input.name)) { + // Mirrors the real driver's shape for a `23505` unique + // violation, so callers exercising `isUniqueViolation`-style + // handling against this fake see the same thing production does. + throw Object.assign( + new Error( + `webhook trigger ${input.name} already exists for this workflow definition`, + ), + { code: "23505" }, + ); + } + const row = newRow(input); + rows.set(row.id, row); + return row; + }, + async ensure(input: CreateWebhookTriggerInput) { + const existing = findByName( + input.tenantId, + input.workflowDefinitionId, + input.name, + ); + if (existing) return existing; + const row = newRow(input); rows.set(row.id, row); return row; }, diff --git a/packages/workflow-catalog/src/connect-github-credential-link.test.ts b/packages/workflow-catalog/src/connect-github-credential-link.test.ts index 5d38af494..c0b5b311b 100644 --- a/packages/workflow-catalog/src/connect-github-credential-link.test.ts +++ b/packages/workflow-catalog/src/connect-github-credential-link.test.ts @@ -139,6 +139,8 @@ describe("the room GitHub connect card reads what its own submit writes", () => githubDescriptor.displayName, ), resolveCodeReviewDefinitionId: async () => "wfd_code_review", + acquireRepoReviewLease: async () => true, + releaseRepoReviewLease: async () => {}, hasRepoGrant: async () => false, mintRepoGrant: async () => {}, createWebhookTrigger: async () => ({ id: "trg_1" }), @@ -184,6 +186,8 @@ describe("the room GitHub connect card reads what its own submit writes", () => log: () => {}, resolveGithubConfig: buildResolveGithubConfig(store, githubDescriptor.id), resolveCodeReviewDefinitionId: async () => "wfd_code_review", + acquireRepoReviewLease: async () => true, + releaseRepoReviewLease: async () => {}, hasRepoGrant: async () => false, mintRepoGrant: async () => {}, createWebhookTrigger: async () => ({ id: "trg_1" }), diff --git a/packages/workflow-catalog/src/connect-github-routes.test.ts b/packages/workflow-catalog/src/connect-github-routes.test.ts index 1677dd00a..fbe4debb0 100644 --- a/packages/workflow-catalog/src/connect-github-routes.test.ts +++ b/packages/workflow-catalog/src/connect-github-routes.test.ts @@ -80,6 +80,8 @@ function buildApp(overrides: Partial = {}) { log: () => {}, resolveGithubConfig: async () => githubConfig, resolveCodeReviewDefinitionId: async () => "wfd_code_review", + acquireRepoReviewLease: async () => true, + releaseRepoReviewLease: async () => {}, hasRepoGrant: async (tenantId, repo) => grants.some((g) => g.tenantId === tenantId && g.repo.id === repo.id), mintRepoGrant: async (tenantId, repo) => { diff --git a/packages/workflow-catalog/src/connect-github-routes.ts b/packages/workflow-catalog/src/connect-github-routes.ts index 4580d851e..1880bf95d 100644 --- a/packages/workflow-catalog/src/connect-github-routes.ts +++ b/packages/workflow-catalog/src/connect-github-routes.ts @@ -84,6 +84,23 @@ export type ConnectGithubRoutesDeps = { * `undefined` when the template's own workflow was never deployed for * this tenant (a create-flow bug, not something this route can fix). */ resolveCodeReviewDefinitionId(tenantId: string): Promise; + /** Acquires the short-lived lease serializing one repo's setup + * (CL-7242) — see `./connect-github-setup.ts`'s + * `ConnectGithubSetupPorts.acquireRepoReviewLease` for why this is + * the actual concurrency backstop, not `hasRepoGrant`/`hasWebhookTrigger` + * below. A host binds this to `@corbits/webhook-triggers`' + * `RepoReviewLeaseStore.acquire`. */ + acquireRepoReviewLease( + tenantId: string, + repo: GitHubRepoSummary, + ): Promise; + /** Releases a lease `acquireRepoReviewLease` won — see + * `ConnectGithubSetupPorts.releaseRepoReviewLease`. A host binds + * this to `RepoReviewLeaseStore.release`. */ + releaseRepoReviewLease( + tenantId: string, + repo: GitHubRepoSummary, + ): Promise; /** True once this repo already has the `repo:` grant — see * `./connect-github-setup.ts`'s `ConnectGithubSetupPorts.hasRepoGrant` * for why this makes a retry between minting the grant and creating @@ -294,6 +311,10 @@ export function createConnectGithubRoutes( const introductionsAlreadyPosted = settingsBefore.selectedRepos.length > 0; const result = await startReviewingRepos(body.repoIds, state.repos, { + acquireRepoReviewLease: (repo) => + deps.acquireRepoReviewLease(tenant.id, repo), + releaseRepoReviewLease: (repo) => + deps.releaseRepoReviewLease(tenant.id, repo), hasRepoGrant: (repo) => deps.hasRepoGrant( tenant.id, diff --git a/packages/workflow-catalog/src/connect-github-setup.test.ts b/packages/workflow-catalog/src/connect-github-setup.test.ts index 275a8d620..bc8ebae6d 100644 --- a/packages/workflow-catalog/src/connect-github-setup.test.ts +++ b/packages/workflow-catalog/src/connect-github-setup.test.ts @@ -18,6 +18,10 @@ function fakePorts() { const createdTriggerRepos: string[] = []; let persistedRepoIds: readonly string[] | undefined; const ports: ConnectGithubSetupPorts = { + async acquireRepoReviewLease() { + return true; + }, + async releaseRepoReviewLease() {}, async hasRepoGrant() { return false; }, @@ -55,9 +59,22 @@ function fakePorts() { function fakeBackedPorts() { const grantedRepoNames = new Set(); const existingTriggerRepoNames = new Set(); + const leasedRepoNames = new Set(); const grantedRepos: string[] = []; const createdTriggerRepos: string[] = []; const ports: ConnectGithubSetupPorts = { + // Synchronous check-and-set, same as the real lease's DB-side + // compare-and-swap: two calls sharing this same `ports` object + // (e.g. via `Promise.all`) race here exactly like two real + // concurrent HTTP requests race the real `RepoReviewLeaseStore`. + async acquireRepoReviewLease(repo) { + if (leasedRepoNames.has(repo.name)) return false; + leasedRepoNames.add(repo.name); + return true; + }, + async releaseRepoReviewLease(repo) { + leasedRepoNames.delete(repo.name); + }, async hasRepoGrant(repo) { return grantedRepoNames.has(repo.name); }, @@ -189,6 +206,35 @@ describe("startReviewingRepos", () => { ]); expect(result.createdTriggerIds).toEqual(["trg_2", "trg_3"]); }); + + test("two concurrent calls for the same repo mint exactly one grant and one trigger (CL-7242)", async () => { + // Reconstructs the audit's own reproduction: `Promise.all` of two + // concurrent `startReviewingRepos` calls sharing the same + // backing state, for the same repo. Before the lease existed, + // both calls' `hasRepoGrant`/`hasWebhookTrigger` reads landed + // before either write, so both minted -- 2 grants, 2 triggers, + // for one repo. `acquireRepoReviewLease`'s synchronous + // check-and-set means only one of the two calls below ever gets + // past the lease for "acme/widgets": the other's per-repo body + // never runs at all. + const fake = fakeBackedPorts(); + + const [first, second] = await Promise.all([ + startReviewingRepos(["1"], REPOS, fake.ports), + startReviewingRepos(["1"], REPOS, fake.ports), + ]); + + expect(fake.grantedRepos).toEqual(["acme/widgets"]); + expect(fake.createdTriggerRepos).toEqual(["acme/widgets"]); + + const totalCreated = + first.createdTriggerIds.length + second.createdTriggerIds.length; + expect(totalCreated).toBe(1); + + const totalSkipped = + first.skippedRepoIds.length + second.skippedRepoIds.length; + expect(totalSkipped).toBe(1); + }); }); describe("webhookTriggerName", () => { diff --git a/packages/workflow-catalog/src/connect-github-setup.ts b/packages/workflow-catalog/src/connect-github-setup.ts index 9f8d54010..c65fca755 100644 --- a/packages/workflow-catalog/src/connect-github-setup.ts +++ b/packages/workflow-catalog/src/connect-github-setup.ts @@ -8,15 +8,40 @@ import type { GitHubRepoSummary } from "@corbits/github-tools"; export interface ConnectGithubSetupPorts { + /** + * Acquires the short-lived lease serializing this repo's setup + * (CL-7242): true means this call now owns it and must run the rest + * of this repo's body below; false means another call currently + * owns it (or very recently did) and this repo must be skipped + * entirely, including `hasRepoGrant`/`hasWebhookTrigger`. This is + * the actual concurrency backstop — `hasRepoGrant` and + * `hasWebhookTrigger` below are a fast path for a *sequential* + * retry (CL-7134), not a lock, and are only safe to reach because + * the lease already ensures two concurrent calls for the same repo + * can never both get here. A host binds this to + * `@corbits/webhook-triggers`' `RepoReviewLeaseStore.acquire`. + */ + acquireRepoReviewLease(repo: GitHubRepoSummary): Promise; + /** + * Releases the lease `acquireRepoReviewLease` won, once this repo's + * body finishes (success or failure) — so a legitimate retry never + * waits out the lease's own staleness window. A host binds this to + * `RepoReviewLeaseStore.release`. + */ + releaseRepoReviewLease(repo: GitHubRepoSummary): Promise; /** * True once this repo already has the `repo:` grant — * checked before minting one, so a retry after a failure between - * minting the grant and creating the trigger never mints a second - * grant for a repo that already has one. The `grant` table - * (`vendor/intx/db`) carries no unique constraint over + * minting the grant and creating the trigger skips straight past + * minting for a repo that already has one (CL-7134's fast path). The + * `grant` table (`vendor/intx/db`) carries no unique constraint over * tenant/resource/action, so this read is the only thing standing - * between a retry and a duplicate row. A host binds this to GET - * `/api/tenants/:id/grants?resource=repo:`. + * between a retry and a duplicate row absent the lease below. A host + * binds this to GET `/api/tenants/:id/grants?resource=repo:`. + * + * Safe to be a plain read, not a conflict-checked one: the lease + * above already guarantees only one caller ever reaches this point + * for a given repo at a time. */ hasRepoGrant(repo: GitHubRepoSummary): Promise; /** @@ -26,13 +51,23 @@ export interface ConnectGithubSetupPorts { * middleware, applied here to a repo instead of a room). A host binds * this to POST `/api/tenants/:id/grants`; this module never touches * drizzle directly. + * + * A plain insert is safe here (CL-7242): the lease above is what + * makes this call-site single-flight per repo, so this never needs + * its own conflict handling against the platform's `grant` table. */ mintRepoGrant(repo: GitHubRepoSummary): Promise; /** * Creates the live `webhook_trigger` row this repo's pull-request-opened * events fire — the onboarding card's start-reviewing step is what - * creates this trigger, for each repo the person picked. A host binds - * this to `@corbits/webhook-triggers`' `WebhookTriggerStore.create`. + * creates this trigger, for each repo the person picked. A host + * binds this to `@corbits/webhook-triggers`' `WebhookTriggerStore.ensure` + * rather than `create`: the lease already makes this call-site + * single-flight per repo, so `ensure`'s own idempotence + * (backed by `webhook_trigger_tenant_definition_name_unique`, our + * own schema, unaffected by CL-7242's vendored-table constraint) is + * pure defense-in-depth — a lease bug degrades to a silent no-op + * here instead of a hard failure or a real duplicate trigger. */ createWebhookTrigger( repo: GitHubRepoSummary, @@ -40,9 +75,10 @@ export interface ConnectGithubSetupPorts { /** * True once this repo already has a live webhook trigger — checked * before creating one, so a retry after a mid-loop failure (a repo - * 1..N-1 already set up, N onward not) never mints a second trigger - * for a repo a prior attempt already finished. A host binds this to a - * read against `@corbits/webhook-triggers`' `WebhookTriggerStore.list`. + * 1..N-1 already set up, N onward not) skips straight past creating + * one for a repo a prior attempt already finished (CL-7134's fast + * path). A host binds this to a read against + * `@corbits/webhook-triggers`' `WebhookTriggerStore.list`. * * This is never cleared on GitHub disconnect: nothing here disables or * deletes a trigger, so re-adding a repo after a reconnect finds its @@ -63,6 +99,14 @@ export interface ConnectGithubSetupPorts { export interface StartReviewingReposResult { readonly createdTriggerIds: readonly string[]; + /** + * Repos this call didn't touch because another call currently owns + * (or very recently owned) their setup lease — not an error, but + * worth surfacing rather than silently dropping: if the lease + * holder crashed, this is the caller's only signal that a repo in + * the selection got no setup work done at all this round. + */ + readonly skippedRepoIds: readonly string[]; } /** @@ -78,7 +122,19 @@ export interface StartReviewingReposResult { * creating the trigger must still create the trigger without re-minting * the grant, and a retry after a failure before the grant was minted * must still mint it. A repo both checks already report true for is - * skipped entirely. + * skipped entirely. (CL-7134.) + * + * `hasRepoGrant`/`hasWebhookTrigger` are a fast path for a *sequential* + * retry, not a lock: two truly concurrent calls for the same repo (a + * double-click, or a client retrying an in-flight request rather than + * a failed one) can both read "not set up yet" before either write + * lands. `acquireRepoReviewLease` is the actual concurrency backstop + * (CL-7242): only the caller that wins the lease enters the rest of a + * repo's body at all, so the checks below never race against a + * concurrent duplicate for the same repo — the lease is released + * (`releaseRepoReviewLease`) as soon as that body finishes, success or + * failure, so a legitimate retry is never blocked by its own prior + * attempt. */ export async function startReviewingRepos( repoIds: readonly string[], @@ -97,20 +153,29 @@ export async function startReviewingRepos( }); const createdTriggerIds: string[] = []; + const skippedRepoIds: string[] = []; for (const repo of selected) { - if (!(await ports.hasRepoGrant(repo))) { - await ports.mintRepoGrant(repo); - } - if (await ports.hasWebhookTrigger(repo)) { + if (!(await ports.acquireRepoReviewLease(repo))) { + skippedRepoIds.push(repo.id); continue; } - const trigger = await ports.createWebhookTrigger(repo); - createdTriggerIds.push(trigger.id); + try { + if (!(await ports.hasRepoGrant(repo))) { + await ports.mintRepoGrant(repo); + } + if (await ports.hasWebhookTrigger(repo)) { + continue; + } + const trigger = await ports.createWebhookTrigger(repo); + createdTriggerIds.push(trigger.id); + } finally { + await ports.releaseRepoReviewLease(repo); + } } await ports.persistSelectedRepos(repoIds); - return { createdTriggerIds }; + return { createdTriggerIds, skippedRepoIds }; } /** diff --git a/scripts/checks/no-product-tenancy.ts b/scripts/checks/no-product-tenancy.ts index 3d6392f05..6cd6f31d5 100644 --- a/scripts/checks/no-product-tenancy.ts +++ b/scripts/checks/no-product-tenancy.ts @@ -79,8 +79,17 @@ const ALLOWLIST: readonly { }, { relPath: "packages/webhook-triggers/src/schema.ts", - maxOccurrences: 1, - tables: ["webhook_trigger"], + maxOccurrences: 2, + tables: [ + "webhook_trigger", + // CL-7242: a short-lived lease serializing the GitHub connect + // card's start-reviewing step, entirely workbench-owned state + // with no relationship to Interchange's own `grant` table — see + // repo-review-lease.ts's doc comment for why the concurrency + // fix lives here rather than as any change to Interchange's + // schema. + "repo_review_lease", + ], }, { relPath: "packages/notify/src/schema.ts", From 183db9a95716ed254318aefa6fd2df264f5252a8 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 12:02:30 -0700 Subject: [PATCH 2/3] Reconcile pre-existing duplicate repo grants without touching Interchange's schema CL-7242's lease prevents new repo: grant duplicates, but a database that already carries duplicates from before this fix ships needs them cleaned up too -- and never so a repo ends up with fewer grants than it had, only ever down to exactly one. reconcileDuplicateRepoGrants is ordinary application code, not a migration: a plain DELETE against Interchange's own `grant` table issued through @intx/db's published `grant` export -- the exact same table object apps/hub/src/index.ts already reads and writes through for every other grant operation in this codebase. No schema is touched (no ALTER TABLE, no index, no constraint), so this creates no re-pin debt the way a DDL delta would (a schema change must be hand-reapplied at every future vendor/intx/* re-pin; a DML statement runs once and leaves nothing behind to reconcile). Scoped to exactly the shape mintRepoGrant writes (system-origin, repo:-prefixed resource), keeping the lowest-id row per (tenant, resource, action) group and deleting the rest -- safe to re-run, since a database with no duplicates left just deletes nothing. scripts/db-setup.ts calls it once per setup run, alongside the installed-package migrations it already sequences, and logs what it removed. --- .../src/reconcile-duplicate-repo-grants.ts | 55 +++++++ .../reconcile-duplicate-repo-grants.test.ts | 135 ++++++++++++++++++ scripts/db-setup.ts | 17 +++ 3 files changed, 207 insertions(+) create mode 100644 packages/workflow-catalog/src/reconcile-duplicate-repo-grants.ts create mode 100644 packages/workflow-catalog/test/reconcile-duplicate-repo-grants.test.ts diff --git a/packages/workflow-catalog/src/reconcile-duplicate-repo-grants.ts b/packages/workflow-catalog/src/reconcile-duplicate-repo-grants.ts new file mode 100644 index 000000000..2b6f66f7d --- /dev/null +++ b/packages/workflow-catalog/src/reconcile-duplicate-repo-grants.ts @@ -0,0 +1,55 @@ +// One-time cleanup for pre-existing duplicate `repo:` +// grants a database can already carry from CL-7242's check-then-act +// race, before this fix's lease (`@corbits/webhook-triggers`' +// `repo_review_lease`) started preventing new ones. +// +// This is ordinary application code, not a migration: a plain DELETE +// against Interchange's own `grant` table, issued through `@intx/db`'s +// published `grant` export — the exact same table object +// `apps/hub/src/index.ts` already reads and writes through for every +// other grant operation in this codebase. No schema is touched (no +// `ALTER TABLE`, no index, no constraint), so this creates no re-pin +// debt: a DDL delta would need hand-reapplying at every future +// `vendor/intx/*` re-pin, but a DML statement runs once and leaves +// nothing behind for a re-pin to reconcile. Deliberately safe to +// re-run — a database with no duplicates left just deletes nothing. +import postgres from "postgres"; + +export interface ReconcileDuplicateRepoGrantsReport { + /** Ids of the duplicate rows removed, kept for the caller to log. */ + readonly removedIds: readonly string[]; +} + +/** + * Keeps the lowest-`id` row per `(tenant_id, resource, action)` among + * `repo:`-prefixed grants with an origin `mintRepoGrant` could plausibly + * have written and deletes the rest. `'creator'` is what the current + * `mintRepoGrantViaHttp` path (`apps/hub/src/native-repo-grants.ts`) + * writes; `'system'` is what the direct-insert path this PR replaced + * used to write, so a database carrying duplicates from either + * generation of this feature gets them cleaned up. Never touches any + * other grant family. + */ +export async function reconcileDuplicateRepoGrants( + databaseUrl: string, +): Promise { + const sql = postgres(databaseUrl, { max: 1, onnotice: () => undefined }); + try { + const removed = await sql<{ id: string }[]>` + DELETE FROM "grant" AS g + USING "grant" AS older + WHERE g.origin IN ('system', 'creator') + AND g.resource LIKE 'repo:%' + AND older.origin IN ('system', 'creator') + AND older.resource LIKE 'repo:%' + AND g.tenant_id = older.tenant_id + AND g.resource = older.resource + AND g.action = older.action + AND older.id < g.id + RETURNING g.id + `; + return { removedIds: removed.map((row) => row.id) }; + } finally { + await sql.end(); + } +} diff --git a/packages/workflow-catalog/test/reconcile-duplicate-repo-grants.test.ts b/packages/workflow-catalog/test/reconcile-duplicate-repo-grants.test.ts new file mode 100644 index 000000000..46b688bf7 --- /dev/null +++ b/packages/workflow-catalog/test/reconcile-duplicate-repo-grants.test.ts @@ -0,0 +1,135 @@ +// DB-gated: skipped when no DATABASE_URL is reachable. Runs against +// its own scratch database. +// +// CL-7242: `reconcileDuplicateRepoGrants` is the one-time cleanup for +// duplicate `repo:` grants a database can already carry +// from before the `repo_review_lease` fix started preventing new +// ones. Proves it keeps exactly one grant per group (never zero, +// matching the "a repo never ends up with fewer grants than it had" +// bar), never touches an unrelated grant, and is safe to re-run. +import { afterAll, beforeAll, describe, expect, test } from "bun:test"; +import postgres from "postgres"; + +import { runMigrations } from "@intx/db"; + +import { e2eDatabaseUrl } from "../../../scripts/e2e/harness"; +import { reconcileDuplicateRepoGrants } from "../src/reconcile-duplicate-repo-grants"; + +function scratchUrlFor(e2eUrl: string): string { + const url = new URL(e2eUrl); + const database = url.pathname.replace(/^\//, ""); + url.pathname = `/${database}_reconcile_duplicate_repo_grants_test`; + return url.toString(); +} + +const databaseUrl = e2eDatabaseUrl(); +const describeIfDb = databaseUrl === undefined ? describe.skip : describe; + +const TENANT_ID = "tnt_reconcile"; +const ROLE_ID = "role_reconcile_member"; + +describeIfDb("reconcileDuplicateRepoGrants", () => { + const scratchUrl = scratchUrlFor( + databaseUrl ?? "postgres://localhost:5432/unused", + ); + const scratchTarget = new URL(scratchUrl); + const scratchDatabase = scratchTarget.pathname.replace(/^\//, ""); + + beforeAll(async () => { + const maintenanceUrl = new URL(scratchUrl); + maintenanceUrl.pathname = "/postgres"; + const maintenance = postgres(maintenanceUrl.toString(), { + max: 1, + onnotice: () => undefined, + }); + try { + await maintenance.unsafe(`DROP DATABASE IF EXISTS "${scratchDatabase}"`); + await maintenance.unsafe(`CREATE DATABASE "${scratchDatabase}"`); + } finally { + await maintenance.end(); + } + const parsed = new URL(scratchUrl); + await runMigrations( + { + host: parsed.hostname, + port: Number(parsed.port || 5432), + user: decodeURIComponent(parsed.username), + password: decodeURIComponent(parsed.password), + database: scratchDatabase, + }, + { schema: "public" }, + ); + const sql = postgres(scratchUrl, { max: 1, onnotice: () => undefined }); + try { + await sql`INSERT INTO "tenant" (id, name, slug, domain) VALUES (${TENANT_ID}, 'Acme', 'acme-reconcile', 'acme-reconcile.example')`; + await sql`INSERT INTO "role" (id, tenant_id, name) VALUES (${ROLE_ID}, ${TENANT_ID}, 'member')`; + } finally { + await sql.end(); + } + }); + + afterAll(async () => { + const maintenanceUrl = new URL(scratchUrl); + maintenanceUrl.pathname = "/postgres"; + const maintenance = postgres(maintenanceUrl.toString(), { + max: 1, + onnotice: () => undefined, + }); + try { + await maintenance.unsafe(`DROP DATABASE IF EXISTS "${scratchDatabase}"`); + } finally { + await maintenance.end(); + } + }); + + test("keeps exactly one grant per duplicated (tenant, resource, action) and leaves other grants alone", async () => { + const sql = postgres(scratchUrl, { max: 1, onnotice: () => undefined }); + try { + // Two duplicate repo grants from the pre-lease race, one pair from + // each generation of `mintRepoGrant`: the direct-insert path this + // PR replaced wrote `'system'`-origin rows; the current + // `mintRepoGrantViaHttp` path writes `'creator'`-origin rows. + await sql` + INSERT INTO "grant" (id, tenant_id, role_id, resource, action, effect, origin, created_at) + VALUES + ('grant_dup_1', ${TENANT_ID}, ${ROLE_ID}, 'repo:acme/widgets', 'read', 'allow', 'system', now() - interval '1 hour'), + ('grant_dup_2', ${TENANT_ID}, ${ROLE_ID}, 'repo:acme/widgets', 'read', 'allow', 'creator', now()) + `; + // A non-duplicated repo grant, a non-repo system grant, and a + // repo grant with an origin this feature never writes -- none + // should ever be touched. + await sql` + INSERT INTO "grant" (id, tenant_id, role_id, resource, action, effect, origin) + VALUES + ('grant_solo', ${TENANT_ID}, ${ROLE_ID}, 'repo:acme/gadgets', 'read', 'allow', 'creator'), + ('grant_unrelated', ${TENANT_ID}, ${ROLE_ID}, 'workbench:room-1', 'read', 'allow', 'system'), + ('grant_other_origin', ${TENANT_ID}, ${ROLE_ID}, 'repo:acme/gizmos', 'read', 'allow', 'role') + `; + + const first = await reconcileDuplicateRepoGrants(scratchUrl); + expect(first.removedIds).toEqual(["grant_dup_2"]); + + const widgetsRows = await sql` + SELECT id FROM "grant" WHERE tenant_id = ${TENANT_ID} AND resource = 'repo:acme/widgets' + `; + expect(widgetsRows).toHaveLength(1); + expect(widgetsRows[0]?.["id"]).toBe("grant_dup_1"); + + const untouchedRows = await sql` + SELECT id FROM "grant" WHERE id IN ('grant_solo', 'grant_unrelated', 'grant_other_origin') + ORDER BY id + `; + expect(untouchedRows.map((r) => r["id"])).toEqual([ + "grant_other_origin", + "grant_solo", + "grant_unrelated", + ]); + + // Idempotent: nothing left to remove on a second run. + const second = await reconcileDuplicateRepoGrants(scratchUrl); + expect(second.removedIds).toEqual([]); + } finally { + await sql.end(); + } + }); +}); diff --git a/scripts/db-setup.ts b/scripts/db-setup.ts index f8189d2d2..d60382ac6 100644 --- a/scripts/db-setup.ts +++ b/scripts/db-setup.ts @@ -37,6 +37,7 @@ import { listWorkbenchLaunchFoldedRunIds, } from "../packages/chat/src/migrations"; import { applyWebhookTriggersMigrations } from "../packages/webhook-triggers/src/migrations"; +import { reconcileDuplicateRepoGrants } from "../packages/workflow-catalog/src/reconcile-duplicate-repo-grants"; import { applyNotifyMigrations } from "../packages/notify/src/migrations"; import { applyRoutineMigrations } from "../packages/routines/src/migrations"; import { @@ -126,6 +127,22 @@ async function applyInstalledPackageMigrations( } } await backfillFoldedRunsFromInstalledPackages(databaseUrl); + + // CL-7242: not a package migration (no schema, no ledger, no DDL at + // all) -- a plain, always-safe-to-re-run DELETE against the + // platform's own `grant` table through @intx/db's published export, + // cleaning up any repo grants CL-7242's race duplicated before this + // fix's lease started preventing new ones. See + // reconcile-duplicate-repo-grants.ts for why this is DML, not DDL, + // and why that distinction is what keeps it off the vendored-delta + // ledger. + const { removedIds } = await reconcileDuplicateRepoGrants(databaseUrl); + if (removedIds.length > 0) { + console.log( + `db-setup: removed ${removedIds.length} duplicate repo grant(s): ` + + removedIds.join(", "), + ); + } } /** From 92b3012cb6dff005dd686d1bbe8431277fb89df7 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 12:02:37 -0700 Subject: [PATCH 3/3] Add startReviewingRepos concurrency test wiring real ports The lease store's own test proves the underlying compare-and-swap works in isolation; this proves startReviewingRepos itself settles on exactly one grant and one trigger when its ports are bound the way apps/hub/src/index.ts actually binds them -- a bare grant insert, a real RepoReviewLeaseStore, a real WebhookTriggerStore.ensure() -- racing two calls for the same repo the way a double-click or a client retrying an in-flight request would, against a scratch database. --- bun.lock | 4 + packages/workflow-catalog/package.json | 6 +- .../test/connect-github-setup.race.test.ts | 218 ++++++++++++++++++ 3 files changed, 227 insertions(+), 1 deletion(-) create mode 100644 packages/workflow-catalog/test/connect-github-setup.race.test.ts diff --git a/bun.lock b/bun.lock index bbddfa72e..f30cae5d4 100644 --- a/bun.lock +++ b/bun.lock @@ -1513,11 +1513,15 @@ "@workbench/hub-client": "workspace:*", "arktype": "catalog:", "hono": "^4.11.9", + "postgres": "catalog:", }, "devDependencies": { "@corbits/workflow-freeze": "workspace:*", + "@intx/crypto": "0.3.0", + "@intx/db": "workspace:*", "@types/bun": "catalog:", "@workbench/connections": "workspace:*", + "drizzle-orm": "catalog:", "typescript": "catalog:", }, }, diff --git a/packages/workflow-catalog/package.json b/packages/workflow-catalog/package.json index 46331677f..b307303ce 100644 --- a/packages/workflow-catalog/package.json +++ b/packages/workflow-catalog/package.json @@ -25,12 +25,16 @@ "@intx/hub-api": "workspace:*", "@workbench/hub-client": "workspace:*", "arktype": "catalog:", - "hono": "^4.11.9" + "hono": "^4.11.9", + "postgres": "catalog:" }, "devDependencies": { "@corbits/workflow-freeze": "workspace:*", + "@intx/crypto": "0.3.0", + "@intx/db": "workspace:*", "@types/bun": "catalog:", "@workbench/connections": "workspace:*", + "drizzle-orm": "catalog:", "typescript": "catalog:" } } diff --git a/packages/workflow-catalog/test/connect-github-setup.race.test.ts b/packages/workflow-catalog/test/connect-github-setup.race.test.ts new file mode 100644 index 000000000..1ca56a228 --- /dev/null +++ b/packages/workflow-catalog/test/connect-github-setup.race.test.ts @@ -0,0 +1,218 @@ +// CL-7242: reconstructs the audit's own reproduction -- two concurrent +// `startReviewingRepos` calls for the same repo -- against real, +// database-backed ports rather than plain fakes. `hasRepoGrant`/ +// `mintRepoGrant` here fake a plain read-then-insert against the +// `grant` table directly, standing in for `apps/hub/src/index.ts`'s +// real binding (HTTP calls through `native-repo-grants.ts`, see +// CL-7242 follow-up "Mint workbench tenants and repo grants via +// Interchange HTTP") without a live hub-api server -- what this test +// actually proves is that `startReviewingRepos`' lease serializes any +// such hasRepoGrant/mintRepoGrant pair correctly, which is exactly +// what makes the real HTTP-bound versions safe too. The lease store's +// own test (`@corbits/webhook-triggers`'s +// `repo-review-lease.drizzle.test.ts`) proves the underlying +// compare-and-swap works; this proves `startReviewingRepos` itself +// settles on exactly one grant and one trigger, racing the two calls +// a double-click or a client retry would produce. DB-gated: skipped +// when DATABASE_URL is unset. Runs against its own scratch database +// (mirroring `@corbits/webhook-triggers`' own `store.drizzle.test.ts`), +// never the developer's or the walking-skeleton suite's. +import { afterAll, beforeAll, describe, expect, test } from "bun:test"; +import { drizzle } from "drizzle-orm/postgres-js"; +import { and, eq } from "drizzle-orm"; +import postgres from "postgres"; + +import { runMigrations } from "@intx/db"; +import * as intxSchema from "@intx/db/schema"; +import { grant as grantTable, role as roleTable } from "@intx/db/schema"; +import { createNoopCredentialCipher } from "@intx/crypto"; +import { + createDrizzleRepoReviewLeaseStore, + createDrizzleWebhookTriggerStore, + applyWebhookTriggersMigrations, +} from "@corbits/webhook-triggers"; +import type { GitHubRepoSummary } from "@corbits/github-tools"; + +import { e2eDatabaseUrl } from "../../../scripts/e2e/harness"; +import { + startReviewingRepos, + webhookTriggerName, + type ConnectGithubSetupPorts, +} from "../src/connect-github-setup"; + +function scratchUrlFor(e2eUrl: string): string { + const url = new URL(e2eUrl); + const database = url.pathname.replace(/^\//, ""); + url.pathname = `/${database}_connect_github_race_test`; + return url.toString(); +} + +const databaseUrl = e2eDatabaseUrl(); +const describeIfDb = databaseUrl === undefined ? describe.skip : describe; + +const REPO: GitHubRepoSummary = { + id: "1", + name: "acme/widgets", +}; +const TENANT_ID = "tnt_race"; +const DEFINITION_ID = "def_code_review"; + +describeIfDb("startReviewingRepos under real concurrency (CL-7242)", () => { + const scratchUrl = scratchUrlFor( + databaseUrl ?? "postgres://localhost:5432/unused", + ); + const scratchTarget = new URL(scratchUrl); + const scratchDatabase = scratchTarget.pathname.replace(/^\//, ""); + + beforeAll(async () => { + const maintenanceUrl = new URL(scratchUrl); + maintenanceUrl.pathname = "/postgres"; + const maintenance = postgres(maintenanceUrl.toString(), { + max: 1, + onnotice: () => undefined, + }); + try { + await maintenance.unsafe(`DROP DATABASE IF EXISTS "${scratchDatabase}"`); + await maintenance.unsafe(`CREATE DATABASE "${scratchDatabase}"`); + } finally { + await maintenance.end(); + } + const parsed = new URL(scratchUrl); + await runMigrations( + { + host: parsed.hostname, + port: Number(parsed.port || 5432), + user: decodeURIComponent(parsed.username), + password: decodeURIComponent(parsed.password), + database: scratchDatabase, + }, + { schema: "public" }, + ); + await applyWebhookTriggersMigrations(scratchUrl); + }); + + afterAll(async () => { + const maintenanceUrl = new URL(scratchUrl); + maintenanceUrl.pathname = "/postgres"; + const maintenance = postgres(maintenanceUrl.toString(), { + max: 1, + onnotice: () => undefined, + }); + try { + await maintenance.unsafe(`DROP DATABASE IF EXISTS "${scratchDatabase}"`); + } finally { + await maintenance.end(); + } + }); + + test("two concurrent calls for the same repo mint exactly one grant and one trigger", async () => { + const client = postgres(scratchUrl, { + max: 10, + onnotice: () => undefined, + }); + try { + const db = drizzle(client, { schema: intxSchema }); + const webhookStore = createDrizzleWebhookTriggerStore( + db, + createNoopCredentialCipher(), + ); + const leaseStore = createDrizzleRepoReviewLeaseStore(db); + + const roleId = "role_race_member"; + await client`INSERT INTO "tenant" (id, name, slug, domain) VALUES (${TENANT_ID}, 'Acme', 'acme-race', 'acme-race.example')`; + await client`INSERT INTO "role" (id, tenant_id, name) VALUES (${roleId}, ${TENANT_ID}, 'member')`; + + // Mirrors apps/hub/src/index.ts's real wiring exactly: the lease + // is acquired first, and mintRepoGrant is a bare insert against + // Interchange's own `grant` table -- no onConflict, no index of + // ours on their table. Only the lease serializes the two racing + // calls below. + const buildPorts = (): ConnectGithubSetupPorts => ({ + acquireRepoReviewLease: (repo) => + leaseStore.acquire(TENANT_ID, repo.name), + releaseRepoReviewLease: (repo) => + leaseStore.release(TENANT_ID, repo.name), + hasRepoGrant: async (repo) => { + const existing = await db.query.grant.findFirst({ + where: and( + eq(grantTable.tenantId, TENANT_ID), + eq(grantTable.resource, `repo:${repo.name}`), + eq(grantTable.action, "read"), + ), + columns: { id: true }, + }); + return existing !== undefined; + }, + mintRepoGrant: async (repo) => { + const memberRole = await db.query.role.findFirst({ + where: and( + eq(roleTable.tenantId, TENANT_ID), + eq(roleTable.name, "member"), + ), + columns: { id: true }, + }); + if (memberRole === undefined) throw new Error("no member role"); + await db.insert(grantTable).values({ + id: `grant_${crypto.randomUUID()}`, + tenantId: TENANT_ID, + roleId: memberRole.id, + resource: `repo:${repo.name}`, + action: "read", + effect: "allow", + origin: "system", + }); + }, + hasWebhookTrigger: async (repo) => { + const triggers = await webhookStore.list(TENANT_ID); + const name = webhookTriggerName(repo); + return triggers.some( + (t) => t.workflowDefinitionId === DEFINITION_ID && t.name === name, + ); + }, + createWebhookTrigger: async (repo) => { + const row = await webhookStore.ensure({ + id: `wht_${crypto.randomUUID()}`, + tenantId: TENANT_ID, + name: webhookTriggerName(repo), + workflowDefinitionId: DEFINITION_ID, + inputTemplate: + "Review the pull request at {{pull_request.html_url}}", + secret: crypto.randomUUID(), + createdBy: "user_1", + }); + return { id: row.id }; + }, + persistSelectedRepos: async () => {}, + }); + + const [first, second] = await Promise.all([ + startReviewingRepos(["1"], [REPO], buildPorts()), + startReviewingRepos(["1"], [REPO], buildPorts()), + ]); + + // Whether the second call gets skipped by the lease or simply + // finds the work already done (CL-7134's fast path) depends on + // real, non-deterministic timing between the two real Postgres + // round trips -- both are correct outcomes. The property this + // test actually cares about is that at most one trigger is ever + // created, and the database never ends up with a duplicate. + const totalCreated = + first.createdTriggerIds.length + second.createdTriggerIds.length; + expect(totalCreated).toBe(1); + + const grantRows = await client` + SELECT id FROM "grant" + WHERE tenant_id = ${TENANT_ID} AND resource = 'repo:acme/widgets' + `; + expect(grantRows).toHaveLength(1); + + const triggerRows = await client` + SELECT id FROM "webhook_triggers"."webhook_trigger" + WHERE tenant_id = ${TENANT_ID} AND workflow_definition_id = ${DEFINITION_ID} + `; + expect(triggerRows).toHaveLength(1); + } finally { + await client.end(); + } + }); +});