Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
22 kB · 500 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501import fs from "node:fs/promises";import os from "node:os";import path from "node:path";import YAML from "yaml";import { afterEach, describe, expect, test } from "vitest";import { loadAdapterCatalog, modelAdapterIdentityJson, privateCheckpointFor,} from "../src/adapters/model-adapters.js";import { loadAgentDeclarations } from "../src/agents/declarations.js";import { ThoughtAgentRuntime } from "../src/agents/runtime.js";import { AgentRunFailure, type AgentRunner } from "../src/agents/types.js";import { canonicalJson, sha256, type JsonObject } from "../src/core/json.js";import { stableKey } from "../src/core/ids.js";import { createDefaultRegistry } from "../src/events/registry.js";import type { JazzThoughtStore } from "../src/jazz/store.js";import { projectTrainingExamples, recordJudgment } from "../src/training/judgments.js";import { testStore } from "./helpers.js";
const roots: string[] = [];const stores: JazzThoughtStore[] = [];const checkpoint = "tinker://private-fixture-checkpoint";const environment = { THOUGHTSTREAM_TEST_ADAPTER_MODEL: checkpoint, THOUGHTSTREAM_TINKER_ALLOWED_MODELS: `Qwen/Qwen3.5-4B,${checkpoint}`,};
afterEach(async () => { await Promise.all(stores.splice(0).map((store) => store.close())); await Promise.all(roots.splice(0).map((root) => fs.rm(root, { recursive: true, force: true })));});
describe("immutable model-adapter startup catalog", () => { test("binds one active release privately and keeps release identity stable across deployment generations", async () => { const first = await fixtureCatalog({ generation: 1 }); const second = await fixtureCatalog({ generation: 2 }); const firstCatalog = await loadAdapterCatalog(first.root, environment); const secondCatalog = await loadAdapterCatalog(second.root, environment); const firstBinding = firstCatalog.selectedByDeclaration.get("conceptualizer@1")!; const secondBinding = secondCatalog.selectedByDeclaration.get("conceptualizer@1")!;
expect(firstBinding.manifestSha256).toBe(secondBinding.manifestSha256); expect(firstCatalog.digest).not.toBe(secondCatalog.digest); expect(firstCatalog.generation).toBe(1); expect(secondCatalog.generation).toBe(2); expect(privateCheckpointFor(firstBinding)).toBe(checkpoint); expect(JSON.stringify(modelAdapterIdentityJson(firstBinding))).not.toContain(checkpoint); expect(JSON.stringify(firstCatalog)).not.toContain(checkpoint); expect(Object.isFrozen(firstBinding)).toBe(true); expect(Object.isFrozen(firstBinding.binding)).toBe(true); expect(() => { (firstBinding as { description: string }).description = "mutated"; }).toThrow(); expect(() => (firstCatalog.selectedByDeclaration as Map<string, unknown>).set("forged@1", {})).toThrow();
const clone = JSON.parse(JSON.stringify(firstBinding)); expect(() => privateCheckpointFor(clone)).toThrow("unavailable or mismatched"); });
test("validates an expected catalog digest without hashing the digest field into itself", async () => { const fixture = await fixtureCatalog({ generation: 3 }); const initial = await loadAdapterCatalog(fixture.root, environment); await writeDeployment(fixture.root, { generation: 3, manifestSha256: fixture.manifestSha256, expectedCatalogSha256: initial.digest, }); await expect(loadAdapterCatalog(fixture.root, environment)).resolves.toMatchObject({ digest: initial.digest, generation: 3, }); });
test.each([ ["missing release", { releaseId: "missing" }, "missing or has a digest mismatch"], ["release digest mismatch", { manifestSha256: "0".repeat(64) }, "missing or has a digest mismatch"], ["unresolved checkpoint", { environment: { THOUGHTSTREAM_TINKER_ALLOWED_MODELS: "Qwen/Qwen3.5-4B" } }, "checkpoint is unresolved"], ["unallowlisted checkpoint", { environment: { ...environment, THOUGHTSTREAM_TINKER_ALLOWED_MODELS: "Qwen/Qwen3.5-4B" } }, "not allowlisted"], ["unallowlisted public base", { baseModel: "totally/fabricated-base" }, "not allowlisted"], ["catalog digest mismatch", { expectedCatalogSha256: "0".repeat(64) }, "catalog digest mismatch"], ["missing process identity", { expectedProcesses: ["consumer-a"] }, "does not authorize this process identity"], ["wrong process identity", { expectedProcesses: ["consumer-a"], processIdentity: "consumer-b" }, "does not authorize this process identity"], ] as Array<[string, FixtureOptions, string]>) ("fails closed for %s", async (_name, options, message) => { const fixture = await fixtureCatalog(options); const selectedEnvironment = { ...(options.environment ?? environment), ...(options.processIdentity ? { THOUGHTSTREAM_ADAPTER_PROCESS_ID: options.processIdentity } : {}), }; await expect(loadAdapterCatalog(fixture.root, selectedEnvironment)).rejects.toThrow(message); });
test.each(["candidate", "retired"] as const)("keeps a %s deployment unbound and rejects an enabled declaration", async (state) => { const fixture = await fixtureCatalog({ state }); const catalog = await loadAdapterCatalog(fixture.root, {}); expect(catalog.selectedByDeclaration.size).toBe(0); await writeAdapterAgent(fixture.root); await expect(loadAgentDeclarations(path.join(fixture.root, "agents"), environment, catalog)).rejects.toThrow( "does not match the loaded deployment catalog", ); });
test("accepts only an explicitly authorized process identity", async () => { const fixture = await fixtureCatalog({ expectedProcesses: ["consumer-a", "consumer-b"] }); await expect(loadAdapterCatalog(fixture.root, { ...environment, THOUGHTSTREAM_ADAPTER_PROCESS_ID: "consumer-b", })).resolves.toMatchObject({ generation: 1 }); });
test("keeps the private checkpoint out of allowlist failures", async () => { const fixture = await fixtureCatalog(); const failed = loadAdapterCatalog(fixture.root, { THOUGHTSTREAM_TEST_ADAPTER_MODEL: checkpoint, THOUGHTSTREAM_TINKER_ALLOWED_MODELS: "Qwen/Qwen3.5-4B", }); await expect(failed).rejects.toThrow("not allowlisted"); await expect(failed).rejects.not.toThrow(checkpoint); });
test("compiles the production conceptualizer against one exact active adapter selection", async () => { const fixture = await fixtureCatalog(); await fs.mkdir(path.join(fixture.root, "agents"), { recursive: true }); await fs.mkdir(path.join(fixture.root, "prompts"), { recursive: true }); const declaration = (await fs.readFile(path.join(process.cwd(), "agents", "conceptualizer.yaml"), "utf8")) .replace("version: 14", "version: 1") .replace(" profile: openai-json-default", " profile: tinker-default") .replace(" model: gpt-4.1-mini", " adapter: { id: fixture-adapter, version: 1 }"); await fs.writeFile(path.join(fixture.root, "agents", "conceptualizer.yaml"), declaration); await fs.copyFile( path.join(process.cwd(), "prompts", "conceptualizer.md"), path.join(fixture.root, "prompts", "conceptualizer.md"), );
const [compiled] = await loadAgentDeclarations(path.join(fixture.root, "agents"), environment); expect(compiled).toMatchObject({ id: "conceptualizer", enabled: true, mode: "pi", model: "Qwen/Qwen3.5-4B", outputEventType: "stream.thought.derived.concept.graph", adapterCatalogGeneration: 1, adapterCatalogDigest: expect.stringMatching(/^[a-f0-9]{64}$/), modelAdapter: { id: "fixture-adapter", version: 1, manifestSha256: fixture.manifestSha256, binding: { checkpointReferenceSha256: sha256(checkpoint) }, }, }); expect(Object.isFrozen(compiled)).toBe(true); expect(Object.isFrozen(compiled!.eventTypes)).toBe(true); expect(() => { compiled!.description = "mutated"; }).toThrow(); expect(privateCheckpointFor(compiled!.modelAdapter!)).toBe(checkpoint);
const store = await testStore(fixture.root); stores.push(store); const runtime = new ThoughtAgentRuntime(store, []); await expect(runtime.registerDeclarations([structuredClone(compiled!)])).rejects.toThrow( "not compiled by the trusted declaration loader", ); await expect(runtime.registerDeclarations([compiled!])).resolves.toBeUndefined(); });
test("persists exact adapter and catalog provenance while raising privacy and enforcing export policy", async () => { const fixture = await fixtureCatalog(); await writeAdapterAgent(fixture.root); const declarations = await loadAgentDeclarations(path.join(fixture.root, "agents"), environment); const declaration = declarations[0]!; const store = await testStore(fixture.root); stores.push(store); const source = (await store.appendEvent({ type: "stream.thought.runtime.notice", schemaVersion: 1, source: "fixture:adapter", sourceKind: "system", externalId: "adapter-provenance", idempotencyKey: "adapter-provenance", occurredAt: "2026-07-25T20:00:00.000Z", actor: "fixture:adapter", correlationId: "adapter-provenance", privacy: "public-source", payload: { text: "Synthetic adapter provenance fixture" }, })).event; const runner: AgentRunner = { mode: "pi", run: async () => ({ summary: "Adapter-backed result", tags: ["adapter"], importance: "normal", confidence: 0.9, model: { provider: "tinker", id: declaration.model!, revision: checkpoint, }, }), }; const runtime = new ThoughtAgentRuntime(store, [runner]);
const [result] = await runtime.consumeBacklog(declarations); expect(result?.error).toBeUndefined(); expect(result?.derivedEvent).toMatchObject({ privacy: "private", payload: { modelAdapter: declaration.modelAdapter, adapterCatalogDigest: declaration.adapterCatalogDigest, adapterCatalogGeneration: 1, }, }); const run = await store.getRun(result!.runId); expect(run).toMatchObject({ status: "completed", privacy: "private", modelAdapter: declaration.modelAdapter, adapterCatalogDigest: declaration.adapterCatalogDigest, adapterCatalogGeneration: 1, checkpointRevision: `sha256:${sha256(checkpoint)}`, contextManifest: { privacy: "private", modelAdapter: declaration.modelAdapter, adapterCatalogDigest: declaration.adapterCatalogDigest, adapterCatalogGeneration: 1, }, }); const completed = (await store.listEvents({ types: ["stream.thought.agent.run.completed"] })) .find((event) => event.payload.runId === run!.id); expect(completed).toMatchObject({ privacy: "private", payload: { modelAdapter: declaration.modelAdapter, adapterCatalogDigest: declaration.adapterCatalogDigest, adapterCatalogGeneration: 1, }, }); expect((await store.getProjection(stableKey("effective-output", run!.id)))?.payload.privacy).toBe("private");
await recordJudgment(store, { runId: run!.id, kind: "accept", criterion: "adapter-provenance", criterionVersion: 1, qualityEligible: true, externalExportEligible: true, sensitiveExternalExportAuthorized: true, }); expect(await projectTrainingExamples(store, { includeSensitivePrivate: true })).toEqual([]); const examples = await projectTrainingExamples(store, { includeSensitivePrivate: true, includeRestrictedModelAdapters: true, }); expect(examples).toHaveLength(1); expect(examples[0]?.provenance).toMatchObject({ modelAdapter: declaration.modelAdapter, adapterCatalogDigest: declaration.adapterCatalogDigest, adapterCatalogGeneration: 1, }); expect(examples[0]?.input.contextManifest).toMatchObject({ modelAdapter: declaration.modelAdapter, adapterCatalogDigest: declaration.adapterCatalogDigest, adapterCatalogGeneration: 1, }); const durable = JSON.stringify({ source, agents: await store.listAgents(), runs: await store.listRuns(), events: await store.listEvents(), traces: await store.listTrace(run!.id), examples, }); expect(durable).not.toContain(checkpoint); });
test("keeps adapter-backed output failures repair-eligible with private original provenance", async () => { const fixture = await fixtureCatalog(); await writeAdapterAgent(fixture.root); const declarations = await loadAgentDeclarations(path.join(fixture.root, "agents"), environment); const declaration = declarations[0]!; const store = await testStore(fixture.root); stores.push(store); await store.appendEvent({ type: "stream.thought.runtime.notice", schemaVersion: 1, source: "fixture:adapter", sourceKind: "system", externalId: "adapter-repair", idempotencyKey: "adapter-repair", occurredAt: "2026-07-25T20:00:00.000Z", actor: "fixture:adapter", correlationId: "adapter-repair", privacy: "public-source", payload: { text: "Synthetic adapter repair fixture" }, }); const runner: AgentRunner = { mode: "pi", run: async () => { throw new AgentRunFailure("Adapter output rejected", { diagnostic: { code: "invalid-final-output", stage: "final-output-validation", reason: "invalid-json", checkpointRevision: checkpoint, outputContract: declaration.outputContract as unknown as JsonObject, assistantMessages: 1, textParts: 1, textChars: 12, textSha256: sha256("invalid-json"), thinkingParts: 0, thinkingChars: 0, thinkingRedacted: true, toolCallParts: 0, otherParts: 0, stopReason: "stop", }, }); }, }; const runtime = new ThoughtAgentRuntime(store, [runner]);
const [result] = await runtime.consumeBacklog(declarations); expect(result).toMatchObject({ error: "Adapter output rejected" }); const run = await store.getRun(result!.runId); const failed = (await store.listEvents({ types: ["stream.thought.agent.run.failed"] }))[0]; const request = (await store.listEvents({ types: ["stream.thought.agent.repair.requested"] }))[0]; expect(run).toMatchObject({ privacy: "private", modelAdapter: declaration.modelAdapter }); expect(failed).toMatchObject({ privacy: "private" }); expect(request).toMatchObject({ privacy: "private", payload: { originalRunId: run!.id, model: { modelAdapter: declaration.modelAdapter, adapterCatalogDigest: declaration.adapterCatalogDigest, adapterCatalogGeneration: 1, }, }, }); expect((await store.getProjection(stableKey("effective-output", run!.id)))?.payload).toMatchObject({ status: "unresolved", privacy: "private", originalModel: { modelAdapter: declaration.modelAdapter, adapterCatalogDigest: declaration.adapterCatalogDigest, adapterCatalogGeneration: 1, }, }); expect(JSON.stringify({ run, failed, request })).not.toContain(checkpoint); });
test("rejects missing, empty, duplicate, and ambiguous catalog structure", async () => { const missingRoot = await temporaryRoot(); await expect(loadAdapterCatalog(missingRoot, environment)).rejects.toThrow("adapters/releases");
const noDeployment = await temporaryRoot(); await fs.mkdir(path.join(noDeployment, "adapters", "releases"), { recursive: true }); await fs.writeFile(path.join(noDeployment, "adapters", "releases", "fixture.yaml"), releaseManifest()); await expect(loadAdapterCatalog(noDeployment, environment)).rejects.toThrow("adapters/deployment.yaml");
const empty = await temporaryRoot(); await fs.mkdir(path.join(empty, "adapters", "releases"), { recursive: true }); await fs.writeFile(path.join(empty, "adapters", "deployment.yaml"), deploymentYaml({ manifestSha256: "0".repeat(64) })); await expect(loadAdapterCatalog(empty, environment)).rejects.toThrow("complete release catalog");
const duplicate = await fixtureCatalog(); await fs.copyFile( path.join(duplicate.root, "adapters", "releases", "fixture.yaml"), path.join(duplicate.root, "adapters", "releases", "duplicate.yaml"), ); await expect(loadAdapterCatalog(duplicate.root, environment)).rejects.toThrow("Duplicate release");
const ambiguous = await fixtureCatalog(); await writeDeployment(ambiguous.root, { manifestSha256: ambiguous.manifestSha256, duplicateSelection: true, }); await expect(loadAdapterCatalog(ambiguous.root, environment)).rejects.toThrow("Ambiguous selected deployment");
const duplicateCapability = await fixtureCatalog(); await fs.writeFile( path.join(duplicateCapability.root, "adapters", "releases", "fixture.yaml"), releaseManifest().replace("capabilities: [conceptualization]", "capabilities: [conceptualization, conceptualization]"), ); await expect(loadAdapterCatalog(duplicateCapability.root, environment)).rejects.toThrow("Duplicate capability");
const duplicateEval = await fixtureCatalog(); await fs.writeFile( path.join(duplicateEval.root, "adapters", "releases", "fixture.yaml"), releaseManifest().replace( `evals: [{ id: fixture-eval, sha256: ${"d".repeat(64)} }]`, `evals: [{ id: fixture-eval, sha256: ${"d".repeat(64)} }, { id: fixture-eval, sha256: ${"e".repeat(64)} }]`, ), ); await expect(loadAdapterCatalog(duplicateEval.root, environment)).rejects.toThrow("Duplicate eval receipt"); });
test("does not register the removed model-adapter lifecycle event surface", () => { const registry = createDefaultRegistry(); for (const type of [ "stream.thought.model.adapter.registered", "stream.thought.model.adapter.activated", "stream.thought.model.adapter.retired", ]) expect(registry.has(type, 1)).toBe(false); });});
interface FixtureOptions { generation?: number; state?: "candidate" | "active" | "retired"; releaseId?: string; manifestSha256?: string; expectedCatalogSha256?: string; expectedProcesses?: string[]; processIdentity?: string; environment?: NodeJS.ProcessEnv; baseModel?: string;}
async function fixtureCatalog(options: FixtureOptions = {}) { const root = await temporaryRoot(); await fs.mkdir(path.join(root, "adapters", "releases"), { recursive: true }); const manifest = releaseManifest(options); const manifestSha256 = sha256(canonicalJson(YAML.parse(manifest))); await fs.writeFile(path.join(root, "adapters", "releases", "fixture.yaml"), manifest); await writeDeployment(root, { ...options, manifestSha256: options.manifestSha256 ?? manifestSha256 }); return { root, manifestSha256 };}
async function writeDeployment(root: string, options: FixtureOptions & { manifestSha256: string; duplicateSelection?: boolean }) { await fs.mkdir(path.join(root, "adapters"), { recursive: true }); await fs.writeFile(path.join(root, "adapters", "deployment.yaml"), deploymentYaml(options));}
function releaseManifest(options: FixtureOptions = {}) { return [ "schemaVersion: 1", "id: fixture-adapter", "version: 1", "description: Fixture learned adapter", "releasedAt: 2026-07-25T20:00:00.000Z", "providerProfile: tinker-default", `baseModel: ${options.baseModel ?? "Qwen/Qwen3.5-4B"}`, "checkpoint: { env: THOUGHTSTREAM_TEST_ADAPTER_MODEL }", `dataset: { id: fixture-dataset, sha256: ${"c".repeat(64)} }`, `evals: [{ id: fixture-eval, sha256: ${"d".repeat(64)} }]`, "capabilities: [conceptualization]", "privacyClass: private", "exportClass: restricted", "", ].join("\n");}
function deploymentYaml(options: FixtureOptions & { manifestSha256: string; duplicateSelection?: boolean }) { const selection = [ " - declaration: { id: conceptualizer, version: 1 }", ` release: { id: ${options.releaseId ?? "fixture-adapter"}, version: 1, manifestSha256: "${options.manifestSha256}" }`, ` state: ${options.state ?? "active"}`, ]; return [ "schemaVersion: 1", `generation: ${options.generation ?? 1}`, ...(options.expectedCatalogSha256 ? [`expectedCatalogSha256: "${options.expectedCatalogSha256}"`] : []), ...(options.expectedProcesses ? ["expectedProcesses:", ...options.expectedProcesses.map((value) => ` - ${value}`)] : []), "selections:", ...selection, ...(options.duplicateSelection ? selection : []), "", ].join("\n");}
async function temporaryRoot() { const root = await fs.mkdtemp(path.join(os.tmpdir(), "thoughtstream-model-adapters-")); roots.push(root); return root;}
async function writeAdapterAgent(root: string) { await fs.mkdir(path.join(root, "agents"), { recursive: true }); await fs.mkdir(path.join(root, "prompts"), { recursive: true }); await fs.writeFile(path.join(root, "prompts", "adapter.md"), "Observe the synthetic fixture.\n"); await fs.writeFile(path.join(root, "agents", "adapter.yaml"), [ "id: conceptualizer", "version: 1", "description: Adapter provenance fixture", "enabled: true", "subscribe: { types: [stream.thought.runtime.notice], sources: [fixture:adapter], privacy: [public-source] }", "context: { maxEvents: 1, maxChars: 2000, strategy: single-event }", "runner:", " kind: pi", " profile: tinker-default", " adapter: { id: fixture-adapter, version: 1 }", " maxOutputTokens: 200", " timeoutMs: 60000", "accounting:", " leaseMs: 70000", " reservation: { inputTokens: 1000, outputTokens: 200 }", " limits: [{ window: hour, maxCalls: 10, maxInputTokens: 10000, maxOutputTokens: 2000 }]", "prompt: prompts/adapter.md", "emit: [stream.thought.derived.topics]", "policy: { tools: [], externalActions: false }", "", ].join("\n"));}