275 lines
9.9 KiB
TypeScript
275 lines
9.9 KiB
TypeScript
import { Inject, Injectable } from '@nestjs/common';
|
|
import {
|
|
and,
|
|
asc,
|
|
companies,
|
|
eq,
|
|
estates,
|
|
hierarchyAuditEvents,
|
|
hierarchyOutbox,
|
|
platformProjects,
|
|
type Db,
|
|
type HIERARCHY_AUDIT_TARGET_KINDS,
|
|
type HIERARCHY_AUDIT_VERBS,
|
|
} from '@mosaicstack/db';
|
|
import { DB } from '../database/database.module.js';
|
|
|
|
/**
|
|
* Hierarchy audit event + outbox machinery (contract 1 §5.2).
|
|
*
|
|
* Every hierarchy mutation writes its semantic audit event AND the event's
|
|
* outbox record on the caller's transaction, so state, event, and outbox
|
|
* commit or roll back together. Events reference their target by an
|
|
* immutable snapshot (id, slug, parent chain at event time), never by a
|
|
* foreign key into the class tables — append-only events survive the
|
|
* deletion of their target. This module exposes no update or delete path
|
|
* for events: append-only is a property of the code surface, witnessed by
|
|
* the integration tests.
|
|
*
|
|
* This is NOT a class-table writer: it touches only the audit/outbox
|
|
* tables, so it does not appear on the writer-coverage allowlist. The
|
|
* hierarchy command repositories (M4-1b-ii) are the allowlisted writers and
|
|
* call into this on their own transactions.
|
|
*/
|
|
|
|
export type HierarchyAuditVerb = (typeof HIERARCHY_AUDIT_VERBS)[number];
|
|
export type HierarchyTargetKind = (typeof HIERARCHY_AUDIT_TARGET_KINDS)[number];
|
|
export type HierarchyNodeKind = Exclude<HierarchyTargetKind, 'grant'>;
|
|
|
|
export interface ParentChainEntry {
|
|
readonly kind: HierarchyNodeKind;
|
|
readonly id: string;
|
|
readonly slug: string;
|
|
}
|
|
|
|
/** Immutable node snapshot at event time; parentChain is root-first. */
|
|
export interface HierarchyNodeSnapshot {
|
|
readonly id: string;
|
|
readonly slug: string;
|
|
readonly name: string;
|
|
readonly parentChain: readonly ParentChainEntry[];
|
|
}
|
|
|
|
export interface AppendHierarchyEventInput {
|
|
readonly actorId: string;
|
|
readonly verb: HierarchyAuditVerb;
|
|
readonly targetKind: HierarchyTargetKind;
|
|
readonly targetId: string;
|
|
/** Node events: HierarchyNodeSnapshot. Grant events: subject/target/role snapshot (contract 2 §4.4). */
|
|
readonly targetSnapshot: Record<string, unknown>;
|
|
/** Present exactly on transfers (CHECK-enforced): source/destination parent { kind, id, slug }. */
|
|
readonly transferFrom?: ParentChainEntry;
|
|
readonly transferTo?: ParentChainEntry;
|
|
readonly correlationId: string;
|
|
/** Prior event in the causal chain (e.g. the delete event causing cascaded grant_revoke events). */
|
|
readonly causationId?: string;
|
|
readonly idempotencyKey: string;
|
|
}
|
|
|
|
export type HierarchyAuditEventRow = typeof hierarchyAuditEvents.$inferSelect;
|
|
export type HierarchyOutboxRow = typeof hierarchyOutbox.$inferSelect;
|
|
|
|
export interface AppendHierarchyEventResult {
|
|
readonly event: HierarchyAuditEventRow;
|
|
/** True when the idempotency key had already committed an identical event (REQ-AUD-001 duplicate suppression). */
|
|
readonly replayed: boolean;
|
|
}
|
|
|
|
type Tx = Pick<Db, 'insert' | 'select'>;
|
|
|
|
export class HierarchyAuditIdempotencyConflictError extends Error {
|
|
constructor(idempotencyKey: string) {
|
|
super(
|
|
`hierarchy audit idempotency key ${idempotencyKey} already exists with different event content`,
|
|
);
|
|
this.name = 'HierarchyAuditIdempotencyConflictError';
|
|
}
|
|
}
|
|
|
|
export class HierarchyNodeNotFoundError extends Error {
|
|
constructor(kind: HierarchyNodeKind, id: string) {
|
|
super(`hierarchy node not found: ${kind} ${id}`);
|
|
this.name = 'HierarchyNodeNotFoundError';
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Append one audit event and its outbox record on the caller's transaction.
|
|
* A duplicate idempotency key with identical semantic content returns the
|
|
* prior event (replayed: true) without inserting anything; a duplicate key
|
|
* with different content throws.
|
|
*/
|
|
export async function appendHierarchyEvent(
|
|
tx: Tx,
|
|
input: AppendHierarchyEventInput,
|
|
): Promise<AppendHierarchyEventResult> {
|
|
const inserted = await tx
|
|
.insert(hierarchyAuditEvents)
|
|
.values({
|
|
actorId: input.actorId,
|
|
verb: input.verb,
|
|
targetKind: input.targetKind,
|
|
targetId: input.targetId,
|
|
targetSnapshot: input.targetSnapshot,
|
|
transferFrom: input.transferFrom ?? null,
|
|
transferTo: input.transferTo ?? null,
|
|
correlationId: input.correlationId,
|
|
causationId: input.causationId ?? null,
|
|
idempotencyKey: input.idempotencyKey,
|
|
})
|
|
.onConflictDoNothing()
|
|
.returning();
|
|
const event = inserted[0];
|
|
if (event) {
|
|
await tx.insert(hierarchyOutbox).values({
|
|
eventId: event.id,
|
|
idempotencyKey: input.idempotencyKey,
|
|
correlationId: input.correlationId,
|
|
});
|
|
return { event, replayed: false };
|
|
}
|
|
|
|
const prior = await tx
|
|
.select()
|
|
.from(hierarchyAuditEvents)
|
|
.where(eq(hierarchyAuditEvents.idempotencyKey, input.idempotencyKey))
|
|
.limit(1);
|
|
const existing = prior[0];
|
|
if (!existing || !sameEvent(existing, input)) {
|
|
throw new HierarchyAuditIdempotencyConflictError(input.idempotencyKey);
|
|
}
|
|
// Event and outbox committed atomically the first time, so the outbox
|
|
// record already exists; a replay inserts nothing.
|
|
return { event: existing, replayed: true };
|
|
}
|
|
|
|
/** Key-order-independent serialization: jsonb does not preserve key order. */
|
|
function canonicalJson(value: unknown): string {
|
|
if (Array.isArray(value)) return `[${value.map(canonicalJson).join(',')}]`;
|
|
if (value !== null && typeof value === 'object') {
|
|
const record = value as Record<string, unknown>;
|
|
const body = Object.keys(record)
|
|
.sort()
|
|
.map((key) => `${JSON.stringify(key)}:${canonicalJson(record[key])}`)
|
|
.join(',');
|
|
return `{${body}}`;
|
|
}
|
|
return JSON.stringify(value);
|
|
}
|
|
|
|
function sameEvent(row: HierarchyAuditEventRow, input: AppendHierarchyEventInput): boolean {
|
|
return (
|
|
row.actorId === input.actorId &&
|
|
row.verb === input.verb &&
|
|
row.targetKind === input.targetKind &&
|
|
row.targetId === input.targetId &&
|
|
row.correlationId === input.correlationId &&
|
|
(row.causationId ?? null) === (input.causationId ?? null) &&
|
|
canonicalJson(row.targetSnapshot) === canonicalJson(input.targetSnapshot) &&
|
|
// Transfer source/destination are semantic content (§5.2): a retry with a
|
|
// different destination must conflict, never silently replay.
|
|
canonicalJson(row.transferFrom ?? null) === canonicalJson(input.transferFrom ?? null) &&
|
|
canonicalJson(row.transferTo ?? null) === canonicalJson(input.transferTo ?? null)
|
|
);
|
|
}
|
|
|
|
/**
|
|
* Build the immutable snapshot for a node: its row plus the parent chain up
|
|
* to the company root, root-first, read on the caller's transaction so the
|
|
* snapshot is consistent with the mutation it audits.
|
|
*/
|
|
export async function buildNodeSnapshot(
|
|
tx: Tx,
|
|
kind: HierarchyNodeKind,
|
|
id: string,
|
|
): Promise<HierarchyNodeSnapshot> {
|
|
if (kind === 'company') {
|
|
const rows = await tx.select().from(companies).where(eq(companies.id, id)).limit(1);
|
|
const row = rows[0];
|
|
if (!row) throw new HierarchyNodeNotFoundError(kind, id);
|
|
return { id: row.id, slug: row.slug, name: row.name, parentChain: [] };
|
|
}
|
|
if (kind === 'estate') {
|
|
const rows = await tx.select().from(estates).where(eq(estates.id, id)).limit(1);
|
|
const row = rows[0];
|
|
if (!row) throw new HierarchyNodeNotFoundError(kind, id);
|
|
const parent = await buildNodeSnapshot(tx, 'company', row.companyId);
|
|
return {
|
|
id: row.id,
|
|
slug: row.slug,
|
|
name: row.name,
|
|
parentChain: [...parent.parentChain, { kind: 'company', id: parent.id, slug: parent.slug }],
|
|
};
|
|
}
|
|
const rows = await tx.select().from(platformProjects).where(eq(platformProjects.id, id)).limit(1);
|
|
const row = rows[0];
|
|
if (!row) throw new HierarchyNodeNotFoundError(kind, id);
|
|
const parent = await buildNodeSnapshot(tx, 'estate', row.estateId);
|
|
return {
|
|
id: row.id,
|
|
slug: row.slug,
|
|
name: row.name,
|
|
parentChain: [...parent.parentChain, { kind: 'estate', id: parent.id, slug: parent.slug }],
|
|
};
|
|
}
|
|
|
|
@Injectable()
|
|
export class HierarchyAuditRepository {
|
|
constructor(@Inject(DB) private readonly db: Db) {}
|
|
|
|
/** Compose an event+outbox append into a caller-owned transaction. */
|
|
append(tx: Tx, input: AppendHierarchyEventInput): Promise<AppendHierarchyEventResult> {
|
|
return appendHierarchyEvent(tx, input);
|
|
}
|
|
|
|
snapshot(tx: Tx, kind: HierarchyNodeKind, id: string): Promise<HierarchyNodeSnapshot> {
|
|
return buildNodeSnapshot(tx, kind, id);
|
|
}
|
|
|
|
/** Per-target ordered event history (REQ-AUD-001 per-target ordering; read-only). */
|
|
async eventsForTarget(targetId: string): Promise<HierarchyAuditEventRow[]> {
|
|
return this.db
|
|
.select()
|
|
.from(hierarchyAuditEvents)
|
|
.where(eq(hierarchyAuditEvents.targetId, targetId))
|
|
.orderBy(asc(hierarchyAuditEvents.seq));
|
|
}
|
|
|
|
/**
|
|
* Claim the oldest pending outbox record (claim-by-CAS: the UPDATE is
|
|
* guarded on status so a lost race returns null and the caller retries).
|
|
*/
|
|
async claimPendingOutbox(): Promise<HierarchyOutboxRow | null> {
|
|
const candidates = await this.db
|
|
.select()
|
|
.from(hierarchyOutbox)
|
|
.where(eq(hierarchyOutbox.status, 'pending'))
|
|
.orderBy(asc(hierarchyOutbox.createdAt))
|
|
.limit(1);
|
|
const candidate = candidates[0];
|
|
if (!candidate) return null;
|
|
const claimed = await this.db
|
|
.update(hierarchyOutbox)
|
|
.set({ status: 'processing', updatedAt: new Date() })
|
|
.where(and(eq(hierarchyOutbox.id, candidate.id), eq(hierarchyOutbox.status, 'pending')))
|
|
.returning();
|
|
return claimed[0] ?? null;
|
|
}
|
|
|
|
async completeOutbox(id: string): Promise<void> {
|
|
const now = new Date();
|
|
await this.db
|
|
.update(hierarchyOutbox)
|
|
.set({ status: 'delivered', deliveredAt: now, updatedAt: now })
|
|
.where(and(eq(hierarchyOutbox.id, id), eq(hierarchyOutbox.status, 'processing')));
|
|
}
|
|
|
|
/** Return a claimed record to pending (delivery failed; it stays replayable). */
|
|
async releaseOutbox(id: string): Promise<void> {
|
|
await this.db
|
|
.update(hierarchyOutbox)
|
|
.set({ status: 'pending', updatedAt: new Date() })
|
|
.where(and(eq(hierarchyOutbox.id, id), eq(hierarchyOutbox.status, 'processing')));
|
|
}
|
|
}
|