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
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,11 @@ This file records what changes **in the product** – process and session state

## [Unreleased]

### Changed
- `pnpm seed:samples` first settles samples an aborted run left behind (new, in progress, or failed in
processing): rejected with a fixed reason and an audit event, never deleted – the list no longer keeps
such leftovers. An approved sample whose export failed is kept (#84).

### Fixed
- `/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
Expand Down
3 changes: 3 additions & 0 deletions docs/technical/deployment-vercel.md
Original file line number Diff line number Diff line change
Expand Up @@ -124,6 +124,9 @@ route; if it is not reachable, the export is queued and the next drain retries i
once with `pnpm samples:record` (`src/features/samples/data/`) – no model call, so the samples work even while the
model provider is overloaded. List and detail label them „Vorbereitetes Beispiel – aufgezeichnete KI-Antwort“. The
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“.

## 7. Jobs without a worker

Expand Down
3 changes: 3 additions & 0 deletions src/features/requests/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ export {
countRequestsByStatus,
countRequestsCreatedBy,
countSamples,
listSampleLeftoverIds,
createRequest,
findDuplicate,
getRequest,
Expand All @@ -14,9 +15,11 @@ export {
recordDuplicateDecision,
recordExportRetry,
recordProcessingFailure,
retireSampleLeftover,
transitionRequest,
type NewRequest,
type RequestRow,
} from "./repository";
export { parseCursor, REQUEST_PAGE_SIZE } from "./cursor";
export { canTransition, InvalidTransition, nextStatus, type ErrorStage, type RequestEvent } from "./status";
export { isSampleLeftover } from "./sample-leftover";
25 changes: 25 additions & 0 deletions src/features/requests/repository.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ import { and, asc, desc, eq, gte, inArray, or, sql, type SQL } from "drizzle-orm
import { requests, type RequestStatus } from "@/db/schema";
import { tenantOf, type TenantTx } from "@/features/tenancy";
import { parseCursor, REQUEST_PAGE_SIZE } from "./cursor";
import { isSampleLeftover } from "./sample-leftover";
import { nextStatus, type RequestEvent } from "./status";

export type RequestRow = typeof requests.$inferSelect;
Expand Down Expand Up @@ -166,6 +167,30 @@ export async function countSamples(tx: TenantTx, statuses: readonly RequestStatu
return row?.count ?? 0;
}

/** Ids of the tenant's samples an aborted seed run left behind (#84) – see `isSampleLeftover`. */
export async function listSampleLeftoverIds(tx: TenantTx): Promise<string[]> {
tenantOf(tx);
const rows = await tx
.select({ id: requests.id })
.from(requests)
.where(
and(
eq(requests.source, "sample"),
or(inArray(requests.status, ["NEW", "PROCESSING"]), and(eq(requests.status, "ERROR"), eq(requests.errorStage, "processing"))),
),
);
return rows.map((row) => row.id);
}

/**
* Settles a sample leftover on a locked row (#84): rejected with `reason`. The rule lives here, not with
* the caller (#92 review) – any other request, or a sample that is no leftover, is left alone (null).
*/
export async function retireSampleLeftover(tx: TenantTx, row: RequestRow, reason: string): Promise<RequestRow | null> {
if (!isSampleLeftover(row)) return null;
return transitionRequest(tx, row, "sample.retired", { rejectionReason: reason, errorStage: null, errorMessage: null, nextRetryAt: null });
}

/** Requests a user created since `since` (upload rate limit, #59 review) – within the tenant. */
export async function countRequestsCreatedBy(tx: TenantTx, userId: string, since: Date): Promise<number> {
tenantOf(tx);
Expand Down
30 changes: 30 additions & 0 deletions src/features/requests/sample-leftover.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
import { describe, expect, it } from "vitest";
import { isSampleLeftover } from "./sample-leftover";

// #84 + #92 review: only what an aborted seed run leaves behind – never a decided sample, never an
// approved one whose export failed (it may already be in the ERP and can be exported again).
const sample = { source: "sample" as const, errorStage: null };

describe("isSampleLeftover", () => {
it.each(["NEW", "PROCESSING"] as const)("counts a sample in %s", (status) => {
expect(isSampleLeftover({ ...sample, status })).toBe(true);
});

it("counts a sample whose processing failed", () => {
expect(isSampleLeftover({ ...sample, status: "ERROR", errorStage: "processing" })).toBe(true);
});

it("never counts an approved sample whose export failed", () => {
expect(isSampleLeftover({ ...sample, status: "ERROR", errorStage: "export" })).toBe(false);
});

it.each(["REVIEW", "APPROVED", "EXPORTED", "REJECTED"] as const)("never counts a decided or settled sample (%s)", (status) => {
expect(isSampleLeftover({ ...sample, status })).toBe(false);
});

it("never counts an upload, whatever its state", () => {
for (const status of ["NEW", "PROCESSING", "ERROR"] as const) {
expect(isSampleLeftover({ source: "upload", status, errorStage: "processing" })).toBe(false);
}
});
});
20 changes: 20 additions & 0 deletions src/features/requests/sample-leftover.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
import type { RequestSource, RequestStatus } from "@/db/schema";

// Its own input type from the schema – importing RequestRow from the repository would close a cycle.
export interface SampleLeftoverCandidate {
source: RequestSource;
status: RequestStatus;
errorStage: "processing" | "export" | null;
}

/**
* A prepared sample (#71) that an aborted seed run left behind (#84). Processing of a sample is never
* queued, so NEW, PROCESSING and ERROR from processing only mean the run broke off. Never a decided sample,
* and never an approved one whose export failed (ERROR from export): it may already be in the ERP and can be
* exported again (#92 review).
*/
export function isSampleLeftover(request: SampleLeftoverCandidate): boolean {
if (request.source !== "sample") return false;
if (request.status === "NEW" || request.status === "PROCESSING") return true;
return request.status === "ERROR" && request.errorStage === "processing";
}
8 changes: 8 additions & 0 deletions src/features/requests/status.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,11 +18,19 @@ describe("request status machine (ADR-0001: NEW → PROCESSING → REVIEW → AP
["NEW", "reject.duplicate", "REJECTED"],
["REVIEW", "reject.duplicate", "REJECTED"],
["ERROR", "reject.duplicate", "REJECTED"],
// A prepared sample left behind by an aborted seed run is settled (#84) – never a decided one.
["NEW", "sample.retired", "REJECTED"],
["PROCESSING", "sample.retired", "REJECTED"],
["ERROR", "sample.retired", "REJECTED"],
] as const)("%s --%s--> %s", (from, event, to) => {
expect(nextStatus(from, event as RequestEvent)).toBe(to);
});

it.each([
["REVIEW", "sample.retired"],
["APPROVED", "sample.retired"],
["EXPORTED", "sample.retired"],
["REJECTED", "sample.retired"],
["REVIEW", "processing.succeeded"],
["EXPORTED", "processing.started"],
["NEW", "approve"],
Expand Down
7 changes: 6 additions & 1 deletion src/features/requests/status.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,8 @@ export type RequestEvent =
| "reject.duplicate"
| "export.succeeded"
| "export.failed"
| "reprocess.export";
| "reprocess.export"
| "sample.retired";

export type ErrorStage = "processing" | "export";

Expand All @@ -31,6 +32,10 @@ const TRANSITIONS: Record<RequestEvent, Partial<Record<RequestStatus, RequestSta
"export.succeeded": { APPROVED: "EXPORTED" },
"export.failed": { APPROVED: "ERROR" },
"reprocess.export": { ERROR: "APPROVED" },
// A prepared sample (#71) left behind by an aborted seed run (#84) – only through
// `retireSampleLeftover`, which checks `isSampleLeftover` (sample; NEW, PROCESSING or ERROR from
// processing – never an approved sample whose export failed, #92 review).
"sample.retired": { NEW: "REJECTED", PROCESSING: "REJECTED", ERROR: "REJECTED" },
};

export class InvalidTransition extends Error {
Expand Down
2 changes: 1 addition & 1 deletion src/features/samples/index.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,3 @@
// 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 { seedSamples, type SampleDeps, type SeededSample } from "./seed";
export { seedSamples, type SampleDeps, type SeededSample, type SeedResult } from "./seed";
36 changes: 33 additions & 3 deletions src/features/samples/seed.ts
Original file line number Diff line number Diff line change
@@ -1,9 +1,10 @@
import { randomUUID } from "node:crypto";
import { recordAudit } from "@/features/audit";
import { exportRequestJob, type ExportDeps } from "@/features/export";
import type { Actor } from "@/features/identity";
import { submitUpload, type IntakeDeps } from "@/features/intake";
import { markProcessingFailed, processRequestJob, type JobSender } from "@/features/jobs";
import { countSamples, type RequestRow } from "@/features/requests";
import { countSamples, listSampleLeftoverIds, lockRequest, retireSampleLeftover, type RequestRow } from "@/features/requests";
import { approveRequest } from "@/features/review";
import { freshSampleMail, recordedAiClient, SAMPLES, sampleRecording, type Sample } from "./samples";

Expand All @@ -20,6 +21,14 @@ export interface SampleDeps extends IntakeDeps {
export: ExportDeps;
}

const LEFTOVER_REASON = "Beispiel durch einen neuen Lauf ersetzt – das Anlegen war abgebrochen.";

export interface SeedResult {
seeded: SeededSample[];
/** Leftovers of aborted runs settled as rejected (#84). */
retired: number;
}

export interface SeededSample {
key: string;
requestId: string;
Expand Down Expand Up @@ -54,7 +63,8 @@ function recordingSender(queue?: JobSender): { sender: JobSender; jobId: () => s
* a sample to the live model; and the export job queued by the approval runs right away, so visitors see
* the ERP reference without doing anything.
*/
export async function seedSamples(deps: SampleDeps, actor: Actor): Promise<SeededSample[]> {
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]));
Expand All @@ -63,7 +73,27 @@ export async function seedSamples(deps: SampleDeps, actor: Actor): Promise<Seede
const status = sample.purpose === "exported" ? await approveAndExport(deps, actor, requestId) : "REVIEW";
seeded.push({ key: sample.key, requestId, status });
}
return seeded;
return { seeded, retired };
}

/**
* Settles samples an aborted run left behind (#84): through the status machine, with the reason visible
* on the request and an audit event – never deleted. Which samples count is decided in the requests
* module (`isSampleLeftover`, re-checked under the row lock).
*/
async function retireLeftovers(deps: SampleDeps, actor: Actor): Promise<number> {
return deps.tenancy.withTenant(actor.companyId, async (tx) => {
let retired = 0;
for (const id of await listSampleLeftoverIds(tx)) {
const request = await lockRequest(tx, id);
if (!request) continue;
const settled = await retireSampleLeftover(tx, request, LEFTOVER_REASON);
if (!settled) continue;
await recordAudit(tx, { actorUserId: actor.userId, action: "request.sample_retired", entityType: "request", entityId: id, data: { from: request.status } });
retired++;
}
return retired;
});
}

async function createSample(deps: SampleDeps, actor: Actor, sample: Sample): Promise<string> {
Expand Down
3 changes: 2 additions & 1 deletion src/seed-samples.ts
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,8 @@ async function main(): Promise<void> {
const cookie = login.headers.getSetCookie().map((line) => line.split(";")[0]).join("; ");
const actor = await getActor(auth, database.db, new Headers({ cookie }));
if (!actor) throw new Error(`${email} has no company – run pnpm seed:demo first`);
const seeded = await seedSamples(deps, actor);
const { seeded, retired } = await seedSamples(deps, actor);
if (retired > 0) console.log(`${email}: ${retired} leftover sample(s) of an aborted run rejected`);
console.log(seeded.length === 0 ? `${email}: samples still in place` : `${email}: ${seeded.map((sample) => `${sample.key} → ${sample.status}`).join(", ")}`);
}
} finally {
Expand Down
34 changes: 30 additions & 4 deletions tests/integration/samples.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -62,8 +62,9 @@ describe("prepared samples", () => {

it("seeds one sample in review and one approved and exported sample; a later delivery of the export job is a no-op", async () => {
const [admin] = actors as [Actor];
const seeded = await seedSamples(deps(), admin);
const { seeded, retired } = await seedSamples(deps(), admin);

expect(retired).toBe(0);
expect(seeded.map(({ key, status }) => [key, status])).toEqual([
["werk-ost", "REVIEW"],
["pumpe-p204", "EXPORTED"],
Expand Down Expand Up @@ -106,15 +107,15 @@ describe("prepared samples", () => {
const [admin] = actors as [Actor];
const before = await samplesOf(admin);

expect(await seedSamples(deps(), admin)).toEqual([]);
expect(await seedSamples(deps(), admin)).toEqual({ seeded: [], retired: 0 });

expect(await samplesOf(admin)).toEqual(before);
});

it("leaves the export to the queued job when the ERP is unreachable – nothing half-written", async () => {
const admin = actors[1]!;

const seeded = await seedSamples(deps({ erpDown: true }), admin);
const { seeded } = await seedSamples(deps({ erpDown: true }), admin);

const exported = seeded.find((sample) => sample.key === "pumpe-p204")!;
expect(exported.status).toBe("APPROVED");
Expand All @@ -138,6 +139,20 @@ describe("prepared samples", () => {
expect(await jobsFor(QUEUES.exportRequest, exported.id)).toBe(1);
});

it("keeps an approved sample whose export failed – it is no leftover (#92 review)", async () => {
const admin = actors[1]!;
const exported = (await samplesOf(admin)).find((sample) => sample.status === "APPROVED")!;
await stack.database.pool.query("delete from pgboss.job where name = $1 and singleton_key = $2", [QUEUES.exportRequest, exported.id]);
await tenancy.withTenant(admin.companyId, async (tx) =>
transitionRequest(tx, (await lockRequest(tx, exported.id))!, "export.failed", { errorStage: "export", errorMessage: "ERP nicht erreichbar." }),
);

const { retired } = await seedSamples(deps(), admin);

expect(retired).toBe(0);
expect(await tenancy.withTenant(admin.companyId, (tx) => getRequest(tx, exported.id))).toMatchObject({ status: "ERROR", errorStage: "export" });
});

it("an aborted run leaves the sample in ERROR, not in progress; the next run replaces it; no live reprocess", async () => {
const admin = actors[2]!;
const brokenStorage = Object.assign(Object.create(storage) as S3BlobStore, {
Expand All @@ -153,11 +168,22 @@ describe("prepared samples", () => {
// Reprocessing a sample would call the live model – refused.
await expect(reprocessRequest({ tenancy, boss }, admin, failed!.id)).rejects.toBeInstanceOf(ReprocessRefused);

const seeded = await seedSamples(deps(), admin);
const { seeded, retired } = await seedSamples(deps(), admin);
expect(seeded.map(({ key, status }) => [key, status])).toEqual([
["werk-ost", "REVIEW"],
["pumpe-p204", "EXPORTED"],
]);

// #84: the leftover is settled – rejected with a reason and an audit event, never deleted.
expect(retired).toBe(1);
expect(await tenancy.withTenant(admin.companyId, (tx) => getRequest(tx, failed!.id))).toMatchObject({
status: "REJECTED",
rejectionReason: "Beispiel durch einen neuen Lauf ersetzt – das Anlegen war abgebrochen.",
errorMessage: null,
});
const events = await tenancy.withTenant(admin.companyId, (tx) => listAuditEvents(tx, "request", failed!.id));
expect(events.filter((event) => event.action === "request.sample_retired")).toHaveLength(1);
expect((await samplesOf(admin)).map((sample) => sample.status).sort()).toEqual(["EXPORTED", "REJECTED", "REVIEW"]);
});

it("the database refuses an unknown request source (migration 0019)", async () => {
Expand Down
Loading