diff --git a/CHANGELOG.md b/CHANGELOG.md index 7b01285..52eea8a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -16,6 +16,8 @@ This file records what changes **in the product** – process and session state such leftovers. An approved sample whose export failed is kept (#84). ### Fixed +- Two `pnpm seed:samples` runs for the same company at the same time no longer interfere: the second waits for + the first and then finds the samples in place, or stops after 20 seconds with a clear message (#93). - `/api/health` reports the AI service as `starting` instead of `failed` when it does not answer in time – on the showcase usually the cold start of the scaled-to-zero container; after 30 s without an answer it reads `failed` again, so a real outage is not hidden (#81). diff --git a/docs/technical/architecture.md b/docs/technical/architecture.md index c6d2387..3bd1002 100644 --- a/docs/technical/architecture.md +++ b/docs/technical/architecture.md @@ -30,7 +30,7 @@ Every new file belongs to one of these modules – otherwise add the module here | `review` | `src/features/review/` | review UI, corrections, approve/reject | authenticated UI | confidential + personal | session, role check, audit | built: review page (fields + status badges + source view for mail, PDF incl. OCR label, XLSX cells, DOCX paragraphs/tables, attachments; line items as a table with per-field status and audited corrections), corrections with history, approve (→ export job) / reject with reason, duplicate decision (confirm or reject as duplicate) | | `export` | `src/features/export/` | ERP port + REST adapter, idempotency | outbound HTTP | confidential | idempotency key, unique export, timeout | built: REST adapter (timeout, error classes, contract validation), export handler under row lock, `drainExports()`, `request_exports` | | `erp-mock` | `src/features/erp-mock/` | simulated ERP REST API | route behind flag | synthetic | disabled unless `ERP_MOCK_ENABLED` | built: idempotent receiver (replay → same reference, 409 on a different body), fault injection, bounded in-memory store, route `/api/erp-mock/v1/quote-requests` | -| `samples` | `src/features/samples/`, entrypoints `src/seed-samples.ts` (`pnpm seed:samples`) and `src/samples-record.ts` (`pnpm samples:record`) | prepared showcase cases: synthetic mails with a recorded AI answer, seeded through intake → processing → approval | operator scripts only | synthetic | no model call when seeding (recording replayed inline, no queued processing job); requests marked `source = 'sample'` (CHECK, migration 0019) | built: two samples per demo company – one in review (uncertain + missing value), one approved with its export queued (#71) | +| `samples` | `src/features/samples/`, entrypoints `src/seed-samples.ts` (`pnpm seed:samples`) and `src/samples-record.ts` (`pnpm samples:record`) | prepared showcase cases: synthetic mails with a recorded AI answer, seeded through intake → processing → approval | operator scripts only | synthetic | no model call when seeding (recording replayed inline, no queued processing job); requests marked `source = 'sample'` (CHECK, migration 0019); one seed run per company at a time (transaction-level advisory lock, #93) | built: two samples per demo company – one in review (uncertain + missing value), one approved with its export queued (#71) | | `identity` | `src/features/identity/` | Better Auth, users, companies, roles | public login route | personal (staff) | rate limit, invite-only | built: Better Auth (invite-only, organization + admin plugins), `authorize()`, audited invite, user management (`/users`: roles, deactivate/reactivate, last-admin rule), seed | | `tenancy` | `src/features/tenancy/` | `withTenant()`, RLS policies | internal | – | forced RLS, `app_rw` without BYPASSRLS | built: `withTenant()`, forced RLS on `app.*`, guard test (every `app` table: `company_id`, forced RLS, only company policies; allow-list empty) | | `audit` | `src/features/audit/` | append-only audit events | internal | personal (staff) | INSERT/SELECT only | partial: `recordAudit()` (append-only enforced by grants) | diff --git a/docs/technical/deployment-vercel.md b/docs/technical/deployment-vercel.md index e7c0d61..764d012 100644 --- a/docs/technical/deployment-vercel.md +++ b/docs/technical/deployment-vercel.md @@ -128,7 +128,9 @@ model provider is overloaded. List and detail label them „Vorbereitetes Beispi step is idempotent: repeat it whenever visitors have decided the samples, and it adds only what is missing. Samples an aborted run left behind (new, in progress, or in error from processing) are first set to rejected with a fixed reason and an audit event – never deleted (#84). An approved sample whose export failed is kept: -it can be exported again with „Erneut verarbeiten“. +it can be exported again with „Erneut verarbeiten“. Only one run per company works at a time (#93): a second +run started meanwhile waits for the first and then finds the samples in place; after 20 seconds it stops with +„Another sample seed run for this company is still in progress“ – run it again later. ## 7. Jobs without a worker diff --git a/src/features/samples/index.ts b/src/features/samples/index.ts index cec5915..99235f9 100644 --- a/src/features/samples/index.ts +++ b/src/features/samples/index.ts @@ -1,3 +1,4 @@ // Public API of the `samples` module (#71): prepared showcase cases from recorded AI answers. export { freshSampleMail, RECORDED_MODEL_PREFIX, recordedAiClient, recordingFile, SAMPLES, sampleMail, sampleRecording, type Sample } from "./samples"; +export { SeedRunBusy } from "./repository"; export { seedSamples, type SampleDeps, type SeededSample, type SeedResult } from "./seed"; diff --git a/src/features/samples/repository.ts b/src/features/samples/repository.ts new file mode 100644 index 0000000..b738e27 --- /dev/null +++ b/src/features/samples/repository.ts @@ -0,0 +1,30 @@ +import { sql } from "drizzle-orm"; +import { tenantOf, type TenantTx } from "@/features/tenancy"; + +/** Another seed run of the same company is still in progress (#93). */ +export class SeedRunBusy extends Error { + constructor() { + super("Another sample seed run for this company is still in progress – run it again once that one has finished."); + this.name = "SeedRunBusy"; + } +} + +// Waiting ends either way: at our `lock_timeout`, or at a lower `statement_timeout` of pool or role. +const GAVE_UP = new Set(["55P03", "57014"]); + +/** + * One seed run per company at a time (#93). A transaction-level advisory lock – like the duplicate + * detection's – held by a transaction that stays open for the whole run: safe behind a transaction + * pooler (Supabase), and released when the run ends or its connection drops. A second run waits up to + * `timeoutMs`, then stops with `SeedRunBusy`. + */ +export async function lockSeedRun(tx: TenantTx, timeoutMs: number): Promise { + await tx.execute(sql`select set_config('lock_timeout', ${String(Math.max(1, Math.floor(timeoutMs)))}, true)`); + try { + await tx.execute(sql`select pg_advisory_xact_lock(hashtextextended(${`seed-samples:${tenantOf(tx)}`}, 0))`); + } catch (error) { + const code = (error as { code?: unknown; cause?: { code?: unknown } }).cause?.code ?? (error as { code?: unknown }).code; + if (typeof code === "string" && GAVE_UP.has(code)) throw new SeedRunBusy(); + throw error; + } +} diff --git a/src/features/samples/seed.ts b/src/features/samples/seed.ts index 27f8900..d8e70c7 100644 --- a/src/features/samples/seed.ts +++ b/src/features/samples/seed.ts @@ -6,6 +6,7 @@ import { submitUpload, type IntakeDeps } from "@/features/intake"; import { markProcessingFailed, processRequestJob, type JobSender } from "@/features/jobs"; import { countSamples, listSampleLeftoverIds, lockRequest, retireSampleLeftover, type RequestRow } from "@/features/requests"; import { approveRequest } from "@/features/review"; +import { lockSeedRun } from "./repository"; import { freshSampleMail, recordedAiClient, SAMPLES, sampleRecording, type Sample } from "./samples"; // While a sample of a purpose is in one of these statuses it still serves that purpose – no new one. @@ -19,8 +20,14 @@ const STILL_SERVING: Record export interface SampleDeps extends IntakeDeps { /** The normal export path (ERP adapter, reviewed values) – the exported sample runs through it. */ export: ExportDeps; + /** How long a run waits for another run of the same company (#93); default 20 s. */ + seedLockTimeoutMs?: number; } +// Below the 30 s `statement_timeout` of pool and role (`src/db/client.ts`, Supabase bootstrap), which would +// otherwise end the wait first (#99 review). +const SEED_LOCK_TIMEOUT_MS = 20_000; + const LEFTOVER_REASON = "Beispiel durch einen neuen Lauf ersetzt – das Anlegen war abgebrochen."; export interface SeedResult { @@ -64,16 +71,21 @@ function recordingSender(queue?: JobSender): { sender: JobSender; jobId: () => s * the ERP reference without doing anything. */ export async function seedSamples(deps: SampleDeps, actor: Actor): Promise { - const retired = await retireLeftovers(deps, actor); - const seeded: SeededSample[] = []; - for (const sample of SAMPLES) { - const serving = await deps.tenancy.withTenant(actor.companyId, (tx) => countSamples(tx, STILL_SERVING[sample.purpose])); - if (serving > 0) continue; - const requestId = await createSample(deps, actor, sample); - const status = sample.purpose === "exported" ? await approveAndExport(deps, actor, requestId) : "REVIEW"; - seeded.push({ key: sample.key, requestId, status }); - } - return { seeded, retired }; + // One run per company at a time (#93): this transaction only holds the lock; every step of the run + // commits in its own transaction, as before. A second run waits, then sees the samples in place. + return deps.tenancy.withTenant(actor.companyId, async (lock) => { + await lockSeedRun(lock, deps.seedLockTimeoutMs ?? SEED_LOCK_TIMEOUT_MS); + const retired = await retireLeftovers(deps, actor); + const seeded: SeededSample[] = []; + for (const sample of SAMPLES) { + const serving = await deps.tenancy.withTenant(actor.companyId, (tx) => countSamples(tx, STILL_SERVING[sample.purpose])); + if (serving > 0) continue; + const requestId = await createSample(deps, actor, sample); + const status = sample.purpose === "exported" ? await approveAndExport(deps, actor, requestId) : "REVIEW"; + seeded.push({ key: sample.key, requestId, status }); + } + return { seeded, retired }; + }); } /** diff --git a/src/seed-samples.ts b/src/seed-samples.ts index 2803040..a241705 100644 --- a/src/seed-samples.ts +++ b/src/seed-samples.ts @@ -18,7 +18,8 @@ async function main(): Promise { const password = process.env.SEED_PASSWORD; if (!password) throw new Error("SEED_PASSWORD is required (the demo admins sign in to seed their company)"); const config = loadConfig(); - const database = createDatabase(config.databaseUrl, { max: 2 }); + // One connection holds the per-company seed lock for the whole run (#93); the steps use the others. + const database = createDatabase(config.databaseUrl, { max: 3 }); const auth = createAuth(database.db, config.auth); const storage = new S3BlobStore(config.storage); const boss = await createJobQueue(config.databaseUrl); diff --git a/tests/integration/samples.test.ts b/tests/integration/samples.test.ts index 83b526d..106bb43 100644 --- a/tests/integration/samples.test.ts +++ b/tests/integration/samples.test.ts @@ -10,7 +10,7 @@ import { getActor, type Actor } from "@/features/identity"; import { QUEUES, ReprocessRefused, reprocessRequest } from "@/features/jobs"; import { getRequest, lockRequest, transitionRequest } from "@/features/requests"; import { currentFieldValues, currentLineItemValues, loadReview } from "@/features/review"; -import { RECORDED_MODEL_PREFIX, seedSamples, type SampleDeps } from "@/features/samples"; +import { RECORDED_MODEL_PREFIX, SeedRunBusy, seedSamples, type SampleDeps } from "@/features/samples"; import { S3BlobStore } from "@/features/storage"; import { createTenancy, type Tenancy } from "@/features/tenancy"; import { companyWithAdmin, createStack, type Stack } from "./helpers/stack"; @@ -51,7 +51,7 @@ describe("prepared samples", () => { tenancy = createTenancy(stack.database.db); storage = new S3BlobStore(loadConfig().storage); boss = await createJobQueue(loadConfig().databaseUrl); - for (let i = 0; i < 3; i++) actors.push((await getActor(stack.auth, stack.database.db, new Headers({ cookie: (await companyWithAdmin(stack)).cookie })))!); + for (let i = 0; i < 5; i++) actors.push((await getActor(stack.auth, stack.database.db, new Headers({ cookie: (await companyWithAdmin(stack)).cookie })))!); }); afterAll(async () => { @@ -186,6 +186,40 @@ describe("prepared samples", () => { expect((await samplesOf(admin)).map((sample) => sample.status).sort()).toEqual(["EXPORTED", "REJECTED", "REVIEW"]); }); + it("serializes concurrent runs for one company: each sample once, nothing of the other run rejected (#93)", async () => { + const admin = actors[3]!; + + const runs = await Promise.all([seedSamples(deps(), admin), seedSamples(deps(), admin)]); + + expect(runs.flatMap((run) => run.seeded.map(({ key }) => key)).sort()).toEqual(["pumpe-p204", "werk-ost"]); + expect(runs.map((run) => run.retired)).toEqual([0, 0]); + expect((await samplesOf(admin)).map((sample) => sample.status).sort()).toEqual(["EXPORTED", "REVIEW"]); + }); + + it("a second run stops with a clear message while the first holds the company too long (#93)", async () => { + const admin = actors[4]!; + let entered!: () => void; + let release!: () => void; + const inside = new Promise((resolve) => (entered = resolve)); + const held = new Promise((resolve) => (release = resolve)); + // The first run stalls while storing its first original – inside its run, so it holds the company. + const slowStorage = Object.assign(Object.create(storage) as S3BlobStore, { + put: async (...args: Parameters) => { + entered(); + await held; + return storage.put(...args); + }, + }); + const first = seedSamples(deps({ storage: slowStorage }), admin); + await inside; + + await expect(seedSamples({ ...deps(), seedLockTimeoutMs: 300 }, admin)).rejects.toBeInstanceOf(SeedRunBusy); + + release(); + expect((await first).seeded.map(({ key }) => key)).toEqual(["werk-ost", "pumpe-p204"]); + expect((await samplesOf(admin)).map((sample) => sample.status).sort()).toEqual(["EXPORTED", "REVIEW"]); + }); + it("the database refuses an unknown request source (migration 0019)", async () => { const [admin] = actors as [Actor];