feat(gateway,cli): agent enrollment command family (M4-4b) (#1483)
ci/woodpecker/push/publish Pipeline was successful
ci/woodpecker/push/publish Pipeline was successful
Co-authored-by: fred <[email protected]>
This commit was merged in pull request #1483.
This commit is contained in:
@@ -0,0 +1,538 @@
|
||||
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<T> =
|
||||
| ({ 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<Db, 'insert' | 'select' | 'update' | 'delete'>;
|
||||
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<string, unknown>;
|
||||
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<NormalizedEnrollment, 'actorId' | 'idempotencyKey' | 'correlationId' | 'digest'>,
|
||||
): 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<EnrollmentResult<{ agent: EnrolledAgentView }>> {
|
||||
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<EnrollmentResult<{ agent: EnrolledAgentView }>> {
|
||||
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<EnrollmentResult<{ agent: EnrolledAgentView }>> {
|
||||
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<EnrollmentResult<{ agent: EnrolledAgentView }>> {
|
||||
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<FenceRow | null> {
|
||||
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<boolean> {
|
||||
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<void> {
|
||||
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<AgentRow> {
|
||||
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<FenceRow | null> {
|
||||
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<string, unknown>;
|
||||
causationId?: string;
|
||||
},
|
||||
): Promise<void> {
|
||||
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<void> {
|
||||
await tx.insert(agentOutbox).values({ eventId, correlationId });
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user