feat(hierarchy): audit event + outbox machinery (M4-1b-i, contract 1 §5.2) #1460
@@ -0,0 +1,106 @@
|
|||||||
|
import { RequestMethod, type Type } from '@nestjs/common';
|
||||||
|
import { describe, expect, it } from 'vitest';
|
||||||
|
import { AppModule } from '../app.module.js';
|
||||||
|
import { HierarchyModule } from '../hierarchy/hierarchy.module.js';
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Hierarchy route-inventory baseline (contract 1 §6.3(a)).
|
||||||
|
*
|
||||||
|
* M4-1b-i ships the audit event + outbox machinery with NO mutation routes:
|
||||||
|
* the hierarchy command family (controllers + DTOs) lands in M4-1b-ii once
|
||||||
|
* contract 2 merges. This witness enumerates every route the AppModule graph
|
||||||
|
* declares and pins that baseline, so a hierarchy route appearing before its
|
||||||
|
* command-family witnesses exist fails here first. When M4-1b-ii lands, this
|
||||||
|
* baseline is replaced by an exact inventory of the command family.
|
||||||
|
*/
|
||||||
|
|
||||||
|
interface RouteEntry {
|
||||||
|
method: string;
|
||||||
|
path: string;
|
||||||
|
controller: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Module-metadata entry: a module class or a DynamicModule-shaped object. */
|
||||||
|
type ModuleEntry =
|
||||||
|
| Type<unknown>
|
||||||
|
| { module: Type<unknown>; imports?: unknown[]; controllers?: Type<unknown>[] };
|
||||||
|
|
||||||
|
function collectControllers(root: ModuleEntry): Type<unknown>[] {
|
||||||
|
const visited = new Set<unknown>();
|
||||||
|
const controllers: Type<unknown>[] = [];
|
||||||
|
const walk = (entry: ModuleEntry | undefined | null): void => {
|
||||||
|
if (!entry || visited.has(entry)) return;
|
||||||
|
visited.add(entry);
|
||||||
|
const moduleClass = typeof entry === 'function' ? entry : entry.module;
|
||||||
|
// Entries with no resolvable class (forwardRef wrappers, async dynamic
|
||||||
|
// modules) carry no decorator metadata to read here.
|
||||||
|
if (typeof moduleClass !== 'function') return;
|
||||||
|
if (visited.has(moduleClass) && typeof entry !== 'function') return;
|
||||||
|
visited.add(moduleClass);
|
||||||
|
// 'controllers' / 'imports' are the metadata keys the @Module decorator writes.
|
||||||
|
const declared = (Reflect.getMetadata('controllers', moduleClass) ?? []) as Type<unknown>[];
|
||||||
|
controllers.push(...declared);
|
||||||
|
if (typeof entry !== 'function' && entry.controllers) controllers.push(...entry.controllers);
|
||||||
|
const imports = [
|
||||||
|
...((Reflect.getMetadata('imports', moduleClass) ?? []) as ModuleEntry[]),
|
||||||
|
...(typeof entry !== 'function' ? ((entry.imports ?? []) as ModuleEntry[]) : []),
|
||||||
|
];
|
||||||
|
for (const imported of imports) walk(imported);
|
||||||
|
};
|
||||||
|
walk(root);
|
||||||
|
return controllers;
|
||||||
|
}
|
||||||
|
|
||||||
|
function routesOf(controller: Type<unknown>): RouteEntry[] {
|
||||||
|
// 'path' on the class is the @Controller prefix; 'path'/'method' on a
|
||||||
|
// handler are written by the @Get/@Post/... route decorators.
|
||||||
|
const base = (Reflect.getMetadata('path', controller) ?? '') as string | string[];
|
||||||
|
const bases = Array.isArray(base) ? base : [base];
|
||||||
|
const routes: RouteEntry[] = [];
|
||||||
|
const prototype = controller.prototype as Record<string, unknown>;
|
||||||
|
for (const name of Object.getOwnPropertyNames(prototype)) {
|
||||||
|
if (name === 'constructor') continue;
|
||||||
|
const handler = Object.getOwnPropertyDescriptor(prototype, name)?.value;
|
||||||
|
if (typeof handler !== 'function') continue;
|
||||||
|
const method = Reflect.getMetadata('method', handler) as number | undefined;
|
||||||
|
if (method === undefined) continue;
|
||||||
|
const sub = (Reflect.getMetadata('path', handler) ?? '/') as string;
|
||||||
|
for (const prefix of bases) {
|
||||||
|
const path = `/${prefix}/${sub}`.replace(/\/+/g, '/').replace(/(.)\/$/, '$1');
|
||||||
|
routes.push({
|
||||||
|
method: RequestMethod[method] ?? String(method),
|
||||||
|
path,
|
||||||
|
controller: controller.name,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return routes;
|
||||||
|
}
|
||||||
|
|
||||||
|
describe('hierarchy route-inventory baseline (§6.3(a))', () => {
|
||||||
|
const inventory = collectControllers(AppModule).flatMap(routesOf);
|
||||||
|
|
||||||
|
it('control: the enumeration sees the known route surface', () => {
|
||||||
|
const paths = inventory.map((r) => `${r.method} ${r.path}`);
|
||||||
|
expect(paths).toContain('GET /health');
|
||||||
|
expect(paths).toContain('POST /api/workspaces');
|
||||||
|
expect(paths).toContain('GET /api/teams');
|
||||||
|
expect(inventory.length).toBeGreaterThan(20);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('declares zero hierarchy mutation routes before M4-1b-ii', () => {
|
||||||
|
const hierarchyRoutes = inventory.filter((r) =>
|
||||||
|
/hierarch|compan|estate|platform[-_]?project/i.test(r.path),
|
||||||
|
);
|
||||||
|
expect(
|
||||||
|
hierarchyRoutes,
|
||||||
|
'a hierarchy route landed without replacing the §6.3(a) baseline with a command-family inventory',
|
||||||
|
).toEqual([]);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('HierarchyModule itself declares no controllers', () => {
|
||||||
|
expect((Reflect.getMetadata('controllers', HierarchyModule) ?? []) as unknown[]).toEqual([]);
|
||||||
|
const hierarchyControllers = collectControllers(HierarchyModule);
|
||||||
|
expect(hierarchyControllers).toEqual([]);
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -24,6 +24,7 @@ import { GCModule } from './gc/gc.module.js';
|
|||||||
import { HarnessModule } from './harness/harness.module.js';
|
import { HarnessModule } from './harness/harness.module.js';
|
||||||
import { ReloadModule } from './reload/reload.module.js';
|
import { ReloadModule } from './reload/reload.module.js';
|
||||||
import { WorkspaceModule } from './workspace/workspace.module.js';
|
import { WorkspaceModule } from './workspace/workspace.module.js';
|
||||||
|
import { HierarchyModule } from './hierarchy/hierarchy.module.js';
|
||||||
import { QueueModule } from './queue/queue.module.js';
|
import { QueueModule } from './queue/queue.module.js';
|
||||||
import { FederationModule } from './federation/federation.module.js';
|
import { FederationModule } from './federation/federation.module.js';
|
||||||
import { ThrottlerGuard, ThrottlerModule } from '@nestjs/throttler';
|
import { ThrottlerGuard, ThrottlerModule } from '@nestjs/throttler';
|
||||||
@@ -65,6 +66,7 @@ const federationEnabled = loadConfig(resolveGatewayConfigPath()).tier === 'feder
|
|||||||
QueueModule,
|
QueueModule,
|
||||||
ReloadModule,
|
ReloadModule,
|
||||||
WorkspaceModule,
|
WorkspaceModule,
|
||||||
|
HierarchyModule,
|
||||||
...(federationEnabled ? [FederationModule] : []),
|
...(federationEnabled ? [FederationModule] : []),
|
||||||
],
|
],
|
||||||
controllers: [HealthController],
|
controllers: [HealthController],
|
||||||
|
|||||||
@@ -0,0 +1,247 @@
|
|||||||
|
import { mkdtemp, rm } from 'node:fs/promises';
|
||||||
|
import { randomUUID } from 'node:crypto';
|
||||||
|
import { tmpdir } from 'node:os';
|
||||||
|
import { join } from 'node:path';
|
||||||
|
import { afterAll, beforeAll, describe, expect, it } from 'vitest';
|
||||||
|
import { Test, type TestingModule } from '@nestjs/testing';
|
||||||
|
import {
|
||||||
|
companies,
|
||||||
|
createPgliteDb,
|
||||||
|
eq,
|
||||||
|
estates,
|
||||||
|
hierarchyAuditEvents,
|
||||||
|
hierarchyOutbox,
|
||||||
|
platformProjects,
|
||||||
|
runPgliteMigrations,
|
||||||
|
type DbHandle,
|
||||||
|
} from '@mosaicstack/db';
|
||||||
|
import { DB } from '../database/database.module.js';
|
||||||
|
import {
|
||||||
|
HierarchyAuditIdempotencyConflictError,
|
||||||
|
HierarchyAuditRepository,
|
||||||
|
HierarchyNodeNotFoundError,
|
||||||
|
type AppendHierarchyEventInput,
|
||||||
|
} from './hierarchy-audit.repository.js';
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Repository-level §6.4 witnesses for the hierarchy audit machinery
|
||||||
|
* (contract 1 §5.2, REQ-AUD-001): same-transaction atomicity of state +
|
||||||
|
* event + outbox, rollback leaving no residue, idempotent replay, snapshot
|
||||||
|
* parent chains, events surviving target deletion, per-target ordering, and
|
||||||
|
* the outbox claim/complete/release CAS. The schema-level constraints are
|
||||||
|
* witnessed in packages/db/src/hierarchy-audit.witness.test.ts.
|
||||||
|
*/
|
||||||
|
describe('hierarchy audit repository integration', (): void => {
|
||||||
|
let dataDir: string;
|
||||||
|
let handle: DbHandle;
|
||||||
|
let moduleRef: TestingModule;
|
||||||
|
let repo: HierarchyAuditRepository;
|
||||||
|
|
||||||
|
const input = (
|
||||||
|
overrides: Partial<AppendHierarchyEventInput> = {},
|
||||||
|
): AppendHierarchyEventInput => ({
|
||||||
|
actorId: 'user-actor',
|
||||||
|
verb: 'create',
|
||||||
|
targetKind: 'company',
|
||||||
|
targetId: randomUUID(),
|
||||||
|
targetSnapshot: { id: 'x', slug: 'x', name: 'x', parentChain: [] },
|
||||||
|
correlationId: 'corr-1',
|
||||||
|
idempotencyKey: `key-${randomUUID()}`,
|
||||||
|
...overrides,
|
||||||
|
});
|
||||||
|
|
||||||
|
beforeAll(async (): Promise<void> => {
|
||||||
|
dataDir = await mkdtemp(join(tmpdir(), 'mosaic-gateway-hierarchy-audit-'));
|
||||||
|
handle = createPgliteDb(dataDir);
|
||||||
|
await runPgliteMigrations(handle);
|
||||||
|
moduleRef = await Test.createTestingModule({
|
||||||
|
providers: [HierarchyAuditRepository, { provide: DB, useValue: handle.db }],
|
||||||
|
}).compile();
|
||||||
|
repo = moduleRef.get(HierarchyAuditRepository);
|
||||||
|
});
|
||||||
|
|
||||||
|
afterAll(async (): Promise<void> => {
|
||||||
|
await moduleRef.close();
|
||||||
|
await handle.close();
|
||||||
|
await rm(dataDir, { recursive: true, force: true });
|
||||||
|
});
|
||||||
|
|
||||||
|
it('commits state, event, and outbox record atomically in one transaction', async () => {
|
||||||
|
const companyId = randomUUID();
|
||||||
|
const key = `key-${randomUUID()}`;
|
||||||
|
await handle.db.transaction(async (tx) => {
|
||||||
|
await tx.insert(companies).values({ id: companyId, name: 'Atomic Co', slug: 'atomic-co' });
|
||||||
|
const snapshot = await repo.snapshot(tx, 'company', companyId);
|
||||||
|
const result = await repo.append(tx, {
|
||||||
|
...input({ targetId: companyId, idempotencyKey: key }),
|
||||||
|
targetSnapshot: { ...snapshot },
|
||||||
|
});
|
||||||
|
expect(result.replayed).toBe(false);
|
||||||
|
expect(result.event.idempotencyKey).toBe(key);
|
||||||
|
});
|
||||||
|
const events = await handle.db
|
||||||
|
.select()
|
||||||
|
.from(hierarchyAuditEvents)
|
||||||
|
.where(eq(hierarchyAuditEvents.idempotencyKey, key));
|
||||||
|
expect(events).toHaveLength(1);
|
||||||
|
const outbox = await handle.db
|
||||||
|
.select()
|
||||||
|
.from(hierarchyOutbox)
|
||||||
|
.where(eq(hierarchyOutbox.eventId, events[0]!.id));
|
||||||
|
expect(outbox).toHaveLength(1);
|
||||||
|
expect(outbox[0]).toMatchObject({
|
||||||
|
status: 'pending',
|
||||||
|
idempotencyKey: key,
|
||||||
|
correlationId: 'corr-1',
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
it('a rolled-back transaction leaves no state, no event, and no outbox record', async () => {
|
||||||
|
const companyId = randomUUID();
|
||||||
|
const key = `key-${randomUUID()}`;
|
||||||
|
await expect(
|
||||||
|
handle.db.transaction(async (tx) => {
|
||||||
|
await tx.insert(companies).values({ id: companyId, name: 'Doomed Co', slug: 'doomed-co' });
|
||||||
|
await repo.append(tx, input({ targetId: companyId, idempotencyKey: key }));
|
||||||
|
throw new Error('deliberate rollback');
|
||||||
|
}),
|
||||||
|
).rejects.toThrow('deliberate rollback');
|
||||||
|
const [companyRows, eventRows, outboxRows] = await Promise.all([
|
||||||
|
handle.db.select().from(companies).where(eq(companies.id, companyId)),
|
||||||
|
handle.db
|
||||||
|
.select()
|
||||||
|
.from(hierarchyAuditEvents)
|
||||||
|
.where(eq(hierarchyAuditEvents.idempotencyKey, key)),
|
||||||
|
handle.db.select().from(hierarchyOutbox).where(eq(hierarchyOutbox.idempotencyKey, key)),
|
||||||
|
]);
|
||||||
|
expect(companyRows).toHaveLength(0);
|
||||||
|
expect(eventRows).toHaveLength(0);
|
||||||
|
expect(outboxRows).toHaveLength(0);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('replays a duplicate idempotency key without inserting a second event or outbox record', async () => {
|
||||||
|
const first = input();
|
||||||
|
const original = await handle.db.transaction(async (tx) => repo.append(tx, first));
|
||||||
|
const replay = await handle.db.transaction(async (tx) => repo.append(tx, first));
|
||||||
|
expect(original.replayed).toBe(false);
|
||||||
|
expect(replay.replayed).toBe(true);
|
||||||
|
expect(replay.event.id).toBe(original.event.id);
|
||||||
|
const outbox = await handle.db
|
||||||
|
.select()
|
||||||
|
.from(hierarchyOutbox)
|
||||||
|
.where(eq(hierarchyOutbox.eventId, original.event.id));
|
||||||
|
expect(outbox).toHaveLength(1);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('throws on a duplicate idempotency key carrying different event content', async () => {
|
||||||
|
const first = input();
|
||||||
|
await handle.db.transaction(async (tx) => repo.append(tx, first));
|
||||||
|
await expect(
|
||||||
|
handle.db.transaction(async (tx) =>
|
||||||
|
repo.append(tx, { ...first, verb: 'rename', targetId: randomUUID() }),
|
||||||
|
),
|
||||||
|
).rejects.toThrow(HierarchyAuditIdempotencyConflictError);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('throws on a duplicate idempotency key whose transfer destination differs', async () => {
|
||||||
|
const from = { kind: 'company' as const, id: randomUUID(), slug: 'src-co' };
|
||||||
|
const to = { kind: 'company' as const, id: randomUUID(), slug: 'dst-co' };
|
||||||
|
const first = input({
|
||||||
|
verb: 'transfer',
|
||||||
|
targetKind: 'estate',
|
||||||
|
transferFrom: from,
|
||||||
|
transferTo: to,
|
||||||
|
});
|
||||||
|
const original = await handle.db.transaction(async (tx) => repo.append(tx, first));
|
||||||
|
expect(original.replayed).toBe(false);
|
||||||
|
// Identical retry replays; a retry re-routed to a different destination must conflict.
|
||||||
|
const replay = await handle.db.transaction(async (tx) => repo.append(tx, first));
|
||||||
|
expect(replay.replayed).toBe(true);
|
||||||
|
await expect(
|
||||||
|
handle.db.transaction(async (tx) =>
|
||||||
|
repo.append(tx, { ...first, transferTo: { ...to, id: randomUUID() } }),
|
||||||
|
),
|
||||||
|
).rejects.toThrow(HierarchyAuditIdempotencyConflictError);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('builds root-first parent chains and rejects unknown nodes', async () => {
|
||||||
|
const companyId = randomUUID();
|
||||||
|
const estateId = randomUUID();
|
||||||
|
const projectId = randomUUID();
|
||||||
|
await handle.db.transaction(async (tx) => {
|
||||||
|
await tx.insert(companies).values({ id: companyId, name: 'Chain Co', slug: 'chain-co' });
|
||||||
|
await tx
|
||||||
|
.insert(estates)
|
||||||
|
.values({ id: estateId, name: 'Chain Estate', slug: 'chain-estate', companyId });
|
||||||
|
await tx
|
||||||
|
.insert(platformProjects)
|
||||||
|
.values({ id: projectId, name: 'Chain Project', slug: 'chain-project', estateId });
|
||||||
|
});
|
||||||
|
const snapshot = await repo.snapshot(handle.db, 'platform_project', projectId);
|
||||||
|
expect(snapshot).toMatchObject({ id: projectId, slug: 'chain-project', name: 'Chain Project' });
|
||||||
|
expect(snapshot.parentChain).toEqual([
|
||||||
|
{ kind: 'company', id: companyId, slug: 'chain-co' },
|
||||||
|
{ kind: 'estate', id: estateId, slug: 'chain-estate' },
|
||||||
|
]);
|
||||||
|
await expect(repo.snapshot(handle.db, 'estate', randomUUID())).rejects.toThrow(
|
||||||
|
HierarchyNodeNotFoundError,
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('keeps events readable, in per-target seq order, after the target row is deleted', async () => {
|
||||||
|
const companyId = randomUUID();
|
||||||
|
await handle.db.transaction(async (tx) => {
|
||||||
|
await tx.insert(companies).values({ id: companyId, name: 'Mortal Co', slug: 'mortal-co' });
|
||||||
|
const snapshot = await repo.snapshot(tx, 'company', companyId);
|
||||||
|
await repo.append(tx, input({ targetId: companyId, targetSnapshot: { ...snapshot } }));
|
||||||
|
});
|
||||||
|
await handle.db.transaction(async (tx) => {
|
||||||
|
const snapshot = await repo.snapshot(tx, 'company', companyId);
|
||||||
|
await repo.append(tx, {
|
||||||
|
...input({ verb: 'delete', targetId: companyId }),
|
||||||
|
targetSnapshot: { ...snapshot },
|
||||||
|
});
|
||||||
|
await tx.delete(companies).where(eq(companies.id, companyId));
|
||||||
|
});
|
||||||
|
const events = await repo.eventsForTarget(companyId);
|
||||||
|
expect(events.map((e) => e.verb)).toEqual(['create', 'delete']);
|
||||||
|
expect(events[1]!.seq).toBeGreaterThan(events[0]!.seq);
|
||||||
|
expect((events[1]!.targetSnapshot as { id: string }).id).toBe(companyId);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('claims the oldest pending outbox record exactly once, completes and releases by CAS', async () => {
|
||||||
|
// Drain records left pending by earlier cases so ordering is deterministic.
|
||||||
|
for (;;) {
|
||||||
|
const drained = await repo.claimPendingOutbox();
|
||||||
|
if (!drained) break;
|
||||||
|
await repo.completeOutbox(drained.id);
|
||||||
|
}
|
||||||
|
const older = await handle.db.transaction(async (tx) => repo.append(tx, input()));
|
||||||
|
const newer = await handle.db.transaction(async (tx) => repo.append(tx, input()));
|
||||||
|
|
||||||
|
const claimed = await repo.claimPendingOutbox();
|
||||||
|
expect(claimed).not.toBeNull();
|
||||||
|
expect(claimed!.eventId).toBe(older.event.id);
|
||||||
|
expect(claimed!.status).toBe('processing');
|
||||||
|
|
||||||
|
// Delivery fails: release returns it to pending and it is claimable again.
|
||||||
|
await repo.releaseOutbox(claimed!.id);
|
||||||
|
const reclaimed = await repo.claimPendingOutbox();
|
||||||
|
expect(reclaimed!.id).toBe(claimed!.id);
|
||||||
|
|
||||||
|
await repo.completeOutbox(reclaimed!.id);
|
||||||
|
const done = await handle.db
|
||||||
|
.select()
|
||||||
|
.from(hierarchyOutbox)
|
||||||
|
.where(eq(hierarchyOutbox.id, reclaimed!.id));
|
||||||
|
expect(done[0]!.status).toBe('delivered');
|
||||||
|
expect(done[0]!.deliveredAt).not.toBeNull();
|
||||||
|
// completeOutbox is CAS-guarded on 'processing': completing again is a no-op.
|
||||||
|
await repo.completeOutbox(reclaimed!.id);
|
||||||
|
|
||||||
|
const second = await repo.claimPendingOutbox();
|
||||||
|
expect(second!.eventId).toBe(newer.event.id);
|
||||||
|
await repo.completeOutbox(second!.id);
|
||||||
|
expect(await repo.claimPendingOutbox()).toBeNull();
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -0,0 +1,274 @@
|
|||||||
|
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')));
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,17 @@
|
|||||||
|
import { Module } from '@nestjs/common';
|
||||||
|
import { HierarchyAuditRepository } from './hierarchy-audit.repository.js';
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Hierarchy (tenancy/authorization structure) feature module.
|
||||||
|
*
|
||||||
|
* M4-1b-i ships the audit event + outbox machinery only (contract 1 §5.2).
|
||||||
|
* The hierarchy command family — controllers, DTOs, and the allowlisted
|
||||||
|
* class-table repositories — lands in M4-1b-ii once contract 2 (RBAC grant
|
||||||
|
* model) merges; until then this module exposes no routes, which the
|
||||||
|
* route-inventory witness asserts.
|
||||||
|
*/
|
||||||
|
@Module({
|
||||||
|
providers: [HierarchyAuditRepository],
|
||||||
|
exports: [HierarchyAuditRepository],
|
||||||
|
})
|
||||||
|
export class HierarchyModule {}
|
||||||
@@ -0,0 +1,40 @@
|
|||||||
|
CREATE TYPE "public"."hierarchy_outbox_status" AS ENUM('pending', 'processing', 'delivered');--> statement-breakpoint
|
||||||
|
CREATE TABLE "hierarchy_audit_events" (
|
||||||
|
"id" uuid PRIMARY KEY DEFAULT gen_random_uuid() NOT NULL,
|
||||||
|
"seq" bigint GENERATED ALWAYS AS IDENTITY (sequence name "hierarchy_audit_events_seq_seq" INCREMENT BY 1 MINVALUE 1 MAXVALUE 9223372036854775807 START WITH 1 CACHE 1),
|
||||||
|
"actor_id" text NOT NULL,
|
||||||
|
"verb" text NOT NULL,
|
||||||
|
"target_kind" text NOT NULL,
|
||||||
|
"target_id" uuid NOT NULL,
|
||||||
|
"target_snapshot" jsonb NOT NULL,
|
||||||
|
"transfer_from" jsonb,
|
||||||
|
"transfer_to" jsonb,
|
||||||
|
"correlation_id" text NOT NULL,
|
||||||
|
"causation_id" uuid,
|
||||||
|
"idempotency_key" text NOT NULL,
|
||||||
|
"occurred_at" timestamp with time zone DEFAULT now() NOT NULL,
|
||||||
|
CONSTRAINT "hierarchy_audit_events_verb_check" CHECK (verb IN ('create', 'rename', 'transfer', 'delete', 'grant_create', 'grant_change', 'grant_revoke')),
|
||||||
|
CONSTRAINT "hierarchy_audit_events_target_kind_check" CHECK (target_kind IN ('company', 'estate', 'platform_project', 'grant')),
|
||||||
|
CONSTRAINT "hierarchy_audit_events_transfer_check" CHECK ((verb = 'transfer') = (transfer_from IS NOT NULL AND transfer_to IS NOT NULL))
|
||||||
|
);
|
||||||
|
--> statement-breakpoint
|
||||||
|
CREATE TABLE "hierarchy_outbox" (
|
||||||
|
"id" uuid PRIMARY KEY DEFAULT gen_random_uuid() NOT NULL,
|
||||||
|
"event_id" uuid NOT NULL,
|
||||||
|
"idempotency_key" text NOT NULL,
|
||||||
|
"correlation_id" text NOT NULL,
|
||||||
|
"status" "hierarchy_outbox_status" DEFAULT 'pending' NOT NULL,
|
||||||
|
"created_at" timestamp with time zone DEFAULT now() NOT NULL,
|
||||||
|
"updated_at" timestamp with time zone DEFAULT now() NOT NULL,
|
||||||
|
"delivered_at" timestamp with time zone
|
||||||
|
);
|
||||||
|
--> statement-breakpoint
|
||||||
|
ALTER TABLE "hierarchy_audit_events" ADD CONSTRAINT "hierarchy_audit_events_causation_id_hierarchy_audit_events_id_fk" FOREIGN KEY ("causation_id") REFERENCES "public"."hierarchy_audit_events"("id") ON DELETE restrict ON UPDATE no action;--> statement-breakpoint
|
||||||
|
ALTER TABLE "hierarchy_outbox" ADD CONSTRAINT "hierarchy_outbox_event_id_hierarchy_audit_events_id_fk" FOREIGN KEY ("event_id") REFERENCES "public"."hierarchy_audit_events"("id") ON DELETE restrict ON UPDATE no action;--> statement-breakpoint
|
||||||
|
CREATE UNIQUE INDEX "hierarchy_audit_events_idempotency_idx" ON "hierarchy_audit_events" USING btree ("idempotency_key");--> statement-breakpoint
|
||||||
|
CREATE UNIQUE INDEX "hierarchy_audit_events_seq_idx" ON "hierarchy_audit_events" USING btree ("seq");--> statement-breakpoint
|
||||||
|
CREATE INDEX "hierarchy_audit_events_target_seq_idx" ON "hierarchy_audit_events" USING btree ("target_id","seq");--> statement-breakpoint
|
||||||
|
CREATE INDEX "hierarchy_audit_events_correlation_idx" ON "hierarchy_audit_events" USING btree ("correlation_id");--> statement-breakpoint
|
||||||
|
CREATE UNIQUE INDEX "hierarchy_outbox_event_idx" ON "hierarchy_outbox" USING btree ("event_id");--> statement-breakpoint
|
||||||
|
CREATE UNIQUE INDEX "hierarchy_outbox_idempotency_idx" ON "hierarchy_outbox" USING btree ("idempotency_key");--> statement-breakpoint
|
||||||
|
CREATE INDEX "hierarchy_outbox_status_created_idx" ON "hierarchy_outbox" USING btree ("status","created_at");
|
||||||
File diff suppressed because it is too large
Load Diff
@@ -134,6 +134,13 @@
|
|||||||
"when": 1787862158838,
|
"when": 1787862158838,
|
||||||
"tag": "0018_clean_cobalt_man",
|
"tag": "0018_clean_cobalt_man",
|
||||||
"breakpoints": true
|
"breakpoints": true
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"idx": 19,
|
||||||
|
"version": "7",
|
||||||
|
"when": 1787880918208,
|
||||||
|
"tag": "0019_volatile_killraven",
|
||||||
|
"breakpoints": true
|
||||||
}
|
}
|
||||||
]
|
]
|
||||||
}
|
}
|
||||||
@@ -0,0 +1,363 @@
|
|||||||
|
/**
|
||||||
|
* Hierarchy audit event + outbox schema witnesses — contract 1 §5.2 and the
|
||||||
|
* schema-level half of §6.4.
|
||||||
|
*
|
||||||
|
* Witnesses the guarantees the tables themselves carry: verb/target-kind/
|
||||||
|
* transfer CHECK constraints, idempotency uniqueness (REQ-AUD-001 duplicate
|
||||||
|
* suppression at the database level), monotonic append order (`seq`),
|
||||||
|
* deletion-safe linkage (no foreign key from the events table into any class
|
||||||
|
* table — events survive the deletion of their target), the causation
|
||||||
|
* self-FK, and the outbox's FK/uniqueness/status shape. The repository-level
|
||||||
|
* half of §6.4 (same-transaction atomicity, rollback, replay) is witnessed in
|
||||||
|
* apps/gateway/src/hierarchy/hierarchy-audit.integration.test.ts.
|
||||||
|
*
|
||||||
|
* Two legs run the same witness body:
|
||||||
|
* - PGlite (WASM Postgres): always runs.
|
||||||
|
* - Real PostgreSQL (§6.8): runs when DATABASE_URL is set — the binding
|
||||||
|
* witness; CI migrates ci-postgres before `pnpm test`.
|
||||||
|
*/
|
||||||
|
import { randomUUID } from 'node:crypto';
|
||||||
|
import { mkdtempSync, rmSync } from 'node:fs';
|
||||||
|
import { tmpdir } from 'node:os';
|
||||||
|
import { join } from 'node:path';
|
||||||
|
import { sql } from 'drizzle-orm';
|
||||||
|
import { afterAll, beforeAll, describe, expect, it } from 'vitest';
|
||||||
|
import { createDb } from './client.js';
|
||||||
|
import { createPgliteDb } from './client-pglite.js';
|
||||||
|
import { runPgliteMigrations } from './migrate.js';
|
||||||
|
import { companies, hierarchyAuditEvents, hierarchyOutbox } from './schema.js';
|
||||||
|
|
||||||
|
type AnyDb = {
|
||||||
|
db: {
|
||||||
|
insert: (t: unknown) => { values: (v: unknown) => Promise<unknown> };
|
||||||
|
execute: (q: unknown) => Promise<{ rows?: unknown[] } | unknown[]>;
|
||||||
|
};
|
||||||
|
close: () => Promise<void>;
|
||||||
|
};
|
||||||
|
|
||||||
|
/** Match a constraint failure anywhere along drizzle's cause chain. */
|
||||||
|
async function expectViolation(p: Promise<unknown>, re: RegExp, label = ''): Promise<void> {
|
||||||
|
let err: unknown;
|
||||||
|
try {
|
||||||
|
await p;
|
||||||
|
} catch (e) {
|
||||||
|
err = e;
|
||||||
|
}
|
||||||
|
expect(err, label || 'expected the statement to be refused').toBeDefined();
|
||||||
|
const messages: string[] = [];
|
||||||
|
let cur: unknown = err;
|
||||||
|
while (cur instanceof Error) {
|
||||||
|
messages.push(cur.message);
|
||||||
|
cur = (cur as { cause?: unknown }).cause;
|
||||||
|
}
|
||||||
|
expect(messages.join(' | '), label).toMatch(re);
|
||||||
|
}
|
||||||
|
|
||||||
|
function rows(res: { rows?: unknown[] } | unknown[]): Record<string, unknown>[] {
|
||||||
|
return (Array.isArray(res) ? res : (res.rows ?? [])) as Record<string, unknown>[];
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Unique per-run prefix so real-PG runs never collide and clean up safely. */
|
||||||
|
const T = `hier-a-${randomUUID().slice(0, 8)}`;
|
||||||
|
|
||||||
|
/** Minimal valid event row; overrides compose the negative cases. */
|
||||||
|
type EventInsert = typeof hierarchyAuditEvents.$inferInsert;
|
||||||
|
|
||||||
|
function eventRow(overrides: Partial<EventInsert> = {}): EventInsert {
|
||||||
|
return {
|
||||||
|
actorId: `${T}-actor`,
|
||||||
|
verb: 'create',
|
||||||
|
targetKind: 'company',
|
||||||
|
targetId: randomUUID(),
|
||||||
|
targetSnapshot: { id: 'x', slug: 'x', name: 'x', parentChain: [] },
|
||||||
|
correlationId: `${T}-corr`,
|
||||||
|
idempotencyKey: `${T}-${randomUUID()}`,
|
||||||
|
...overrides,
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
function witnessSuite(getHandle: () => AnyDb): void {
|
||||||
|
const db = () => getHandle().db as unknown as ReturnType<typeof createDb>['db'];
|
||||||
|
|
||||||
|
afterAll(async () => {
|
||||||
|
const d = db();
|
||||||
|
await d.execute(sql`DELETE FROM hierarchy_outbox WHERE idempotency_key LIKE ${T + '%'}`);
|
||||||
|
// Caused events first: the causation self-FK is RESTRICT.
|
||||||
|
await d.execute(
|
||||||
|
sql`DELETE FROM hierarchy_audit_events WHERE idempotency_key LIKE ${T + '%'} AND causation_id IS NOT NULL`,
|
||||||
|
);
|
||||||
|
await d.execute(sql`DELETE FROM hierarchy_audit_events WHERE idempotency_key LIKE ${T + '%'}`);
|
||||||
|
await d.execute(sql`DELETE FROM companies WHERE slug LIKE ${T + '%'}`);
|
||||||
|
});
|
||||||
|
|
||||||
|
// ── CHECK constraints ──────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
it('accepts every declared verb and refuses an undeclared one', async () => {
|
||||||
|
for (const verb of [
|
||||||
|
'create',
|
||||||
|
'rename',
|
||||||
|
'delete',
|
||||||
|
'grant_create',
|
||||||
|
'grant_change',
|
||||||
|
'grant_revoke',
|
||||||
|
]) {
|
||||||
|
await db().insert(hierarchyAuditEvents).values(eventRow({ verb }));
|
||||||
|
}
|
||||||
|
await expectViolation(
|
||||||
|
db()
|
||||||
|
.insert(hierarchyAuditEvents)
|
||||||
|
.values(eventRow({ verb: 'update' })),
|
||||||
|
/verb_check|violates check/i,
|
||||||
|
'undeclared verb must be refused',
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('refuses an undeclared target kind', async () => {
|
||||||
|
await expectViolation(
|
||||||
|
db()
|
||||||
|
.insert(hierarchyAuditEvents)
|
||||||
|
.values(eventRow({ targetKind: 'workspace' })),
|
||||||
|
/target_kind_check|violates check/i,
|
||||||
|
'workspace is not an audited target kind (workspace mutation is SOT-side)',
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('requires transfer snapshots exactly on transfers', async () => {
|
||||||
|
const parent = { kind: 'company', id: randomUUID(), slug: 'p' };
|
||||||
|
await db()
|
||||||
|
.insert(hierarchyAuditEvents)
|
||||||
|
.values(
|
||||||
|
eventRow({
|
||||||
|
verb: 'transfer',
|
||||||
|
targetKind: 'estate',
|
||||||
|
transferFrom: parent,
|
||||||
|
transferTo: { ...parent, id: randomUUID() },
|
||||||
|
}),
|
||||||
|
);
|
||||||
|
await expectViolation(
|
||||||
|
db()
|
||||||
|
.insert(hierarchyAuditEvents)
|
||||||
|
.values(eventRow({ verb: 'transfer' })),
|
||||||
|
/transfer_check|violates check/i,
|
||||||
|
'transfer without source/destination snapshots must be refused',
|
||||||
|
);
|
||||||
|
await expectViolation(
|
||||||
|
db()
|
||||||
|
.insert(hierarchyAuditEvents)
|
||||||
|
.values(eventRow({ verb: 'transfer', transferFrom: parent })),
|
||||||
|
/transfer_check|violates check/i,
|
||||||
|
'transfer with only the source snapshot must be refused',
|
||||||
|
);
|
||||||
|
await expectViolation(
|
||||||
|
db()
|
||||||
|
.insert(hierarchyAuditEvents)
|
||||||
|
.values(eventRow({ verb: 'create', transferFrom: parent, transferTo: parent })),
|
||||||
|
/transfer_check|violates check/i,
|
||||||
|
'non-transfer with transfer snapshots must be refused',
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
|
// ── Idempotency and ordering ───────────────────────────────────────────────
|
||||||
|
|
||||||
|
it('refuses a duplicate idempotency key', async () => {
|
||||||
|
const key = `${T}-dup-${randomUUID()}`;
|
||||||
|
await db()
|
||||||
|
.insert(hierarchyAuditEvents)
|
||||||
|
.values(eventRow({ idempotencyKey: key }));
|
||||||
|
await expectViolation(
|
||||||
|
db()
|
||||||
|
.insert(hierarchyAuditEvents)
|
||||||
|
.values(eventRow({ idempotencyKey: key })),
|
||||||
|
/duplicate key|unique/i,
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('assigns strictly increasing seq in insert order for one target', async () => {
|
||||||
|
const targetId = randomUUID();
|
||||||
|
const k1 = `${T}-seq-1-${randomUUID()}`;
|
||||||
|
const k2 = `${T}-seq-2-${randomUUID()}`;
|
||||||
|
await db()
|
||||||
|
.insert(hierarchyAuditEvents)
|
||||||
|
.values(eventRow({ targetId, idempotencyKey: k1 }));
|
||||||
|
await db()
|
||||||
|
.insert(hierarchyAuditEvents)
|
||||||
|
.values(eventRow({ targetId, verb: 'rename', idempotencyKey: k2 }));
|
||||||
|
const res = rows(
|
||||||
|
await db().execute(
|
||||||
|
sql`SELECT idempotency_key, seq FROM hierarchy_audit_events WHERE target_id = ${targetId} ORDER BY seq ASC`,
|
||||||
|
),
|
||||||
|
);
|
||||||
|
expect(res.map((r) => r['idempotency_key'])).toEqual([k1, k2]);
|
||||||
|
expect(Number(res[1]!['seq'])).toBeGreaterThan(Number(res[0]!['seq']));
|
||||||
|
});
|
||||||
|
|
||||||
|
// ── Deletion-safe linkage (§5.2) ───────────────────────────────────────────
|
||||||
|
|
||||||
|
it('has no foreign key into any class table, and events survive target deletion', async () => {
|
||||||
|
const fks = rows(
|
||||||
|
await db().execute(sql`
|
||||||
|
SELECT ccu.table_name AS referenced_table
|
||||||
|
FROM information_schema.table_constraints tc
|
||||||
|
JOIN information_schema.constraint_column_usage ccu
|
||||||
|
ON ccu.constraint_name = tc.constraint_name AND ccu.constraint_schema = tc.constraint_schema
|
||||||
|
WHERE tc.constraint_type = 'FOREIGN KEY' AND tc.table_name = 'hierarchy_audit_events'
|
||||||
|
`),
|
||||||
|
);
|
||||||
|
// The causation self-FK is the ONLY foreign key on the events table.
|
||||||
|
expect([...new Set(fks.map((r) => r['referenced_table']))]).toEqual(['hierarchy_audit_events']);
|
||||||
|
|
||||||
|
const companyId = randomUUID();
|
||||||
|
await db()
|
||||||
|
.insert(companies)
|
||||||
|
.values({ id: companyId, name: 'Doomed', slug: `${T}-doomed` });
|
||||||
|
const key = `${T}-survive-${randomUUID()}`;
|
||||||
|
await db()
|
||||||
|
.insert(hierarchyAuditEvents)
|
||||||
|
.values(
|
||||||
|
eventRow({
|
||||||
|
verb: 'delete',
|
||||||
|
targetId: companyId,
|
||||||
|
targetSnapshot: { id: companyId, slug: `${T}-doomed`, name: 'Doomed', parentChain: [] },
|
||||||
|
idempotencyKey: key,
|
||||||
|
}),
|
||||||
|
);
|
||||||
|
await db().execute(sql`DELETE FROM companies WHERE id = ${companyId}`);
|
||||||
|
const after = rows(
|
||||||
|
await db().execute(
|
||||||
|
sql`SELECT target_snapshot FROM hierarchy_audit_events WHERE idempotency_key = ${key}`,
|
||||||
|
),
|
||||||
|
);
|
||||||
|
expect(after).toHaveLength(1);
|
||||||
|
expect((after[0]!['target_snapshot'] as { id: string }).id).toBe(companyId);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('enforces the causation self-FK and RESTRICTs deleting a cause', async () => {
|
||||||
|
await expectViolation(
|
||||||
|
db()
|
||||||
|
.insert(hierarchyAuditEvents)
|
||||||
|
.values(eventRow({ causationId: randomUUID() })),
|
||||||
|
/foreign key/i,
|
||||||
|
'causation must reference an existing event',
|
||||||
|
);
|
||||||
|
const causeKey = `${T}-cause-${randomUUID()}`;
|
||||||
|
await db()
|
||||||
|
.insert(hierarchyAuditEvents)
|
||||||
|
.values(eventRow({ verb: 'delete', idempotencyKey: causeKey }));
|
||||||
|
const cause = rows(
|
||||||
|
await db().execute(
|
||||||
|
sql`SELECT id FROM hierarchy_audit_events WHERE idempotency_key = ${causeKey}`,
|
||||||
|
),
|
||||||
|
)[0]!;
|
||||||
|
await db()
|
||||||
|
.insert(hierarchyAuditEvents)
|
||||||
|
.values(
|
||||||
|
eventRow({ verb: 'grant_revoke', targetKind: 'grant', causationId: cause['id'] as string }),
|
||||||
|
);
|
||||||
|
await expectViolation(
|
||||||
|
db().execute(sql`DELETE FROM hierarchy_audit_events WHERE id = ${cause['id'] as string}`),
|
||||||
|
/foreign key/i,
|
||||||
|
'a cause with dependent events must not be deletable',
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
|
// ── Outbox shape ───────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
it('outbox rows require an existing event, one outbox row per event, unique idempotency', async () => {
|
||||||
|
await expectViolation(
|
||||||
|
db()
|
||||||
|
.insert(hierarchyOutbox)
|
||||||
|
.values({
|
||||||
|
eventId: randomUUID(),
|
||||||
|
idempotencyKey: `${T}-ob-${randomUUID()}`,
|
||||||
|
correlationId: `${T}-corr`,
|
||||||
|
}),
|
||||||
|
/foreign key/i,
|
||||||
|
'outbox must reference an existing event',
|
||||||
|
);
|
||||||
|
const key = `${T}-ob-${randomUUID()}`;
|
||||||
|
await db()
|
||||||
|
.insert(hierarchyAuditEvents)
|
||||||
|
.values(eventRow({ idempotencyKey: key }));
|
||||||
|
const event = rows(
|
||||||
|
await db().execute(sql`SELECT id FROM hierarchy_audit_events WHERE idempotency_key = ${key}`),
|
||||||
|
)[0]!;
|
||||||
|
const eventId = event['id'] as string;
|
||||||
|
await db()
|
||||||
|
.insert(hierarchyOutbox)
|
||||||
|
.values({ eventId, idempotencyKey: key, correlationId: `${T}-corr` });
|
||||||
|
await expectViolation(
|
||||||
|
db()
|
||||||
|
.insert(hierarchyOutbox)
|
||||||
|
.values({ eventId, idempotencyKey: `${T}-ob2-${randomUUID()}`, correlationId: `${T}-c` }),
|
||||||
|
/duplicate key|unique/i,
|
||||||
|
'one outbox record per event',
|
||||||
|
);
|
||||||
|
await expectViolation(
|
||||||
|
db().execute(
|
||||||
|
sql`INSERT INTO hierarchy_outbox (event_id, idempotency_key, correlation_id, status)
|
||||||
|
VALUES (${eventId}, ${`${T}-ob3-${randomUUID()}`}, 'c', 'failed')`,
|
||||||
|
),
|
||||||
|
/invalid input value for enum|22P02/i,
|
||||||
|
'status outside pending/processing/delivered must be refused',
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('outbox FK RESTRICTs event deletion while the outbox row exists', async () => {
|
||||||
|
const key = `${T}-obr-${randomUUID()}`;
|
||||||
|
await db()
|
||||||
|
.insert(hierarchyAuditEvents)
|
||||||
|
.values(eventRow({ idempotencyKey: key }));
|
||||||
|
const event = rows(
|
||||||
|
await db().execute(sql`SELECT id FROM hierarchy_audit_events WHERE idempotency_key = ${key}`),
|
||||||
|
)[0]!;
|
||||||
|
await db()
|
||||||
|
.insert(hierarchyOutbox)
|
||||||
|
.values({
|
||||||
|
eventId: event['id'] as string,
|
||||||
|
idempotencyKey: key,
|
||||||
|
correlationId: `${T}-corr`,
|
||||||
|
});
|
||||||
|
await expectViolation(
|
||||||
|
db().execute(sql`DELETE FROM hierarchy_audit_events WHERE id = ${event['id'] as string}`),
|
||||||
|
/foreign key/i,
|
||||||
|
);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
// ── Leg 1: PGlite (always runs — local witness signal) ───────────────────────
|
||||||
|
|
||||||
|
describe('hierarchy audit witnesses — PGlite', () => {
|
||||||
|
let dir: string;
|
||||||
|
let handle: ReturnType<typeof createPgliteDb>;
|
||||||
|
|
||||||
|
beforeAll(async () => {
|
||||||
|
dir = mkdtempSync(join(tmpdir(), 'hier-audit-witness-'));
|
||||||
|
handle = createPgliteDb(dir);
|
||||||
|
await runPgliteMigrations(handle);
|
||||||
|
});
|
||||||
|
|
||||||
|
afterAll(async () => {
|
||||||
|
await handle.close();
|
||||||
|
rmSync(dir, { recursive: true, force: true });
|
||||||
|
});
|
||||||
|
|
||||||
|
witnessSuite(() => handle as unknown as AnyDb);
|
||||||
|
});
|
||||||
|
|
||||||
|
// ── Leg 2: real PostgreSQL (§6.8 — binding witness, ci-postgres in CI) ───────
|
||||||
|
|
||||||
|
const hasPostgres = Boolean(process.env['DATABASE_URL']);
|
||||||
|
|
||||||
|
describe.skipIf(!hasPostgres)('hierarchy audit witnesses — real PostgreSQL', () => {
|
||||||
|
let handle: ReturnType<typeof createDb>;
|
||||||
|
|
||||||
|
beforeAll(() => {
|
||||||
|
handle = createDb(process.env['DATABASE_URL']!);
|
||||||
|
});
|
||||||
|
|
||||||
|
afterAll(async () => {
|
||||||
|
await handle.close();
|
||||||
|
});
|
||||||
|
|
||||||
|
witnessSuite(() => handle as unknown as AnyDb);
|
||||||
|
});
|
||||||
@@ -123,18 +123,20 @@
|
|||||||
* tracked for M4-1b consideration); the counterfactual — flagging every
|
* tracked for M4-1b consideration); the counterfactual — flagging every
|
||||||
* bare identifier call — false-positives on essentially all callback
|
* bare identifier call — false-positives on essentially all callback
|
||||||
* code. Reviews of modules touching db handles carry this residual.
|
* code. Reviews of modules touching db handles carry this residual.
|
||||||
* - The scan perimeter is <root>/<pkg>/src for the three roots; production
|
* - The scan perimeter is the full <root>/<pkg> tree for the three roots
|
||||||
* TS outside a src/ directory (e.g. packages/mosaic/framework/**) is not
|
* (build output, tool caches, and dot-directories excluded), so
|
||||||
* scanned (verified free of db/driver/execute references at review time).
|
* production TS outside src/ — package configs, e2e helpers,
|
||||||
* Files excluded from the scan — test files and out-of-src modules — are
|
* packages/mosaic/framework/** — is scanned and conduit-visible
|
||||||
* also invisible as import-graph CONDUITS: test files are emitted to
|
* (widened from src/-only in M4-1b-i; the widened set was measured free
|
||||||
* dist, so a production module could launder a symbol or capability
|
* of every trigger token at the time). Files excluded from the scan —
|
||||||
* through a re-export in one. Importing a test module from production
|
* test files — are still invisible as import-graph CONDUITS: test files
|
||||||
* code is anomalous and review-visible; the blind spot is accepted as a
|
* are emitted to dist, so a production module could launder a symbol or
|
||||||
* residual, not closed.
|
* capability through a re-export in one. Importing a test module from
|
||||||
|
* production code is anomalous and review-visible; that blind spot is
|
||||||
|
* accepted as a residual, not closed.
|
||||||
*
|
*
|
||||||
* The writer allowlist names hierarchy command/repository modules ONLY. It is
|
* The writer allowlist names hierarchy command/repository modules ONLY. It is
|
||||||
* empty today: the hierarchy command family (M4-1b) has not landed, so no
|
* empty today: the hierarchy command family (M4-1b-ii) has not landed, so no
|
||||||
* production module may write the class tables. The infrastructure register
|
* production module may write the class tables. The infrastructure register
|
||||||
* holds legitimate non-hierarchy raw execution; registered modules are exempt
|
* holds legitimate non-hierarchy raw execution; registered modules are exempt
|
||||||
* from prong (iii) only — prongs (i) and (ii) apply to them with no
|
* from prong (iii) only — prongs (i) and (ii) apply to them with no
|
||||||
@@ -179,7 +181,8 @@ const CLASS_TABLES = [
|
|||||||
|
|
||||||
/**
|
/**
|
||||||
* Writer allowlist (§6.3b): hierarchy command/repository modules only.
|
* Writer allowlist (§6.3b): hierarchy command/repository modules only.
|
||||||
* EMPTY until the hierarchy command family lands (M4-1b). Adding a module
|
* EMPTY until the hierarchy command family lands (M4-1b-ii; M4-1b-i ships
|
||||||
|
* only the audit/outbox machinery, which writes no class table). Adding a module
|
||||||
* here is a contract-conformance decision reviewed under §5.1 — the module
|
* here is a contract-conformance decision reviewed under §5.1 — the module
|
||||||
* must be part of the Gateway hierarchy command path, and it must not export
|
* must be part of the Gateway hierarchy command path, and it must not export
|
||||||
* a function that executes caller-supplied SQL.
|
* a function that executes caller-supplied SQL.
|
||||||
@@ -291,20 +294,23 @@ function isTestPath(rel: string): boolean {
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** Directory names excluded from the walk: build output and tool caches only. */
|
||||||
|
const EXCLUDED_DIRS = new Set(['node_modules', 'dist', 'build', 'coverage', 'test-results']);
|
||||||
|
|
||||||
function collectSources(): string[] {
|
function collectSources(): string[] {
|
||||||
const files: string[] = [];
|
const files: string[] = [];
|
||||||
for (const root of SCAN_ROOTS) {
|
for (const root of SCAN_ROOTS) {
|
||||||
const rootDir = join(REPO_ROOT, root);
|
const rootDir = join(REPO_ROOT, root);
|
||||||
if (!existsSync(rootDir)) continue;
|
if (!existsSync(rootDir)) continue;
|
||||||
for (const pkg of readdirSync(rootDir)) {
|
for (const pkg of readdirSync(rootDir)) {
|
||||||
const srcDir = join(rootDir, pkg, 'src');
|
const pkgDir = join(rootDir, pkg);
|
||||||
if (!existsSync(srcDir) || !statSync(srcDir).isDirectory()) continue;
|
if (!statSync(pkgDir).isDirectory()) continue;
|
||||||
const walk = (dir: string): void => {
|
const walk = (dir: string): void => {
|
||||||
for (const entry of readdirSync(dir)) {
|
for (const entry of readdirSync(dir)) {
|
||||||
const full = join(dir, entry);
|
const full = join(dir, entry);
|
||||||
const st = statSync(full);
|
const st = statSync(full);
|
||||||
if (st.isDirectory()) {
|
if (st.isDirectory()) {
|
||||||
if (entry === 'node_modules' || entry === 'dist') continue;
|
if (EXCLUDED_DIRS.has(entry) || entry.startsWith('.')) continue;
|
||||||
walk(full);
|
walk(full);
|
||||||
} else if (EXTENSIONS.has(full.slice(full.lastIndexOf('.')))) {
|
} else if (EXTENSIONS.has(full.slice(full.lastIndexOf('.')))) {
|
||||||
const rel = relative(REPO_ROOT, full).split(sep).join('/');
|
const rel = relative(REPO_ROOT, full).split(sep).join('/');
|
||||||
@@ -312,7 +318,7 @@ function collectSources(): string[] {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
walk(srcDir);
|
walk(pkgDir);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return files.sort();
|
return files.sort();
|
||||||
|
|||||||
@@ -4,6 +4,7 @@
|
|||||||
*/
|
*/
|
||||||
|
|
||||||
import { sql } from 'drizzle-orm';
|
import { sql } from 'drizzle-orm';
|
||||||
|
import type { AnyPgColumn } from 'drizzle-orm/pg-core';
|
||||||
import {
|
import {
|
||||||
pgTable,
|
pgTable,
|
||||||
pgEnum,
|
pgEnum,
|
||||||
@@ -1152,3 +1153,111 @@ export const hierarchyGrants = pgTable(
|
|||||||
index('hierarchy_grants_granted_by_idx').on(t.grantedBy),
|
index('hierarchy_grants_granted_by_idx').on(t.grantedBy),
|
||||||
],
|
],
|
||||||
);
|
);
|
||||||
|
|
||||||
|
// ─── Hierarchy audit events + outbox (contract 1 §5.2) ──────────────────────
|
||||||
|
// NOT part of the record class (the class is exactly the five tables above).
|
||||||
|
// Append-only semantic audit log for hierarchy mutations, with a dedicated
|
||||||
|
// transactional outbox — hierarchy events are not workspace-scoped rows and
|
||||||
|
// do not ride the workspace outbox. Deletion-safe linkage: events reference
|
||||||
|
// their target by an immutable snapshot (id, slug, parent chain at event
|
||||||
|
// time), never by a foreign key into the class tables, so append-only events
|
||||||
|
// survive the deletion of their target. Append-only is enforced at the
|
||||||
|
// application layer (the hierarchy audit repository exposes no update/delete
|
||||||
|
// path for events); REQ-AUD-001's INSERT/SELECT-only database role is a
|
||||||
|
// deployment concern outside this schema.
|
||||||
|
|
||||||
|
export const HIERARCHY_AUDIT_VERBS = [
|
||||||
|
'create',
|
||||||
|
'rename',
|
||||||
|
'transfer',
|
||||||
|
'delete',
|
||||||
|
'grant_create',
|
||||||
|
'grant_change',
|
||||||
|
'grant_revoke',
|
||||||
|
] as const;
|
||||||
|
|
||||||
|
export const HIERARCHY_AUDIT_TARGET_KINDS = [
|
||||||
|
'company',
|
||||||
|
'estate',
|
||||||
|
'platform_project',
|
||||||
|
'grant',
|
||||||
|
] as const;
|
||||||
|
|
||||||
|
export const hierarchyAuditEvents = pgTable(
|
||||||
|
'hierarchy_audit_events',
|
||||||
|
{
|
||||||
|
id: uuid('id').primaryKey().defaultRandom(),
|
||||||
|
// Global append order; per-target ordering (REQ-AUD-001) is a filter on
|
||||||
|
// target_id ordered by seq.
|
||||||
|
seq: bigint('seq', { mode: 'number' }).notNull().generatedAlwaysAsIdentity(),
|
||||||
|
// No FK: audit events outlive every principal and every target (§5.2).
|
||||||
|
actorId: text('actor_id').notNull(),
|
||||||
|
verb: text('verb').notNull(),
|
||||||
|
targetKind: text('target_kind').notNull(),
|
||||||
|
targetId: uuid('target_id').notNull(),
|
||||||
|
// Immutable snapshot at event time. Node events: { id, slug, name,
|
||||||
|
// parentChain: [{ kind, id, slug }, …] root-first }. Grant events:
|
||||||
|
// { id, subject: { userId | teamId }, target: { kind, id }, role }
|
||||||
|
// (subject and role per contract 2 §4.4).
|
||||||
|
targetSnapshot: jsonb('target_snapshot').notNull(),
|
||||||
|
// Present exactly on transfers: snapshot of the source/destination
|
||||||
|
// parent ({ kind, id, slug }), CHECK-enforced below.
|
||||||
|
transferFrom: jsonb('transfer_from'),
|
||||||
|
transferTo: jsonb('transfer_to'),
|
||||||
|
correlationId: text('correlation_id').notNull(),
|
||||||
|
// Prior event in the causal chain (e.g. cascaded grant_revoke events
|
||||||
|
// caused by a node delete). Self-FK RESTRICT keeps the chain intact.
|
||||||
|
causationId: uuid('causation_id').references((): AnyPgColumn => hierarchyAuditEvents.id, {
|
||||||
|
onDelete: 'restrict',
|
||||||
|
}),
|
||||||
|
idempotencyKey: text('idempotency_key').notNull(),
|
||||||
|
occurredAt: timestamp('occurred_at', { withTimezone: true }).notNull().defaultNow(),
|
||||||
|
},
|
||||||
|
(t) => [
|
||||||
|
uniqueIndex('hierarchy_audit_events_idempotency_idx').on(t.idempotencyKey),
|
||||||
|
uniqueIndex('hierarchy_audit_events_seq_idx').on(t.seq),
|
||||||
|
index('hierarchy_audit_events_target_seq_idx').on(t.targetId, t.seq),
|
||||||
|
index('hierarchy_audit_events_correlation_idx').on(t.correlationId),
|
||||||
|
check(
|
||||||
|
'hierarchy_audit_events_verb_check',
|
||||||
|
sql`verb IN ('create', 'rename', 'transfer', 'delete', 'grant_create', 'grant_change', 'grant_revoke')`,
|
||||||
|
),
|
||||||
|
check(
|
||||||
|
'hierarchy_audit_events_target_kind_check',
|
||||||
|
sql`target_kind IN ('company', 'estate', 'platform_project', 'grant')`,
|
||||||
|
),
|
||||||
|
check(
|
||||||
|
'hierarchy_audit_events_transfer_check',
|
||||||
|
sql`(verb = 'transfer') = (transfer_from IS NOT NULL AND transfer_to IS NOT NULL)`,
|
||||||
|
),
|
||||||
|
],
|
||||||
|
);
|
||||||
|
|
||||||
|
export const hierarchyOutboxStatusEnum = pgEnum('hierarchy_outbox_status', [
|
||||||
|
'pending',
|
||||||
|
'processing',
|
||||||
|
'delivered',
|
||||||
|
]);
|
||||||
|
|
||||||
|
export const hierarchyOutbox = pgTable(
|
||||||
|
'hierarchy_outbox',
|
||||||
|
{
|
||||||
|
id: uuid('id').primaryKey().defaultRandom(),
|
||||||
|
// FK into the append-only events table (not a class table): never
|
||||||
|
// dangles, so RESTRICT is safe and keeps event/outbox integrity.
|
||||||
|
eventId: uuid('event_id')
|
||||||
|
.notNull()
|
||||||
|
.references(() => hierarchyAuditEvents.id, { onDelete: 'restrict' }),
|
||||||
|
idempotencyKey: text('idempotency_key').notNull(),
|
||||||
|
correlationId: text('correlation_id').notNull(),
|
||||||
|
status: hierarchyOutboxStatusEnum('status').notNull().default('pending'),
|
||||||
|
createdAt: timestamp('created_at', { withTimezone: true }).notNull().defaultNow(),
|
||||||
|
updatedAt: timestamp('updated_at', { withTimezone: true }).notNull().defaultNow(),
|
||||||
|
deliveredAt: timestamp('delivered_at', { withTimezone: true }),
|
||||||
|
},
|
||||||
|
(t) => [
|
||||||
|
uniqueIndex('hierarchy_outbox_event_idx').on(t.eventId),
|
||||||
|
uniqueIndex('hierarchy_outbox_idempotency_idx').on(t.idempotencyKey),
|
||||||
|
index('hierarchy_outbox_status_created_idx').on(t.status, t.createdAt),
|
||||||
|
],
|
||||||
|
);
|
||||||
|
|||||||
Reference in New Issue
Block a user