ceo-api/src/events/outbox.service.ts
valentinbvro eb652446f6 feat: Sprint 1 -- Projection Kernel (read models pentru dashboard)
Repara criteriul de acceptare "dashboardurile citesc read models, nu
interogheaza haotic toate modulele".

- projection_processed_events (unique: name+version+event_id) = garantia de
  idempotency. Checkpointul e doar optimizare de scanare, nu corectitudine:
  rescanam cu un safety lag de 60s si sarim ce s-a aplicat deja, ca sa nu
  pierdem evenimente comise dupa unul cu created_at mai mare.
- ProjectionRegistry cu ProjectionDefinition (name, version, subscribedEvents,
  rebuildStrategy, apply, rebuildTenant).
- executive_dashboard foloseste rebuildStrategy 'canonical-tables':
  recalculeaza contorii din tabelele canonice, nu incrementeaza. Contorii
  incrementali pot deriva daca un eveniment se pierde sau se dubleaza;
  recalcularea e corecta prin constructie si idempotenta natural. Costul e o
  interogare per eveniment relevant -- ok la volumul actual.
- GET /v1/dashboard/executive citeste proiectia. Cand proiectia inca n-a rulat
  pentru workspace, intoarce status 'building' explicit, NU zerouri care ar
  parea date reale.
- POST /v1/dashboard/executive/rebuild -- owner/admin, auditat.
- Serviciile de taskuri/organizatii/segmente emit acum evenimente in outbox
  (in aceeasi tranzactie cu scrierea). Fara ele proiectia era cod mort.
- Campurile din spec care depind de module neconstruite (documents,
  transactions, approvals) raman 0 explicit, nu inventate.
2026-07-29 12:03:34 +02:00

56 lines
1.9 KiB
TypeScript

import { Injectable } from '@nestjs/common';
import type { NodePgDatabase } from 'drizzle-orm/node-postgres';
import { db } from '../db/client';
import { outboxEvents } from '../db/schema';
import * as schema from '../db/schema';
export interface OutboxEventInput {
tenantId: string;
/** Fara workspace, proiectiile nu stiu in ce partitie sa scrie si sar evenimentul. */
workspaceId?: string;
eventType: string;
subjectId?: string;
aggregateType?: string;
actorId?: string;
payload: Record<string, unknown>;
correlationId?: string;
/** Ce comanda/eveniment a produs acest eveniment. */
causationId?: string;
/** C0-C4; controleaza catre ce servicii poate fi distribuit. */
classification?: string;
provenance?: Record<string, unknown>;
eventVersion?: number;
}
type Transaction = Parameters<Parameters<NodePgDatabase<typeof schema>['transaction']>[0]>[0];
/**
* Blueprint 9.1 / principiul 2.1 "events before intelligence": orice scriere de
* domeniu care trebuie sa produca un eveniment scrie in aceeasi tranzactie in
* outbox_events. Publicarea efectiva (dispatch) se face separat, async, de catre
* OutboxDispatcher -- niciodata inline in request path.
*/
@Injectable()
export class OutboxService {
async withTransaction<T>(fn: (tx: Transaction) => Promise<T>): Promise<T> {
return db.transaction(fn);
}
async record(tx: Transaction, event: OutboxEventInput): Promise<void> {
await tx.insert(outboxEvents).values({
tenantId: event.tenantId,
workspaceId: event.workspaceId,
eventType: event.eventType,
eventVersion: event.eventVersion ?? 1,
occurredAt: new Date(),
actorId: event.actorId,
aggregateType: event.aggregateType,
subjectId: event.subjectId,
payload: event.payload,
correlationId: event.correlationId,
causationId: event.causationId,
classification: event.classification ?? 'c2',
provenance: event.provenance ?? {},
});
}
}