Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
12 kB · 197 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198import 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<typeof releaseManifestSchema>;type DeploymentCatalog = z.infer<typeof deploymentCatalogSchema>;
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<ModelAdapterReleaseIdentity, "id" | "version">): 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<string, ModelAdapterIdentity>;}
const privateCheckpointBindings = new WeakMap<ModelAdapterIdentity, string>();
export async function loadAdapterCatalog(projectRoot: string, environment: NodeJS.ProcessEnv = process.env): Promise<LoadedAdapterCatalog> { 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<string, ModelAdapterReleaseIdentity>(); 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<string, ModelAdapterIdentity>(); const seenDeclarations = new Set<string>(); 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<T>(value: T): T { if (value && typeof value === "object") { Object.freeze(value); for (const child of Object.values(value as Record<string, unknown>)) deepFreeze(child); } return value; }
function readonlyMap<K, V>(source: Map<K, V>): ReadonlyMap<K, V> { 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<string>(); for (const value of values) { if (seen.has(value)) return value; seen.add(value); } return undefined;}