Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand Down
2 changes: 1 addition & 1 deletion docs/technical/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -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) |
Expand Down
4 changes: 3 additions & 1 deletion docs/technical/deployment-vercel.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
1 change: 1 addition & 0 deletions src/features/samples/index.ts
Original file line number Diff line number Diff line change
@@ -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";
30 changes: 30 additions & 0 deletions src/features/samples/repository.ts
Original file line number Diff line number Diff line change
@@ -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<void> {
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;
}
}
32 changes: 22 additions & 10 deletions src/features/samples/seed.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -19,8 +20,14 @@ const STILL_SERVING: Record<Sample["purpose"], readonly RequestRow["status"][]>
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 {
Expand Down Expand Up @@ -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<SeedResult> {
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 };
});
}

/**
Expand Down
3 changes: 2 additions & 1 deletion src/seed-samples.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,8 @@ async function main(): Promise<void> {
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);
Expand Down
38 changes: 36 additions & 2 deletions tests/integration/samples.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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 () => {
Expand Down Expand Up @@ -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<void>((resolve) => (entered = resolve));
const held = new Promise<void>((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<S3BlobStore["put"]>) => {
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];

Expand Down
Loading