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
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,10 @@ This file records what changes **in the product** – process and session state
existing `auth` schema only when the app owns it.

### Fixed
- Request detail: while a failed attempt waits for its retry, the page shows the same facts as the list –
last error, attempts and the next retry – instead of "Die Dokumente werden gerade ausgewertet". On the
showcase, opening the list or a request picks up due retries (at most once per 30 s and never while
the previous one still runs) instead of waiting for the daily cron (#70).
- AI service: model calls answered with 429 or 503 (provider overload) are retried up to three attempts
in total with backoff instead of failing the extraction at once; all attempts together stay within
the model timeout, and timeouts are not retried (#69).
Expand Down
4 changes: 2 additions & 2 deletions docs/technical/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,7 @@ Deliberately accepted risks – without an entry here a deviation counts as a de
| Exception | Why accepted | Owner | Expires |
|---|---|---|---|
| No RLS on the `identity` (Better Auth) and `pgboss` schemas | Not company-owned business data; reachable only by server code (ADR-0001 D7) | Fluory | 2026-12-31 (review at M3) |
| Showcase without unattended retries (Vercel Hobby cron once/day) | Showcase only; production runs a worker (D2). Jobs run via `after()` on upload/approval/reprocess and `/api/jobs/drain` (#59) | Fluory | when a production-like demo is needed |
| Showcase without unattended retries (Vercel Hobby cron once/day) | Showcase only; production runs a worker (D2). Jobs run via `after()` on upload/approval/reprocess, on views of the request list and detail (at most once per 30 s per instance, never while its previous page-view drain runs, #70) and `/api/jobs/drain` (#59) | Fluory | when a production-like demo is needed |
| Better Auth admin plugin mounted without any holder of its admin role | ADR-0001 D6 names the plugin; decided in #30: kept – its `banned` field implements deactivation (sign-in blocked by the plugin). Nobody holds `platform-admin`, so `/api/auth/admin/*` rejects every caller (tested); user management runs through `identity` | Fluory | 2026-12-31 (review at M3) |
| Upload cap per person only where configured | Local and CI run without `UPLOAD_MAX_PER_HOUR`; the showcase refuses to start without it (#59); concurrent uploads may pass the check together – a cost cap, not an exact quota | Fluory | 2026-12-31 (review at M3) |
| `.msg` uploads checked by OLE signature only | Structure check of Outlook messages needs a CFB parser; files are served only as attachments with `nosniff` and parsed later by the stateless AI service | Fluory | with #23 (MSG parsing) |
Expand All @@ -64,7 +64,7 @@ worker ─► jobs.drain ─► extraction ─► AI service (bytes in, segments
─► requests(REVIEW) + fields + audit ── one transaction
review ─► corrections + approve ─► requests(APPROVED) + export job + audit
worker ─► export ─► ERP (Idempotency-Key) ─► requests(EXPORTED)
showcase (no worker): after() / cron ─► /api/jobs/drain ─► the same drain round (src/job-drain.ts)
showcase (no worker): after() (upload, approval, reprocess, page view) / cron ─► the same drain round (src/job-drain.ts)
failure at any step ─► retry with backoff ─► dead letter ─► requests(ERROR, visible cause)
```

Expand Down
41 changes: 40 additions & 1 deletion src/app/_server/drain-request.test.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
import { after } from "next/server";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { captureLogs } from "@/features/observability";
import { handleDrainRequest, isAuthorized, processBudgetMs, scheduleAfterResponse } from "./drain-request";
import { createRunGate, handleDrainRequest, isAuthorized, processBudgetMs, scheduleAfterResponse } from "./drain-request";

// `after()` is the Next.js runtime boundary (it needs a request scope): replaced by a recorder.
vi.mock("next/server", async (original) => ({ ...(await original<typeof import("next/server")>()), after: vi.fn() }));
Expand Down Expand Up @@ -145,3 +145,42 @@ describe("processBudgetMs (#59 review)", () => {
expect(processBudgetMs(50_000, 1_000, 90_000)).toBe(0);
});
});

// #70: page views pick up due retries on the showcase – at most once per interval and instance, and
// never while this instance's previous drain still runs (a drain can outlast the interval, #73 review).
describe("createRunGate", () => {
it("allows the first start and then nothing until the interval has passed", () => {
const gate = createRunGate(30_000, 300_000);

expect(gate.tryStart(1_000)).toBe(true);
gate.finish();
expect(gate.tryStart(1_001)).toBe(false);
expect(gate.tryStart(30_999)).toBe(false);
expect(gate.tryStart(31_000)).toBe(true);
});

it("blocks while the previous run is still going, even after the interval", () => {
const gate = createRunGate(30_000, 300_000);

expect(gate.tryStart(0)).toBe(true);
expect(gate.tryStart(90_000)).toBe(false);
gate.finish();
expect(gate.tryStart(90_001)).toBe(true);
});

it("treats a run as gone after the stale limit, so a lost run cannot block page views forever", () => {
const gate = createRunGate(30_000, 300_000);

expect(gate.tryStart(0)).toBe(true);
expect(gate.tryStart(299_999)).toBe(false);
expect(gate.tryStart(300_000)).toBe(true);
});

it("keeps separate state per gate", () => {
const first = createRunGate(30_000, 300_000);
const second = createRunGate(30_000, 300_000);

expect(first.tryStart(0)).toBe(true);
expect(second.tryStart(0)).toBe(true);
});
});
32 changes: 32 additions & 0 deletions src/app/_server/drain-request.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,38 @@ export function processBudgetMs(budgetMs: number, invokedAt: number, now: number
return Math.max(0, budgetMs - (now - invokedAt));
}

export interface RunGate {
/** `true` when a run may start now; the caller must call `finish()` when it ends. */
tryStart(now?: number): boolean;
finish(): void;
}

/**
* Page views may pick up due retries on the showcase (#70), but a burst of views must not start a burst
* of drains: at most one start per `minIntervalMs`, and none while the previous run still goes – a drain
* can outlast the interval (#73 review). A run older than `staleAfterMs` (the function limit) counts as
* gone, so a lost `finish()` cannot block page views forever. State lives in the closure, so it is per
* function instance – a cost guard, not a global lock (pg-boss keeps the jobs themselves exactly-once).
*/
export function createRunGate(minIntervalMs: number, staleAfterMs: number): RunGate {
let lastStart: number | undefined;
let running = false;
return {
tryStart(now = Date.now()) {
if (lastStart !== undefined) {
const since = now - lastStart;
if (since < minIntervalMs || (running && since < staleAfterMs)) return false;
}
lastStart = now;
running = true;
return true;
},
finish() {
running = false;
},
};
}

/**
* Runs `run` via `after()` once the response is sent – only when `enabled` (JOB_DRAIN_INLINE=true).
* Failures never reach the user's response; they are logged by error class only.
Expand Down
74 changes: 74 additions & 0 deletions src/app/_server/drain.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,74 @@
import { after } from "next/server";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";

// Page-view trigger of the serverless drain (#70, #73 review): the wiring between the inline-drain switch,
// the run gate and `after()`. External boundaries are replaced – `after()` (needs a request scope), the
// runtime (database, storage) and the drain round (pg-boss); tests/integration/drain.test.ts drains real
// queues with the same round.
vi.mock("next/server", async (original) => ({ ...(await original<typeof import("next/server")>()), after: vi.fn() }));
const runtime = { drainInline: true };
vi.mock("./runtime", () => ({
getRuntime: () => ({ config: { jobs: { drainInline: runtime.drainInline } }, tenancy: {}, storage: {} }),
getJobClient: async () => ({}),
}));
const round = { processing: { processed: 0, failed: 0, deadLettered: 0 }, exports: { exported: 0, failed: 0, deadLettered: 0 } };
const drainRound = vi.fn(async () => round);
vi.mock("@/job-drain", () => ({ buildJobDeps: () => ({}), drainRound: () => drainRound(), handledJobs: () => 0 }));

const afterMock = vi.mocked(after);
const runScheduled = (index: number) => (afterMock.mock.calls[index]![0] as () => Promise<void>)();

describe("drainOnPageView", () => {
let drainOnPageView: () => void;

beforeEach(async () => {
vi.useFakeTimers({ now: new Date("2026-09-27T08:00:00Z") });
vi.resetModules(); // the gate lives in the module: a fresh one per test
afterMock.mockReset();
drainRound.mockClear();
runtime.drainInline = true;
({ drainOnPageView } = await import("./drain"));
});
afterEach(() => vi.useRealTimers());

it("does nothing without JOB_DRAIN_INLINE (local, worker-based production)", () => {
runtime.drainInline = false;

drainOnPageView();

expect(afterMock).not.toHaveBeenCalled();
});

it("schedules one drain after the response and ignores a second view right away", async () => {
drainOnPageView();
drainOnPageView();

expect(afterMock).toHaveBeenCalledTimes(1);
expect(drainRound).not.toHaveBeenCalled(); // not before the response is sent
await runScheduled(0);
expect(drainRound).toHaveBeenCalledTimes(1);
});

it("starts no second drain while the first still runs, and the next one once it finished", async () => {
drainOnPageView();
vi.advanceTimersByTime(60_000);
drainOnPageView();
expect(afterMock).toHaveBeenCalledTimes(1);

await runScheduled(0);
drainOnPageView();

expect(afterMock).toHaveBeenCalledTimes(2);
});

it("frees the gate when the drain fails, so later views can retry", async () => {
drainRound.mockRejectedValueOnce(new Error("synthetic"));
drainOnPageView();
await runScheduled(0);
vi.advanceTimersByTime(30_000);

drainOnPageView();

expect(afterMock).toHaveBeenCalledTimes(2);
});
});
23 changes: 22 additions & 1 deletion src/app/_server/drain.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
import { SERVERLESS_DRAIN } from "@/config/env";
import { buildJobDeps, drainRound, handledJobs, type DrainRoundResult, type JobDeps } from "@/job-drain";
import { logEvent } from "@/features/observability";
import { processBudgetMs, scheduleAfterResponse } from "./drain-request";
import { createRunGate, processBudgetMs, scheduleAfterResponse } from "./drain-request";
import { getJobClient, getRuntime } from "./runtime";

// Serverless drain of the web process (#59, ADR-0001 D2): the showcase has no worker, so the route
Expand Down Expand Up @@ -44,3 +44,24 @@ export async function drainNow(invokedAt: number = Date.now()): Promise<DrainRou
export function drainAfterResponse(invokedAt: number): void {
scheduleAfterResponse(getRuntime().config.jobs.drainInline, () => drainNow(invokedAt));
}

// Stale after the function limit (maxDuration 300 s of the pages): by then the run is over either way.
const pageViewGate = createRunGate(30_000, 300_000);

/**
* Opening the request list or a request (#70): with JOB_DRAIN_INLINE=true a due retry is picked up by
* normal use instead of waiting for the daily cron – at most once per 30 s and function instance, and
* never while this instance's previous page-view drain still runs.
*/
export function drainOnPageView(): void {
const { drainInline } = getRuntime().config.jobs;
if (!drainInline || !pageViewGate.tryStart()) return;
const invokedAt = Date.now();
scheduleAfterResponse(drainInline, async () => {
try {
await drainNow(invokedAt);
} finally {
pageViewGate.finish();
}
});
}
13 changes: 12 additions & 1 deletion src/app/requests/[id]/page.tsx
Original file line number Diff line number Diff line change
@@ -1,7 +1,10 @@
import Link from "next/link";
import { notFound, redirect } from "next/navigation";
import { drainOnPageView } from "@/app/_server/drain";
import { getRuntime, requestActor } from "@/app/_server/runtime";
import { duplicateDecidable, loadReview, REJECTION_REASON_MAX, type ReviewField } from "@/features/review";
import { processingNotice } from "../processing-notice";
import { requestRowView } from "../row-view";
import { StatusPill } from "../status-pill";
import { DONE_MESSAGES, ERROR_MESSAGES, messageFor } from "./messages";
import { approveAction, confirmNotDuplicateAction, correctFieldAction, rejectAction, rejectAsDuplicateAction } from "./actions";
Expand All @@ -15,7 +18,8 @@ const UUID = /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i;
const dateFormat = new Intl.DateTimeFormat("de-DE", { day: "2-digit", month: "2-digit", year: "numeric", hour: "2-digit", minute: "2-digit", timeZone: "Europe/Berlin" });
const NO_FIELDS: Record<string, string> = {
NEW: "Die Anfrage wartet auf die Verarbeitung. Erkannte Angaben erscheinen danach hier.",
PROCESSING: "Die Dokumente werden gerade ausgewertet.",
// The processing state itself (running, failed attempt, next retry) is the notice above (#70).
PROCESSING: "Erkannte Angaben erscheinen hier nach der Auswertung.",
};

// Review page (#8, #25), layout after design prototype A (#52): fields, positions and documents on the
Expand All @@ -33,6 +37,7 @@ export default async function RequestPage({
if (!UUID.test(id)) notFound();
const view = await loadReview(getRuntime().tenancy, actor, id);
if (!view) notFound();
drainOnPageView();
const { request, fields, lineItems, documents, skippedDocuments, documentNotes, exportRecord } = view;
const query = await searchParams;
// `?field=<key>` selects a header field, `?field=<key>&item=<n>` a line-item field (#25). Without a (valid)
Expand All @@ -50,6 +55,7 @@ export default async function RequestPage({
const canDecideDuplicate = duplicateDecidable(request);
const done = messageFor(DONE_MESSAGES, query.done);
const error = messageFor(ERROR_MESSAGES, query.error);
const notice = processingNotice(request.status, requestRowView(request, exportRecord ?? undefined), (date) => dateFormat.format(date));
const openItems = lineItems.reduce((count, item) => count + item.fields.filter(needsAttention).length, fields.filter(needsAttention).length);
// The anchor brings the panel into view on narrow screens, where it sits below the tables.
const fieldHref = (field: ReviewField) =>
Expand Down Expand Up @@ -84,6 +90,11 @@ export default async function RequestPage({
<strong>Fehler:</strong> {request.errorMessage}
</p>
)}
{notice && (
<p className={notice.tone === "warning" ? "callout callout-warn" : "callout callout-info"} role="status" data-testid="processing-notice">
{notice.text}
</p>
)}
{request.status === "REJECTED" && request.rejectionReason && (
<p className="callout callout-grey">
<strong>Abgelehnt:</strong> {request.rejectionReason}
Expand Down
2 changes: 2 additions & 0 deletions src/app/requests/page.tsx
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import Link from "next/link";
import { redirect } from "next/navigation";
import { drainOnPageView } from "@/app/_server/drain";
import { getRuntime, requestActor } from "@/app/_server/runtime";
import { listExportRecords } from "@/features/export";
import { listRequests, parseCursor, type RequestFilter } from "@/features/requests";
Expand Down Expand Up @@ -42,6 +43,7 @@ function pageHref(filter: RequestFilter, after?: string): string {
export default async function RequestsPage({ searchParams }: { searchParams: Promise<Record<string, string | undefined>> }) {
const actor = await requestActor();
if (!actor) redirect("/login");
drainOnPageView();
const query = await searchParams;
const filter = filterOf(query);
const after = parseCursor(query.after) ?? undefined;
Expand Down
49 changes: 49 additions & 0 deletions src/app/requests/processing-notice.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,49 @@
import { describe, expect, it } from "vitest";
import { processingNotice } from "./processing-notice";
import { requestRowView, type RowRequest } from "./row-view";

// #70: the detail page tells the same story as the list – it is derived from the same row view (#73
// review), so a failed attempt waiting for its retry is never "being evaluated".
const time = (date: Date) => date.toISOString().slice(11, 16);
const now = new Date("2026-09-27T08:00:00Z");
const base: RowRequest = { status: "PROCESSING", attempts: 0, errorStage: null, errorMessage: null, nextRetryAt: null };
const notice = (request: Partial<RowRequest>) => {
const row = { ...base, ...request };
return processingNotice(row.status, requestRowView(row, undefined), time, now);
};

describe("processingNotice", () => {
it("reports ongoing work while no attempt has failed", () => {
expect(notice({})).toEqual({ tone: "info", text: "Die Dokumente werden gerade ausgewertet." });
});

it("shows the last error, the attempts and the next retry of a failed attempt", () => {
expect(notice({ attempts: 2, errorMessage: "Der KI-Dienst ist nicht erreichbar.", nextRetryAt: new Date("2026-09-27T08:15:00Z") })).toEqual({
tone: "warning",
text: "Letzter Versuch fehlgeschlagen: Der KI-Dienst ist nicht erreichbar. Bisher 2 Versuche. Nächster Versuch: 08:15.",
});
});

it("says a retry is due once its time has passed", () => {
expect(notice({ attempts: 1, errorMessage: "Der KI-Dienst ist nicht erreichbar.", nextRetryAt: new Date("2026-09-27T07:55:00Z") })?.text).toBe(
"Letzter Versuch fehlgeschlagen: Der KI-Dienst ist nicht erreichbar. Bisher 1 Versuch. Nächster Versuch fällig seit 07:55.",
);
});

it("promises no retry when none is scheduled (e.g. attempts used up) – like the list, which shows none", () => {
expect(notice({ attempts: 3, errorMessage: "Zeitüberschreitung beim KI-Dienst.", nextRetryAt: null })?.text).toBe(
"Letzter Versuch fehlgeschlagen: Zeitüberschreitung beim KI-Dienst. Bisher 3 Versuche.",
);
});

it("shows a failed attempt of a request that is still NEW, as the list does", () => {
expect(notice({ status: "NEW", attempts: 1, errorMessage: "Der KI-Dienst ist nicht erreichbar." })?.tone).toBe("warning");
});

it("has nothing to say for a NEW request without a failed attempt or outside of processing", () => {
expect(notice({ status: "NEW" })).toBeNull();
for (const status of ["REVIEW", "APPROVED", "EXPORTED", "REJECTED", "ERROR"]) {
expect(notice({ status, attempts: 1, errorMessage: "egal" })).toBeNull();
}
});
});
21 changes: 21 additions & 0 deletions src/app/requests/processing-notice.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
import type { RowView } from "./row-view";

// What the detail page says while a request is being processed (#70). Derived from the list's row view
// (#73 review), so list and detail show the same facts: last error, attempts, next retry.
export interface ProcessingNotice {
tone: "info" | "warning";
text: string;
}

export function processingNotice(status: string, row: RowView, formatTime: (date: Date) => string, now: Date = new Date()): ProcessingNotice | null {
if (status !== "NEW" && status !== "PROCESSING") return null;
if (!row.error) return status === "PROCESSING" ? { tone: "info", text: "Die Dokumente werden gerade ausgewertet." } : null;
const attempts = `Bisher ${row.attempts} ${row.attempts === 1 ? "Versuch" : "Versuche"}.`;
// No scheduled retry (e.g. attempts used up) → no promise; the list shows none either.
const next = !row.nextRetryAt
? ""
: row.nextRetryAt.getTime() <= now.getTime()
? ` Nächster Versuch fällig seit ${formatTime(row.nextRetryAt)}.`
: ` Nächster Versuch: ${formatTime(row.nextRetryAt)}.`;
return { tone: "warning", text: `Letzter Versuch fehlgeschlagen: ${row.error} ${attempts}${next}` };
}
Loading