diff --git a/CHANGELOG.md b/CHANGELOG.md index 5716bae..06f3cba 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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). diff --git a/docs/technical/architecture.md b/docs/technical/architecture.md index 915c726..a97e7bc 100644 --- a/docs/technical/architecture.md +++ b/docs/technical/architecture.md @@ -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) | @@ -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) ``` diff --git a/src/app/_server/drain-request.test.ts b/src/app/_server/drain-request.test.ts index 05ec164..04eefa7 100644 --- a/src/app/_server/drain-request.test.ts +++ b/src/app/_server/drain-request.test.ts @@ -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()), after: vi.fn() })); @@ -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); + }); +}); diff --git a/src/app/_server/drain-request.ts b/src/app/_server/drain-request.ts index c9a8a52..fcfa1b8 100644 --- a/src/app/_server/drain-request.ts +++ b/src/app/_server/drain-request.ts @@ -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. diff --git a/src/app/_server/drain.test.ts b/src/app/_server/drain.test.ts new file mode 100644 index 0000000..21161cb --- /dev/null +++ b/src/app/_server/drain.test.ts @@ -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()), 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)(); + +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); + }); +}); diff --git a/src/app/_server/drain.ts b/src/app/_server/drain.ts index e586fff..89f48c6 100644 --- a/src/app/_server/drain.ts +++ b/src/app/_server/drain.ts @@ -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 @@ -44,3 +44,24 @@ export async function drainNow(invokedAt: number = Date.now()): Promise 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(); + } + }); +} diff --git a/src/app/requests/[id]/page.tsx b/src/app/requests/[id]/page.tsx index e10064b..62689aa 100644 --- a/src/app/requests/[id]/page.tsx +++ b/src/app/requests/[id]/page.tsx @@ -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"; @@ -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 = { 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 @@ -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=` selects a header field, `?field=&item=` a line-item field (#25). Without a (valid) @@ -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) => @@ -84,6 +90,11 @@ export default async function RequestPage({ Fehler: {request.errorMessage}

)} + {notice && ( +

+ {notice.text} +

+ )} {request.status === "REJECTED" && request.rejectionReason && (

Abgelehnt: {request.rejectionReason} diff --git a/src/app/requests/page.tsx b/src/app/requests/page.tsx index c63a92e..cb66dfd 100644 --- a/src/app/requests/page.tsx +++ b/src/app/requests/page.tsx @@ -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"; @@ -42,6 +43,7 @@ function pageHref(filter: RequestFilter, after?: string): string { export default async function RequestsPage({ searchParams }: { searchParams: Promise> }) { const actor = await requestActor(); if (!actor) redirect("/login"); + drainOnPageView(); const query = await searchParams; const filter = filterOf(query); const after = parseCursor(query.after) ?? undefined; diff --git a/src/app/requests/processing-notice.test.ts b/src/app/requests/processing-notice.test.ts new file mode 100644 index 0000000..b3bcc29 --- /dev/null +++ b/src/app/requests/processing-notice.test.ts @@ -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) => { + 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(); + } + }); +}); diff --git a/src/app/requests/processing-notice.ts b/src/app/requests/processing-notice.ts new file mode 100644 index 0000000..9c0ce5d --- /dev/null +++ b/src/app/requests/processing-notice.ts @@ -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}` }; +}