import { createHash, randomUUID } from 'node:crypto'; import { Inject, Injectable, Logger } from '@nestjs/common'; import { seal } from '@mosaicstack/auth'; import { agentAuditEvents, agentIdempotencyFence, agentOutbox, agents, and, eq, providerCredentials, users, type Db, } from '@mosaicstack/db'; import { DB } from '../database/database.module.js'; import type { HarnessRegistry } from '../harness/harness.registry.js'; import { HARNESS_REGISTRY } from '../harness/harness.tokens.js'; /** * Agent enrollment command repository (design * docs/plans/2026-08-29-agent-enrollment-command-design.md §3; contract 5 §4 * envelope; contract 3 §4.3 idempotency fence, ratified via §7 item 4). * * The ONLY writer of the enrollment family's tables (`agent_audit_events`, * `agent_outbox`, `agent_idempotency_fence`) and the only path that sets * `agents.harness`/`agents.enrolled_at`. Every enroll runs one transaction: * fence check → (replay | credential handling → agent insert → fence insert → * audit event + outbox), so state, fence, event, and outbox commit or roll * back together (§3.1 rule 6). * * Authorization (v1, §3.1 rule 4) is the AuthGuard-authenticated actor — no * hierarchy grant is consulted because v1 enrollment binds no hierarchy node. * The recorded fence authorization scope is therefore the constant * platform-user identity domain (§3.1 rule 5). * * Never-echo (§3.1 rule 1): the credential value reaches exactly one sink — * the sealed store write — and appears in no result, audit payload, outbox * row, or log line. Log lines here carry correlation ids and error names * only, never request fields. * * The single-write helper methods (writeSealedCredential, insertAgentRow, * insertFenceRow, appendEvent, insertOutboxRow) are ordinary decomposition; * the atomicity witnesses (§5.6) spy on them to inject faults at each write * point without any test-only production switch. */ export const ENROLLMENT_OPERATION = 'agent.enroll'; /** §3.1 rule 5: v1 authorization is grant-free, so the scope is the authenticated-user identity domain. */ const AUTHORIZATION_SCOPE = 'platform-user'; /** The single bounded collision shape (§3.1 rule 5): constant, identifying no record. */ const CONFLICT_MESSAGE = 'idempotency conflict'; /** One fixed message for every not_found cause — missing and unauthorized are indistinguishable (§3.2). */ const NOT_FOUND_MESSAGE = 'agent not found'; /** Closed per-family error enum (§3.3). 401 is produced by AuthGuard; 403 folds to not_found (§3.2). */ export type EnrollmentErrorCode = | 'validation_failed' | 'authentication_failed' | 'authorization_refused' | 'not_found' | 'conflict' | 'precondition_failed' | 'internal_fault'; export interface EnrollmentFailure { readonly ok: false; readonly error: EnrollmentErrorCode; readonly message: string; /** Refusals carry the correlation id too (contract 5 §4.3 end-to-end traceability). */ readonly correlationId: string; } export type EnrollmentResult = | ({ readonly ok: true; readonly correlationId: string } & T) | EnrollmentFailure; /** The persisted agent row; the table stores no credential material (§3.1 rule 1). */ export interface EnrolledAgentView { readonly id: string; readonly name: string; readonly provider: string; readonly model: string; readonly status: string; readonly harness: string | null; readonly persona: string | null; readonly ownerId: string | null; readonly enrolledAt: string | null; readonly createdAt: string; } export interface EnrollCredentialInput { readonly mode: 'reference' | 'intake'; readonly type?: 'api_key'; readonly value?: string; } export interface EnrollAgentInput { readonly actorId: string; readonly harness: string; readonly name: string; readonly persona?: string | null; readonly model: string; readonly provider: string; readonly credential: EnrollCredentialInput; readonly idempotencyKey: string; readonly correlationId?: string; /** Defense in depth below the DTO: anything but 'actor-bound' is refused (seed-only rule). */ readonly replayMode?: string; } type Tx = Pick; type AgentRow = typeof agents.$inferSelect; type FenceRow = typeof agentIdempotencyFence.$inferSelect; /** Raised inside the transaction when the fence insert lost a same-key race (§3.1 rule 5 concurrency). */ class ConcurrentEnrollmentError extends Error { constructor() { super('concurrent enrollment lost the fence race'); this.name = 'ConcurrentEnrollmentError'; } } function agentView(row: AgentRow): EnrolledAgentView { return { id: row.id, name: row.name, provider: row.provider, model: row.model, status: row.status, harness: row.harness, persona: row.systemPrompt, ownerId: row.ownerId, enrolledAt: row.enrolledAt?.toISOString() ?? null, createdAt: row.createdAt.toISOString(), }; } /** Key-order-independent serialization (jsonb precedent in hierarchy-audit). */ 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; const body = Object.keys(record) .sort() .map((key) => `${JSON.stringify(key)}:${canonicalJson(record[key])}`) .join(','); return `{${body}}`; } return JSON.stringify(value); } interface NormalizedEnrollment { readonly actorId: string; readonly harness: string; readonly name: string; readonly persona: string | null; readonly model: string; readonly provider: string; readonly credential: EnrollCredentialInput; readonly idempotencyKey: string; readonly correlationId: string; readonly digest: string; } /** * Canonicalized-payload digest (§3.1 rule 5). The input EXCLUDES the * credential value by construction: it covers mode and declared type only — * plaintext never reaches the hash. */ function digestOf( input: Omit, ): string { const canonical = canonicalJson({ harness: input.harness, name: input.name, persona: input.persona, model: input.model, provider: input.provider, credential: { mode: input.credential.mode, type: input.credential.type ?? null }, }); return createHash('sha256').update(canonical).digest('hex'); } @Injectable() export class EnrollmentRepository { private readonly logger = new Logger(EnrollmentRepository.name); constructor( @Inject(DB) private readonly db: Db, @Inject(HARNESS_REGISTRY) private readonly registry: HarnessRegistry, ) {} async enroll(input: EnrollAgentInput): Promise> { const correlationId = input.correlationId ?? randomUUID(); const fail = (error: EnrollmentErrorCode, message: string): EnrollmentFailure => ({ ok: false, error, message, correlationId, }); const harness = input.harness.trim(); const name = input.name.trim(); if (harness.length === 0) return fail('validation_failed', 'harness must be non-empty'); if (name.length === 0 || name.length > 200) { return fail('validation_failed', 'name must be non-empty and at most 200 characters'); } if (input.replayMode !== undefined && input.replayMode !== 'actor-bound') { // Seed-only rule (contract 3 §4.3): refused with nothing executed and no fence row. return fail('validation_failed', 'replayMode must be actor-bound'); } if (input.credential.mode === 'reference') { if (input.credential.type !== undefined || input.credential.value !== undefined) { return fail('validation_failed', 'a reference credential carries no type or value'); } } else if ( input.credential.type !== 'api_key' || typeof input.credential.value !== 'string' || input.credential.value.length === 0 ) { return fail('validation_failed', 'an intake credential requires type api_key and a value'); } // Syntactic validity ends above; a well-formed name the live registry // does not know is a precondition failure (§3.1 table). if (!this.registry.has(harness)) { return fail('precondition_failed', 'harness is not registered'); } const normalized: NormalizedEnrollment = { actorId: input.actorId, harness, name, persona: input.persona ?? null, model: input.model, provider: input.provider, credential: input.credential, idempotencyKey: input.idempotencyKey, correlationId, digest: digestOf({ harness, name, persona: input.persona ?? null, model: input.model, provider: input.provider, credential: input.credential, }), }; // Two attempts: a fence-race loser's transaction rolls back and the retry // resolves through the replay path against the winner's committed row — // or executes afresh if the winner aborted (§3.1 rule 5 concurrency). A // unique-violation race never surfaces as an unhandled internal fault. for (let attempt = 0; attempt < 2; attempt += 1) { try { return await this.db.transaction(async (tx) => this.enrollTx(tx, normalized)); } catch (error) { if (error instanceof ConcurrentEnrollmentError && attempt === 0) continue; if (error instanceof ConcurrentEnrollmentError) { return fail('conflict', CONFLICT_MESSAGE); } // §4.4 fail-closed: whatever broke, the transaction rolled back and // the refusal is the internal-fault class — no fallback write or read. this.logger.error( `agent.enroll failed closed (correlation=${correlationId}): ${ error instanceof Error ? error.name : 'unknown error' }`, ); return fail('internal_fault', 'internal fault'); } } return fail('internal_fault', 'internal fault'); } private async enrollTx( tx: Tx, input: NormalizedEnrollment, ): Promise> { const fence = await this.fenceFor(tx, input.idempotencyKey); if (fence) return this.replay(tx, fence, input); if (input.credential.mode === 'reference') { // §3.1 rule 3: the reference must resolve for (actor, provider). const existing = await tx .select({ id: providerCredentials.id }) .from(providerCredentials) .where( and( eq(providerCredentials.userId, input.actorId), eq(providerCredentials.provider, input.provider), ), ) .limit(1); if (existing.length === 0) { return { ok: false, error: 'precondition_failed', message: 'credential reference does not resolve', correlationId: input.correlationId, }; } } else { // §3.1 rule 2: sealed-store write inside THIS transaction — a later // failure rolls it back, leaving no orphan credential. await this.writeSealedCredential( tx, input.actorId, input.provider, input.credential.value as string, ); } const agentRow = await this.insertAgentRow(tx, input); const fenceRow = await this.insertFenceRow(tx, input, agentRow.id); if (!fenceRow) { // A same-(operation, key) winner committed first; abandon our writes. throw new ConcurrentEnrollmentError(); } await this.appendEvent(tx, { eventType: 'agent.enrolled', actorId: input.actorId, agentId: agentRow.id, correlationId: input.correlationId, // §3.1 rule 6 payload: harness, provider, name, credentialMode — no credential material. payload: { harness: input.harness, provider: input.provider, name: input.name, credentialMode: input.credential.mode, }, }); return { ok: true, correlationId: input.correlationId, agent: agentView(agentRow) }; } /** * Replay path (§3.1 rule 5): a fresh submission of a recorded * (operation, key). The actor is re-authorized exactly as a fresh * submission (v1: authenticated actor — the guard already ran); then mode, * scope, digest, and recorded-actor equality; then target-result read * authority (owner or admin) on the referenced agent. ANY failure refuses * with the single bounded conflict shape — constant, identifying no record. * A passing replay executes nothing and appends only the non-mutation * access event (with its outbox record — one outbox row per event). */ private async replay( tx: Tx, fence: FenceRow, input: NormalizedEnrollment, ): Promise> { const collision: EnrollmentFailure = { ok: false, error: 'conflict', message: CONFLICT_MESSAGE, correlationId: input.correlationId, }; if (fence.replayMode !== 'actor-bound') return collision; if (fence.authorizationScope !== AUTHORIZATION_SCOPE) return collision; if (fence.payloadDigest !== input.digest) return collision; if (fence.actorId !== input.actorId) return collision; const rows = await tx.select().from(agents).where(eq(agents.id, fence.outcomeAgentId)).limit(1); const agentRow = rows[0]; if (!agentRow) return collision; const authorized = agentRow.ownerId === input.actorId || (await this.isPlatformAdmin(tx, input.actorId)); if (!authorized) return collision; await this.appendEvent(tx, { eventType: 'agent.enrollment.replayed', actorId: input.actorId, agentId: agentRow.id, correlationId: input.correlationId, payload: { fenceId: fence.id }, }); return { ok: true, correlationId: input.correlationId, agent: agentView(agentRow) }; } /** * agent.enrollment.get (§3.2): owner-or-admin read. Unauthorized and * missing fold to the same not_found wire shape (no existence oracle). */ async getEnrollment( actorId: string, agentId: string, correlationId?: string, ): Promise> { const resolvedCorrelation = correlationId ?? randomUUID(); try { const rows = await this.db.select().from(agents).where(eq(agents.id, agentId)).limit(1); const row = rows[0]; if (row) { const authorized = row.ownerId === actorId || (await this.isPlatformAdmin(this.db, actorId)); if (authorized) { return { ok: true, correlationId: resolvedCorrelation, agent: agentView(row) }; } } return { ok: false, error: 'not_found', message: NOT_FOUND_MESSAGE, correlationId: resolvedCorrelation, }; } catch (error) { this.logger.error( `agent.enrollment.get failed closed (correlation=${resolvedCorrelation}): ${ error instanceof Error ? error.name : 'unknown error' }`, ); return { ok: false, error: 'internal_fault', message: 'internal fault', correlationId: resolvedCorrelation, }; } } private async fenceFor(tx: Tx, idempotencyKey: string): Promise { const rows = await tx .select() .from(agentIdempotencyFence) .where( and( eq(agentIdempotencyFence.operation, ENROLLMENT_OPERATION), eq(agentIdempotencyFence.idempotencyKey, idempotencyKey), ), ) .limit(1); return rows[0] ?? null; } private async isPlatformAdmin(tx: Tx, actorId: string): Promise { const rows = await tx .select({ role: users.role }) .from(users) .where(eq(users.id, actorId)) .limit(1); return rows[0]?.role === 'admin'; } /** * Sealed intake write, mirroring ProviderCredentialsService.store semantics * (seal-at-rest, one row per (userId, provider)) but on the enrollment * transaction (§3.1 rule 2). The plaintext exists only in this frame. */ async writeSealedCredential( tx: Tx, userId: string, provider: string, value: string, ): Promise { const encryptedValue = seal(value); await tx .insert(providerCredentials) .values({ userId, provider, credentialType: 'api_key', encryptedValue, metadata: null }) .onConflictDoUpdate({ target: [providerCredentials.userId, providerCredentials.provider], set: { credentialType: 'api_key', encryptedValue, metadata: null, updatedAt: new Date(), }, }); } async insertAgentRow(tx: Tx, input: NormalizedEnrollment): Promise { const rows = await tx .insert(agents) .values({ name: input.name, provider: input.provider, model: input.model, harness: input.harness, systemPrompt: input.persona, // §3.1 rule 4: owner is the authenticated actor; is_system stays default false. ownerId: input.actorId, enrolledAt: new Date(), }) .returning(); const row = rows[0]; if (!row) throw new Error('agent insert returned no row'); return row; } async insertFenceRow( tx: Tx, input: NormalizedEnrollment, outcomeAgentId: string, ): Promise { const rows = await tx .insert(agentIdempotencyFence) .values({ operation: ENROLLMENT_OPERATION, idempotencyKey: input.idempotencyKey, actorId: input.actorId, authorizationScope: AUTHORIZATION_SCOPE, payloadDigest: input.digest, replayMode: 'actor-bound', outcomeAgentId, }) .onConflictDoNothing() .returning(); return rows[0] ?? null; } /** Append one audit event and its outbox record on the caller's transaction (one outbox row per event). */ async appendEvent( tx: Tx, input: { eventType: 'agent.enrolled' | 'agent.enrollment.replayed'; actorId: string; agentId: string; correlationId: string; payload: Record; causationId?: string; }, ): Promise { const inserted = await tx .insert(agentAuditEvents) .values({ eventType: input.eventType, actorId: input.actorId, agentId: input.agentId, correlationId: input.correlationId, causationId: input.causationId ?? null, payload: input.payload, }) .returning(); const event = inserted[0]; if (!event) throw new Error('agent audit event insert returned no row'); await this.insertOutboxRow(tx, event.id, input.correlationId); } async insertOutboxRow(tx: Tx, eventId: string, correlationId: string): Promise { await tx.insert(agentOutbox).values({ eventId, correlationId }); } }