// applyEntityEvent — die EINZIGE Schreib-Logik für r.entity-Tabellen aus // Stored-Events. Beide Aufrufer benutzen sie: // // - createEventStoreExecutor (live, im Write-TX) — übergibt das gerade // appendete StoredEvent direkt; sensitive Felder liegen als Tabellen- // Ciphertext im Payload (boot-validiert pii/encrypted, #967). // - rebuildProjection via ImplicitProjection (replay, im Rebuild-TX) — // übergibt dasselbe StoredEvent; apply kopiert Ciphertext byte-gleich // zurück, bidx wird unten neu berechnet. // // Live==Rebuild-Equivalence ist damit by-construction für ALLE Felder — // eine geänderte Schreib-Logik muss nur an EINER Stelle gepflegt werden, // kein Sync-Contract mehr. Einzige legitime Divergenz: Crypto-Shredding // (DEK erased → Wert unlesbar, bidx NULL). Load-bearing Test: // db/__tests__/implicit-projection-equivalence.integration.test.ts. // // Tenant-Isolation: applyEntityEvent erwartet einen rohen DbRunner (TX // oder pool), KEINEN TenantDb-Wrapper. Schutz kommt aus zwei Quellen: // 1. Live-Pfad ruft VOR der Schreibung loadById (tenant-scoped) für // update/delete/restore — die aggregateId ist also schon tenant- // validiert bevor wir hier ankommen. // 2. Bei create wird tenantId explizit aus event.tenantId gesetzt, also // nie über den TenantDb-Wrapper-Default abgeleitet. // Damit ist der TenantDb-Wrapper-Loss in dieser Funktion funktional ohne // Sicherheitslücke. // // Auto-Verben: // .created → INSERT // .updated → UPDATE WHERE id=aggregateId // .deleted → soft-delete-UPDATE wenn entity.softDelete, sonst hard-DELETE // .restored → undelete-UPDATE (nur bei softDelete sinnvoll) // .forgotten → hard-DELETE immer (Art.17-Purge, auch bei softDelete); // via executor.forget(). Rebuild replayt created→forgotten // → Row weg, rebuild-safe (was ein direktes deleteMany nicht ist). // // Domain-Events (r.defineEvent) auf demselben Aggregate werden hier NICHT // behandelt — die liefen im Live-Pfad nie durch den Executor und müssen // von expliziten r.projection-apply-Handlern oder r.multiStreamProjection // behandelt werden. ImplicitProjection registriert die Auto-Verben. // // Return-Shape: ApplyResult mit `kind` + optionaler `row`. // - "applied" → Schreibung lief durch. `row` enthält die geschriebene // Row für create/update/soft-delete/restore. Bei hard-delete ist // `row` null (DELETE-Statements geben keine returning-Row her). // - "skipped" → Event ist kein Auto-Verb (Domain-Event auf demselben // Aggregate). Caller no-op. import { collectLookupableFields, computeBlindIndexValues } from "../crypto/blind-index"; import { deleteMany, insertOne, updateMany } from "../db/query"; import type { EntityDefinition } from "../engine/types"; import { InternalError } from "../errors"; import type { StoredEvent } from "../event-store"; import type { DbRow, DbRunner } from "./connection"; import type { TableColumns } from "./dialect"; // biome-ignore lint/suspicious/noExplicitAny: Drizzle-Tabellen sind generisch typed; framework code erasiert die Spalten-Union absichtlich. type Table = TableColumns; export type AutoVerb = "created" | "updated" | "deleted" | "restored" | "forgotten"; const AUTO_VERBS: readonly AutoVerb[] = ["created", "updated", "deleted", "restored", "forgotten"]; export type ApplyResult = | { readonly kind: "applied"; readonly verb: AutoVerb; readonly row: DbRow | null } | { readonly kind: "skipped" }; /** Parsed event.type → AutoVerb wenn das Event eines der Auto-Verben * auf dem gegebenen Aggregate ist. null sonst (Domain-Event). */ export function parseAutoVerb(event: StoredEvent): AutoVerb | null { const prefix = `${event.aggregateType}.`; if (!event.type.startsWith(prefix)) return null; const verb = event.type.slice(prefix.length); // @cast-boundary: verb is validated against the AutoVerb list before the cast. return (AUTO_VERBS as readonly string[]).includes(verb) ? (verb as AutoVerb) : null; } /** Idempotente Anwendung eines Auto-Events auf die Entity-Tabelle. * Wird sowohl beim Live-Append (innerhalb der Write-TX) als auch beim * Rebuild (innerhalb der Rebuild-TX) gerufen — identische Logik. */ export async function applyEntityEvent( event: StoredEvent, table: Table, entity: EntityDefinition, tx: DbRunner, ): Promise { const verb = parseAutoVerb(event); if (verb === null) return { kind: "skipped" }; const softDelete = entity.softDelete ?? false; switch (verb) { case "created": { // tenantId-Resolution explizit, nicht via Spread-Reihenfolge: // Live-Pfad nutzt tx=db.raw (kein TenantDb-Wrapper-Auto-Inject), // beim Replay erst recht keiner. Default = event.tenantId; payload // gewinnt NUR wenn gültig string mit length > 0 (seedTenantMembership- // Pfad: Operator schreibt im Ziel-Tenant, Event im Operator-Tenant). // Pinst durch db/__tests__/apply-entity-event-tenant.integration.ts. // // Fail-loud wenn payload.tenantId gesetzt aber invalid (leer/null/ // non-string): das ist tenant-isolation-kritisch — silent fallback // auf event.tenantId würde eine Bug-payload in den Operator-Tenant // schreiben statt zu failen, was Cross-Tenant-Datendrift erzeugt. const payloadTenantId = event.payload["tenantId"]; let tenantId: string; if (payloadTenantId === undefined) { tenantId = event.tenantId; } else if (typeof payloadTenantId === "string" && payloadTenantId.length > 0) { tenantId = payloadTenantId; } else { throw new InternalError({ message: `applyEntityEvent: payload.tenantId set but invalid (${JSON.stringify(payloadTenantId)}). Tenant-isolation-kritisch: silent fallback auf event.tenantId würde Cross-Tenant-Drift erzeugen.`, }); } const row = await insertOne(tx, table, { ...event.payload, // Blind-Index-Spalten aus dem Payload-Wert (ciphertext → decrypt → // HMAC, plaintext → HMAC, erased → NULL). Hier statt im Executor, // damit Live-Write und Rebuild denselben Wert produzieren und der // deterministische HMAC nie im Event-Log landet. ...(await computeBlindIndexValues(event.payload, collectLookupableFields(entity))), tenantId, id: event.aggregateId, version: event.version, insertedAt: event.createdAt, insertedById: event.createdBy, }); return { kind: "applied", verb, row: row ?? null }; } case "updated": { // payload-Shape: { changes, previous } — siehe event-store-executor.ts. const rawChanges = (event.payload["changes"] ?? {}) as Record; // @cast-boundary engine-payload // Plural file fields (files/images) have no entity-table column — the // array of UUIDs lives in the event payload only. Strip them from the // UPDATE so Postgres doesn't error with "column X does not exist". const changes: Record = {}; for (const [key, value] of Object.entries(rawChanges)) { const fieldType = entity.fields[key]?.type; if (fieldType === "files" || fieldType === "images") continue; changes[key] = value; } const rows = await updateMany( tx, table, { ...changes, // bidx nur für die geänderten lookupable-Felder (siehe created). ...(await computeBlindIndexValues(changes, collectLookupableFields(entity))), version: event.version, modifiedAt: event.createdAt, modifiedById: event.createdBy, }, { id: event.aggregateId }, ); return { kind: "applied", verb, row: rows[0] ?? null }; } case "deleted": { if (softDelete) { const rows = await updateMany( tx, table, { isDeleted: true, deletedAt: event.createdAt, deletedById: event.createdBy, version: event.version, modifiedAt: event.createdAt, modifiedById: event.createdBy, }, { id: event.aggregateId }, ); return { kind: "applied", verb, row: rows[0] ?? null }; } // Hard-Delete: DELETE-Statement gibt keine returning-Row her und // der Live-Pfad nutzt eh `existing` (pre-delete-Snapshot) für die // Response. Beim Replay ist das fine, der Caller braucht die Row // nicht weiter. await deleteMany(tx, table, { id: event.aggregateId }); return { kind: "applied", verb, row: null }; } case "restored": { // Restore ist nur bei softDelete sinnvoll. Hard-Delete-Entities sollten // keine restored-Events erhalten — falls doch, defensive skip. if (!softDelete) return { kind: "skipped" }; const rows = await updateMany( tx, table, { isDeleted: false, deletedAt: null, deletedById: null, version: event.version, modifiedAt: event.createdAt, modifiedById: event.createdBy, }, { id: event.aggregateId }, ); return { kind: "applied", verb, row: rows[0] ?? null }; } case "forgotten": { // Hard-delete regardless of softDelete: forget/purge (Art. 17) removes the // row entirely. On rebuild the aggregate replays created → forgotten, so // the row ends up gone — the rebuild-safe erasure the soft-delete verb // (which only flips isDeleted) cannot provide. await deleteMany(tx, table, { id: event.aggregateId }); return { kind: "applied", verb, row: null }; } } }