import fs from "node:fs/promises"; import path from "node:path"; import YAML from "yaml"; import { z } from "zod"; import { canonicalJson, sha256, type JsonObject } from "../core/json.js"; import { BUILTIN_PROVIDER_PROFILE_IDS, createBuiltinProviderProfileResolver } from "../agents/provider-profiles.js"; export type ModelAdapterPrivacyClass = "public" | "private" | "sensitive"; export type ModelAdapterExportClass = "public" | "restricted" | "forbidden"; export interface ModelAdapterDatasetReceipt { id: string; sha256: string; } export interface ModelAdapterReleaseIdentity { id: string; version: number; description: string; releasedAt: string; manifestSha256: string; checkpointSelector: { kind: "env"; reference: string }; providerProfile: string; baseModel: string; dataset: ModelAdapterDatasetReceipt; evals: ModelAdapterDatasetReceipt[]; capabilities: string[]; privacyClass: ModelAdapterPrivacyClass; exportClass: ModelAdapterExportClass; sourceManifest?: ModelAdapterDatasetReceipt | undefined; } export interface ModelAdapterIdentity extends ModelAdapterReleaseIdentity { binding: { checkpointReferenceSha256: string }; } const digest = z.string().regex(/^[a-f0-9]{64}$/); const identifier = z.string().min(1).max(200).regex(/^[a-z0-9][a-z0-9._/-]*$/); const publicModelId = z.string().min(3).max(300).regex(/^[A-Za-z0-9][A-Za-z0-9._/-]*$/); const env = z.string().regex(/^[A-Z_][A-Z0-9_]*$/); const receipt = z.object({ id: identifier, sha256: digest }).strict(); const releaseManifestSchema = z.object({ schemaVersion: z.literal(1), id: z.string().min(1).max(100).regex(/^[a-z0-9][a-z0-9-]*$/), version: z.number().int().positive(), description: z.string().min(1).max(1000), releasedAt: z.iso.datetime(), providerProfile: z.enum(BUILTIN_PROVIDER_PROFILE_IDS), baseModel: publicModelId, checkpoint: z.object({ env }).strict(), dataset: receipt, evals: z.array(receipt).min(1).max(32), capabilities: z.array(identifier).min(1).max(32), privacyClass: z.enum(["public", "private", "sensitive"]), exportClass: z.enum(["public", "restricted", "forbidden"]), sourceManifest: receipt.optional(), }).strict().superRefine((value, context) => { const duplicateCapability = findDuplicate(value.capabilities); if (duplicateCapability) context.addIssue({ code: "custom", path: ["capabilities"], message: `Duplicate capability ${duplicateCapability}` }); const duplicateEval = findDuplicate(value.evals.map((entry) => entry.id)); if (duplicateEval) context.addIssue({ code: "custom", path: ["evals"], message: `Duplicate eval receipt ${duplicateEval}` }); }); const deploymentCatalogSchema = z.object({ schemaVersion: z.literal(1), generation: z.number().int().positive(), expectedCatalogSha256: digest.optional(), expectedProcesses: z.array(z.string().min(1).max(200)).max(128).optional(), selections: z.array(z.object({ declaration: z.object({ id: z.string().min(1), version: z.number().int().positive() }).strict(), release: z.object({ id: z.string().min(1), version: z.number().int().positive(), manifestSha256: digest }).strict(), state: z.enum(["candidate", "active", "retired"]), }).strict()).max(256), }).strict().superRefine((value, context) => { const duplicateProcess = findDuplicate(value.expectedProcesses ?? []); if (duplicateProcess) context.addIssue({ code: "custom", path: ["expectedProcesses"], message: `Duplicate expected process ${duplicateProcess}` }); }); type ReleaseManifest = z.infer; type DeploymentCatalog = z.infer; export function modelAdapterReleaseJson(identity: ModelAdapterReleaseIdentity): JsonObject { const { sourceManifest: _sourceManifest, ...publicIdentity } = identity; return { ...publicIdentity, checkpointSelector: { ...identity.checkpointSelector }, dataset: { ...identity.dataset }, evals: identity.evals.map((x) => ({ ...x })), capabilities: [...identity.capabilities], ...(identity.sourceManifest ? { sourceManifest: { ...identity.sourceManifest } } : {}) }; } export function modelAdapterIdentityJson(identity: ModelAdapterIdentity): JsonObject { return { ...modelAdapterReleaseJson(identity), binding: { ...identity.binding } }; } export function modelAdapterIdentitySha256(identity: ModelAdapterIdentity): string { return sha256(canonicalJson(modelAdapterIdentityJson(identity))); } export function adapterReleaseKey(identity: Pick): string { return `${identity.id}@${identity.version}`; } export const modelAdapterReleaseIdentitySchema = z.object({ id: z.string(), version: z.number().int().positive(), description: z.string(), releasedAt: z.iso.datetime(), manifestSha256: digest, checkpointSelector: z.object({ kind: z.literal("env"), reference: env }).strict(), providerProfile: z.enum(BUILTIN_PROVIDER_PROFILE_IDS), baseModel: publicModelId, dataset: receipt, evals: z.array(receipt), capabilities: z.array(z.string()), privacyClass: z.enum(["public", "private", "sensitive"]), exportClass: z.enum(["public", "restricted", "forbidden"]), sourceManifest: receipt.optional(), }).strict(); export const modelAdapterIdentitySchema = modelAdapterReleaseIdentitySchema.extend({ binding: z.object({ checkpointReferenceSha256: digest }).strict() }).strict(); export interface LoadedAdapterCatalog { digest: string; generation: number; releases: readonly ModelAdapterReleaseIdentity[]; selectedByDeclaration: ReadonlyMap; } const privateCheckpointBindings = new WeakMap(); export async function loadAdapterCatalog(projectRoot: string, environment: NodeJS.ProcessEnv = process.env): Promise { const root = path.join(projectRoot, "adapters"); const releaseDirectory = path.join(root, "releases"); const deploymentPath = path.join(root, "deployment.yaml"); const releaseDirectoryExists = await fs.stat(releaseDirectory).then((stat) => stat.isDirectory()).catch(() => false); if (!releaseDirectoryExists) throw new Error("Adapter startup requires adapters/releases"); await fs.stat(deploymentPath).catch(() => { throw new Error("Adapter startup requires adapters/deployment.yaml"); }); const entries = await fs.readdir(releaseDirectory, { withFileTypes: true }); if (!entries.length) throw new Error("Adapter startup requires a complete release catalog"); const releases = new Map(); for (const entry of entries.sort((a, b) => a.name.localeCompare(b.name))) { if (!entry.isFile() || !/\.ya?ml$/i.test(entry.name)) continue; const manifest = releaseManifestSchema.parse(YAML.parse(await fs.readFile(path.join(releaseDirectory, entry.name), "utf8"))) as ReleaseManifest; const identity: ModelAdapterReleaseIdentity = deepFreeze({ id: manifest.id, version: manifest.version, description: manifest.description, releasedAt: manifest.releasedAt, manifestSha256: sha256(canonicalJson(manifest as unknown as JsonObject)), checkpointSelector: { kind: "env", reference: manifest.checkpoint.env }, providerProfile: manifest.providerProfile, baseModel: manifest.baseModel, dataset: { ...manifest.dataset }, evals: manifest.evals.map((x) => ({ ...x })), capabilities: [...manifest.capabilities].sort(), privacyClass: manifest.privacyClass, exportClass: manifest.exportClass, ...(manifest.sourceManifest ? { sourceManifest: { ...manifest.sourceManifest } } : {}), }); if (releases.has(adapterReleaseKey(identity))) throw new Error(`Duplicate release ${adapterReleaseKey(identity)}`); releases.set(adapterReleaseKey(identity), identity); } if (!releases.size) throw new Error("Adapter startup requires at least one release manifest"); const deployment = deploymentCatalogSchema.parse(YAML.parse(await fs.readFile(deploymentPath, "utf8"))) as DeploymentCatalog; const deploymentIdentity: JsonObject = { schemaVersion: deployment.schemaVersion, generation: deployment.generation, ...(deployment.expectedProcesses ? { expectedProcesses: [...deployment.expectedProcesses].sort() } : {}), selections: [...deployment.selections].sort((left, right) => ( left.declaration.id.localeCompare(right.declaration.id) || left.declaration.version - right.declaration.version )) as unknown as JsonObject["selections"], }; const bundle: JsonObject = { releases: [...releases.values()].map(modelAdapterReleaseJson).sort((a, b) => ( String(a.id).localeCompare(String(b.id)) || Number(a.version) - Number(b.version) )), deployment: deploymentIdentity, }; const catalogDigest = sha256(canonicalJson(bundle)); if (deployment.expectedCatalogSha256 && deployment.expectedCatalogSha256 !== catalogDigest) throw new Error("Deployment catalog digest mismatch"); if (deployment.expectedProcesses?.length) { const processIdentity = environment.THOUGHTSTREAM_ADAPTER_PROCESS_ID?.trim(); if (!processIdentity || !deployment.expectedProcesses.includes(processIdentity)) { throw new Error("Adapter deployment does not authorize this process identity"); } } const selectedByDeclaration = new Map(); const seenDeclarations = new Set(); for (const selected of deployment.selections) { const declarationKey = `${selected.declaration.id}@${selected.declaration.version}`; if (seenDeclarations.has(declarationKey)) throw new Error(`Ambiguous selected deployment ${declarationKey}`); seenDeclarations.add(declarationKey); const release = releases.get(`${selected.release.id}@${selected.release.version}`); if (!release || release.manifestSha256 !== selected.release.manifestSha256) throw new Error(`Selected release is missing or has a digest mismatch: ${selected.release.id}@${selected.release.version}`); if (selected.state !== "active") continue; const checkpoint = environment[release.checkpointSelector.reference]?.trim(); if (!checkpoint) throw new Error(`Selected release checkpoint is unresolved: ${release.id}@${release.version}`); const providerProfiles = createBuiltinProviderProfileResolver(environment); providerProfiles.resolve(release.providerProfile, release.baseModel); providerProfiles.resolve(release.providerProfile, checkpoint); const binding = deepFreeze({ ...release, binding: { checkpointReferenceSha256: sha256(checkpoint) } }); privateCheckpointBindings.set(binding, checkpoint); selectedByDeclaration.set(declarationKey, binding); } return deepFreeze({ digest: catalogDigest, generation: deployment.generation, releases: [...releases.values()], selectedByDeclaration: readonlyMap(selectedByDeclaration), }); } export function privateCheckpointFor(identity: ModelAdapterIdentity): string { const checkpoint = privateCheckpointBindings.get(identity); if (!checkpoint || sha256(checkpoint) !== identity.binding.checkpointReferenceSha256) { throw new Error(`Adapter binding checkpoint is unavailable or mismatched: ${adapterReleaseKey(identity)}`); } return checkpoint; } function deepFreeze(value: T): T { if (value && typeof value === "object") { Object.freeze(value); for (const child of Object.values(value as Record)) deepFreeze(child); } return value; } function readonlyMap(source: Map): ReadonlyMap { return new Proxy(source, { get(target, property) { if (property === "set" || property === "delete" || property === "clear") return undefined; const value = Reflect.get(target, property, target); return typeof value === "function" ? value.bind(target) : value; }, }); } function findDuplicate(values: string[]): string | undefined { const seen = new Set(); for (const value of values) { if (seen.has(value)) return value; seen.add(value); } return undefined; }