Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
26 kB · 662 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663import fs from "node:fs/promises";import path from "node:path";import { randomUUID } from "node:crypto";import { afterEach, describe, expect, test } from "vitest";import { outputContractForDeclaration, outputContractIdentityJson } from "../src/agents/output-contracts.js";import type { ThoughtAgentDeclaration } from "../src/agents/types.js";import type { JsonObject } from "../src/core/json.js";import type { JazzThoughtStore } from "../src/jazz/store.js";import { stableKey } from "../src/core/ids.js";import { rebuildEffectiveOutputForRun } from "../src/projections/effective-output.js";import { projectTrainingExamples, recordJudgment, writeTrainingJsonl } from "../src/training/judgments.js";import { temporaryProject, testStore } from "./helpers.js";
const stores: JazzThoughtStore[] = [];const roots: string[] = [];
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("training judgments", () => { test("exports public synthetic judgments through the provenance-complete v3 format in owner-only files", async () => { const fixture = await completedRunFixture("public-source"); const { root, store, runId } = fixture; const replacement: JsonObject = { summary: "Corrected public output", tags: ["right"], importance: "normal", confidence: 0.9, }; const judgment = await recordJudgment(store, { runId, kind: "correct", criterion: "topic-fidelity", criterionVersion: 1, replacementOutput: replacement, qualityEligible: true, externalExportEligible: true, notes: "Synthetic public correction.", });
expect(judgment).toMatchObject({ schemaVersion: 2, privacy: "public-source", payload: { qualityEligible: true, externalExportEligible: true }, }); const examples = await projectTrainingExamples(store); expect(examples).toHaveLength(1); expect(examples[0]).toMatchObject({ format: "thoughtstream.training-example.v3", kind: "correct", judgment: { criterion: "topic-fidelity", criterionVersion: 1 }, input: { event: { type: "stream.thought.runtime.notice", schemaVersion: 1, sourceKind: "system", privacy: "public-source", }, contextManifest: { agentRole: "standard" }, }, trajectory: [{ sequence: 1, type: "pi.message_end" }], chosen: replacement, rejected: fixture.original, provenance: { agentVersion: 3, provider: "tinker", model: "Qwen/Qwen3.5-4B", executionAdapterRevision: "checkpoint-fixture", }, }); const exported = JSON.stringify(examples); for (const sentinel of fixture.privateMetadataSentinels) expect(exported).not.toContain(sentinel); for (const forbiddenKey of [ "externalId", "idempotencyKey", "correlationId", "payloadHash", "actor", "runId", "outputEventId", "eventId", "payloadSha256", "createdAt", ]) expect(exported).not.toContain(`\"${forbiddenKey}\"`);
const destination = path.join(root, "exports", "training.jsonl"); const manifest = await writeTrainingJsonl(destination, examples); expect(manifest).toMatchObject({ format: "thoughtstream.training-dataset-manifest.v4", datasetId: expect.stringMatching(/^sha256:/), examples: 1, kinds: { correct: 1 }, models: ["tinker:Qwen/Qwen3.5-4B@execution:checkpoint-fixture"], }); expect(JSON.parse((await fs.readFile(destination, "utf8")).trim())).toMatchObject({ format: "thoughtstream.training-example.v3", chosen: { summary: "Corrected public output" }, }); const storedManifest = JSON.parse(await fs.readFile(`${destination}.manifest.json`, "utf8")) as Record<string, unknown>; expect(storedManifest).toMatchObject({ format: "thoughtstream.training-dataset-manifest.v4", examples: 1 }); expect(storedManifest).not.toHaveProperty("judgmentEventIds"); expect((await fs.stat(destination)).mode & 0o777).toBe(0o600); expect((await fs.stat(`${destination}.manifest.json`)).mode & 0o777).toBe(0o600); });
test("keeps schema-invalid and output-mismatched direct correction rows inert", async () => { const fixture = await completedRunFixture("public-source"); const base = { type: "stream.thought.judgment.training-example", schemaVersion: 2, sourceKind: "system" as const, occurredAt: "2026-07-15T00:00:03.000Z", actor: "corrupt-history-fixture", rootEventId: fixture.sourceEventId, parentEventId: fixture.outputEventId, correlationId: fixture.runId, privacy: "public-source" as const, createdByRuntime: "historical-fixture", }; await fixture.store.appendEvent({ ...base, source: "judgment:invalid-direct-correction", externalId: "invalid-direct-correction", idempotencyKey: "invalid-direct-correction", payload: { runId: fixture.runId, outputEventId: fixture.outputEventId, kind: "correct", criterion: "historical-corruption", criterionVersion: 1, qualityEligible: true, externalExportEligible: false, replacementOutput: { summary: "Missing required fields" }, }, }); await fixture.store.appendEvent({ ...base, source: "judgment:mismatched-direct-correction", externalId: "mismatched-direct-correction", idempotencyKey: "mismatched-direct-correction", payload: { runId: fixture.runId, outputEventId: "evt_not_the_run_output", kind: "correct", criterion: "historical-corruption", criterionVersion: 2, qualityEligible: true, externalExportEligible: false, replacementOutput: { summary: "Valid shape, wrong target", tags: [], importance: "low", confidence: 1, }, }, });
const projection = await rebuildEffectiveOutputForRun(fixture.store, fixture.runId); expect(projection).toMatchObject({ id: stableKey("effective-output", fixture.runId), projectionVersion: 2, lastEventId: fixture.outputEventId, payload: { status: "original", outputEventId: fixture.outputEventId, structuredOutput: fixture.original, }, }); });
test("keeps private quality judgments out of default export and refuses explicit private data inside Git", async () => { const fixture = await completedRunFixture("sensitive"); const { root, store, runId } = fixture; await recordJudgment(store, { runId, kind: "accept", criterion: "telegram-reaction", criterionVersion: 1, qualityEligible: true, externalExportEligible: false, }); await expect(recordJudgment(store, { runId, kind: "accept", criterion: "operator-reviewed-private-export", criterionVersion: 1, qualityEligible: true, externalExportEligible: true, })).rejects.toThrow("requires explicit authorization"); await recordJudgment(store, { runId, kind: "accept", criterion: "operator-reviewed-private-export", criterionVersion: 1, qualityEligible: true, externalExportEligible: true, sensitiveExternalExportAuthorized: true, });
const defaultExamples = await projectTrainingExamples(store); expect(defaultExamples).toEqual([]); const defaultSerialized = JSON.stringify(defaultExamples); for (const sentinel of [fixture.sourceSentinel, fixture.replySentinel, ...fixture.privateMetadataSentinels]) { expect(defaultSerialized).not.toContain(sentinel); }
const privateExamples = await projectTrainingExamples(store, { includeSensitivePrivate: true }); expect(privateExamples).toHaveLength(1); expect(privateExamples[0]?.input.event.privacy).toBe("sensitive"); const privateSerialized = JSON.stringify(privateExamples); expect(privateSerialized).toContain(fixture.replySentinel); expect(privateSerialized).not.toContain(fixture.sourceSentinel); for (const sentinel of fixture.privateMetadataSentinels) expect(privateSerialized).not.toContain(sentinel);
const gitRoot = path.join(root, "public-repository"); await fs.mkdir(path.join(gitRoot, ".git"), { recursive: true }); const refused = path.join(gitRoot, "datasets", "private.jsonl"); await expect(writeTrainingJsonl(refused, privateExamples, { authorizeSensitivePrivateExport: true, publicContentRoots: [], })).rejects.toThrow("inside a Git worktree"); await expect(fs.stat(refused)).rejects.toMatchObject({ code: "ENOENT" }); await expect(fs.stat(`${refused}.manifest.json`)).rejects.toMatchObject({ code: "ENOENT" });
const publicRoot = path.join(root, "published-content"); const publicDestination = path.join(publicRoot, "private.jsonl"); await expect(writeTrainingJsonl(publicDestination, privateExamples, { authorizeSensitivePrivateExport: true, publicContentRoots: [publicRoot], })).rejects.toThrow("inside a public-content root"); await expect(fs.stat(publicDestination)).rejects.toMatchObject({ code: "ENOENT" });
const privateDestination = path.join(root, "private-exports", "private.jsonl"); await expect(writeTrainingJsonl(privateDestination, privateExamples)).rejects.toThrow("requires explicit authorization"); await writeTrainingJsonl(privateDestination, privateExamples, { authorizeSensitivePrivateExport: true, publicContentRoots: [], }); expect((await fs.stat(privateDestination)).mode & 0o777).toBe(0o600); });
test("anchors sensitive dataset and manifest writes when the approved parent path is swapped", async () => { const fixture = await completedRunFixture("sensitive"); await recordJudgment(fixture.store, { runId: fixture.runId, kind: "accept", criterion: "approved-private-export", criterionVersion: 1, qualityEligible: true, externalExportEligible: true, sensitiveExternalExportAuthorized: true, }); const examples = await projectTrainingExamples(fixture.store, { includeSensitivePrivate: true }); const sentinel = fixture.replySentinel; const approvedParent = path.join(fixture.root, "approved-private-parent"); const displaced = path.join(fixture.root, "approved-parent-inode"); const publicRoot = path.join(fixture.root, "synthetic-public-repository"); await fs.mkdir(path.join(publicRoot, ".git"), { recursive: true }); const destination = path.join(approvedParent, "private.jsonl"); let swapped = false;
await expect(writeTrainingJsonl(destination, examples, { authorizeSensitivePrivateExport: true, publicContentRoots: [publicRoot], beforeFinalize: async () => { if (swapped) return; swapped = true; await fs.rename(approvedParent, displaced); await fs.symlink(publicRoot, approvedParent, "dir"); }, })).rejects.toThrow("parent changed during write");
expect(await fs.readdir(publicRoot)).toEqual([".git"]); await expect(fs.stat(path.join(displaced, "private.jsonl"))).rejects.toMatchObject({ code: "ENOENT" }); expect(JSON.stringify(await fs.readdir(publicRoot))).not.toContain(sentinel); });
test("requires explicit authorization when a public trigger produces a sensitive output", async () => { const fixture = await completedRunFixture("public-source", "sensitive"); await expect(recordJudgment(fixture.store, { runId: fixture.runId, kind: "accept", criterion: "mixed-state-output", criterionVersion: 1, qualityEligible: true, externalExportEligible: true, })).rejects.toThrow("requires explicit authorization"); await recordJudgment(fixture.store, { runId: fixture.runId, kind: "accept", criterion: "mixed-state-output", criterionVersion: 1, qualityEligible: true, externalExportEligible: true, sensitiveExternalExportAuthorized: true, }); expect(await projectTrainingExamples(fixture.store)).toEqual([]); expect(await projectTrainingExamples(fixture.store, { includeSensitivePrivate: true })).toHaveLength(1); });
test("requires explicit authority for private state on the compared side of a pair", async () => { const fixture = await completedRunFixture("public-source"); const compared = await appendComparedRun(fixture, "public", "sensitive-pair", "sensitive"); await expect(recordJudgment(fixture.store, { runId: fixture.runId, comparedRunId: compared.runId, kind: "prefer", criterion: "pairwise-sensitive", criterionVersion: 1, qualityEligible: true, externalExportEligible: true, })).rejects.toThrow("requires explicit authorization"); const judgment = await recordJudgment(fixture.store, { runId: fixture.runId, comparedRunId: compared.runId, kind: "prefer", criterion: "pairwise-sensitive", criterionVersion: 1, qualityEligible: true, externalExportEligible: true, sensitiveExternalExportAuthorized: true, }); expect(judgment.privacy).toBe("sensitive"); expect(await projectTrainingExamples(fixture.store)).toEqual([]); const examples = await projectTrainingExamples(fixture.store, { includeSensitivePrivate: true }); expect(examples).toHaveLength(1); expect(examples[0]).toMatchObject({ input: { event: { privacy: "sensitive" } }, comparedProvenance: { modelAdapter: compared.modelAdapter }, }); });
test("rejects pairwise export when the primary run uses a forbidden adapter", async () => { const fixture = await completedRunFixture("public-source"); const compared = await appendComparedRun(fixture, "public", "primary-forbidden-control"); const primary = await fixture.store.getRun(fixture.runId); if (!primary) throw new Error("Missing primary run"); primary.modelAdapter = { ...compared.modelAdapter, id: "primary-forbidden-adapter", manifestSha256: "4".repeat(64), exportClass: "forbidden", binding: { checkpointReferenceSha256: "5".repeat(64) }, }; await fixture.store.upsertRun(primary); await recordJudgment(fixture.store, { runId: fixture.runId, comparedRunId: compared.runId, kind: "prefer", criterion: "pairwise-primary-forbidden", criterionVersion: 1, qualityEligible: true, externalExportEligible: true, }); expect(await projectTrainingExamples(fixture.store, { includeRestrictedModelAdapters: true, })).toEqual([]); });
test("gates both sides of pairwise export and emits exact compared execution and learned-adapter provenance", async () => { const fixture = await completedRunFixture("public-source"); const forbidden = await appendComparedRun(fixture, "forbidden", "forbidden"); await recordJudgment(fixture.store, { runId: fixture.runId, comparedRunId: forbidden.runId, kind: "prefer", criterion: "pairwise-forbidden", criterionVersion: 1, qualityEligible: true, externalExportEligible: true, }); expect(await projectTrainingExamples(fixture.store, { includeRestrictedModelAdapters: true, })).toEqual([]);
const restricted = await appendComparedRun(fixture, "restricted", "restricted"); await recordJudgment(fixture.store, { runId: fixture.runId, comparedRunId: restricted.runId, kind: "prefer", criterion: "pairwise-restricted", criterionVersion: 1, qualityEligible: true, externalExportEligible: true, }); expect(await projectTrainingExamples(fixture.store)).toEqual([]); const [example] = await projectTrainingExamples(fixture.store, { includeRestrictedModelAdapters: true, }); expect(example).toMatchObject({ format: "thoughtstream.training-example.v3", kind: "prefer", provenance: { executionAdapterRevision: "checkpoint-fixture", model: "Qwen/Qwen3.5-4B", }, comparedProvenance: { executionAdapterRevision: "compared-execution-restricted", model: "Qwen/Qwen3.5-35B-A3B-Base", modelAdapter: restricted.modelAdapter, adapterCatalogDigest: "7".repeat(64), adapterCatalogGeneration: 1, }, comparedTrajectory: [{ sequence: 1, type: "pi.compared-restricted" }], }); const manifest = await writeTrainingJsonl(path.join(fixture.root, "pairwise", "training.jsonl"), [example!]); expect(manifest.models).toHaveLength(2); expect(manifest.models).toEqual(expect.arrayContaining([ expect.stringContaining("@execution:checkpoint-fixture"), expect.stringContaining(`@execution:compared-execution-restricted@model-adapter:${restricted.modelAdapter.id}@1#${restricted.modelAdapter.manifestSha256}#${restricted.modelAdapter.binding.checkpointReferenceSha256}@catalog:${"7".repeat(64)}:g1`), ])); });
test("does not reinterpret a legacy private exportEligible judgment as declassification", async () => { const fixture = await completedRunFixture("private"); const legacy = await appendLegacyJudgment(fixture, "private"); expect(legacy.schemaVersion).toBe(1); expect(await projectTrainingExamples(fixture.store)).toEqual([]); expect(await projectTrainingExamples(fixture.store, { includeSensitivePrivate: true })).toEqual([]); });
test("keeps legacy public synthetic judgments export-compatible", async () => { const fixture = await completedRunFixture("public-source"); await appendLegacyJudgment(fixture, "public-source"); const examples = await projectTrainingExamples(fixture.store); expect(examples).toHaveLength(1); expect(examples[0]).toMatchObject({ format: "thoughtstream.training-example.v3", kind: "accept", input: { event: { privacy: "public-source" } }, chosen: fixture.original, }); });});
async function appendLegacyJudgment( fixture: Awaited<ReturnType<typeof completedRunFixture>>, privacy: "private" | "public-source",) { const identity = `legacy-${randomUUID()}`; const result = await fixture.store.appendEvent({ type: "stream.thought.judgment.training-example", schemaVersion: 1, source: "judgment:legacy-fixture", sourceKind: "system", externalId: identity, idempotencyKey: identity, occurredAt: "2026-07-15T00:00:03.000Z", actor: "operator:legacy-fixture", rootEventId: fixture.sourceEventId, parentEventId: fixture.outputEventId, correlationId: fixture.runId, privacy, payload: { runId: fixture.runId, outputEventId: fixture.outputEventId, kind: "accept", criterion: "legacy-export", criterionVersion: 1, exportEligible: true, }, }); return result.event;}
async function appendComparedRun( fixture: Awaited<ReturnType<typeof completedRunFixture>>, exportClass: "public" | "restricted" | "forbidden", label: string, privacy: "public-source" | "private" | "sensitive" = "public-source",) { const primary = await fixture.store.getRun(fixture.runId); if (!primary) throw new Error("Missing primary fixture run"); const outputContract = outputContractForDeclaration({} as ThoughtAgentDeclaration); const outputBody: JsonObject = { summary: `Compared ${label} output`, tags: [label], importance: "normal", confidence: 0.4, }; const output = await fixture.store.appendEvent({ type: "stream.thought.derived.topics", schemaVersion: 1, source: `agent:compared-${label}`, sourceKind: "agent", externalId: `compared-output-${label}`, idempotencyKey: `compared-output-${label}`, occurredAt: "2026-07-15T00:00:01.000Z", actor: `compared-${label}`, rootEventId: fixture.sourceEventId, parentEventId: fixture.sourceEventId, correlationId: `compared-${label}`, privacy, payload: { outputContract: outputContractIdentityJson(outputContract), structuredOutput: outputBody, }, }); const release = { id: `compared-${label}-adapter`, version: 1, description: `Compared ${label} adapter`, releasedAt: "2026-07-25T20:00:00.000Z", manifestSha256: label === "forbidden" ? "d".repeat(64) : "e".repeat(64), checkpointSelector: { kind: "env" as const, reference: "THOUGHTSTREAM_TEST_ADAPTER_MODEL" }, providerProfile: "tinker-default", baseModel: "Qwen/Qwen3.5-35B-A3B-Base", dataset: { id: `compared-${label}-dataset`, sha256: "f".repeat(64) }, evals: [{ id: `compared-${label}-eval`, sha256: "1".repeat(64) }], capabilities: [`compared-${label}`], privacyClass: "public" as const, exportClass, }; const modelAdapter = { ...release, binding: { checkpointReferenceSha256: label === "forbidden" ? "2".repeat(64) : "3".repeat(64) }, }; const runId = `run-compared-${label}-${randomUUID()}`; await fixture.store.upsertRun({ id: runId, executionKey: `execution-compared-${label}`, triggerEventId: primary.triggerEventId, agentId: `compared-${label}`, agentVersion: 1, status: "completed", inputEventIds: [primary.triggerEventId], outputEventIds: [output.event.id], attempt: 1, provider: "tinker", model: release.baseModel, privacy, executionAdapterRevision: `compared-execution-${label}`, modelAdapter, adapterCatalogDigest: label === "forbidden" ? "6".repeat(64) : "7".repeat(64), adapterCatalogGeneration: 1, promptHash: `prompt-${label}`, contextManifest: { agentRole: "standard", outputContract: outputContractIdentityJson(outputContract), contextStrategy: "single-event", }, result: outputBody, createdAt: "2026-07-15T00:00:00.000Z", completedAt: "2026-07-15T00:00:02.000Z", updatedAt: "2026-07-15T00:00:02.000Z", }); await fixture.store.appendTrace({ id: `trace-compared-${label}`, runId, sequence: 1, type: `pi.compared-${label}`, payload: {}, createdAt: "2026-07-15T00:00:01.500Z", }); return { runId, modelAdapter };}
async function completedRunFixture( privacy: "private" | "sensitive" | "public-source", outputPrivacy: "private" | "sensitive" | "public-source" = privacy,) { const root = await temporaryProject(); roots.push(root); const store = await testStore(root); stores.push(store); const sourceSentinel = `source-${randomUUID()}`; const replySentinel = `reply-${randomUUID()}`; const credentialSentinel = `credential-${randomUUID()}`; const routeSentinel = `route-${randomUUID()}`; const externalSentinel = `external-${randomUUID()}`; const correlationSentinel = `correlation-${randomUUID()}`; const source = await store.appendEvent({ type: "stream.thought.runtime.notice", schemaVersion: 1, source: routeSentinel, sourceKind: "system", externalId: externalSentinel, idempotencyKey: `idempotency-${randomUUID()}`, occurredAt: "2026-07-15T00:00:00.000Z", actor: `actor-${randomUUID()}`, correlationId: correlationSentinel, privacy, payload: { text: sourceSentinel, credential: credentialSentinel }, createdByRuntime: "test", }); const outputContract = outputContractForDeclaration({} as ThoughtAgentDeclaration); const original: JsonObject = { summary: privacy === "public-source" ? "Original public output" : replySentinel, tags: ["wrong"], importance: "normal", confidence: 0.6, }; const output = await store.appendEvent({ type: "stream.thought.derived.topics", schemaVersion: 1, source: `agent-${randomUUID()}`, sourceKind: "agent", externalId: `output-${randomUUID()}`, idempotencyKey: `output-key-${randomUUID()}`, occurredAt: "2026-07-15T00:00:01.000Z", actor: `output-actor-${randomUUID()}`, rootEventId: source.event.id, parentEventId: source.event.id, correlationId: `run-correlation-${randomUUID()}`, privacy: outputPrivacy, payload: { outputContract: outputContractIdentityJson(outputContract), structuredOutput: original, }, createdByRuntime: "test", }); const runId = `run-${randomUUID()}`; await store.upsertRun({ id: runId, executionKey: `execution-${randomUUID()}`, triggerEventId: source.event.id, agentId: `agent-${randomUUID()}`, agentVersion: 3, status: "completed", inputEventIds: [source.event.id], outputEventIds: [output.event.id], attempt: 1, provider: "tinker", model: "Qwen/Qwen3.5-4B", executionAdapterRevision: "checkpoint-fixture", promptHash: "prompt-hash", contextManifest: { eventIds: [source.event.id], totalChars: 42, agentRole: "standard", outputContract: outputContractIdentityJson(outputContract), credential: credentialSentinel, route: routeSentinel, contextStrategy: "single-event", }, result: original, createdAt: "2026-07-15T00:00:00.000Z", completedAt: "2026-07-15T00:00:02.000Z", updatedAt: "2026-07-15T00:00:02.000Z", }); await store.appendTrace({ id: `trace-${randomUUID()}`, runId, sequence: 1, type: "pi.message_end", payload: { credential: credentialSentinel, source: sourceSentinel, route: routeSentinel }, createdAt: "2026-07-15T00:00:01.500Z", }); return { root, store, runId, sourceEventId: source.event.id, outputEventId: output.event.id, original, sourceSentinel, replySentinel, privateMetadataSentinels: [credentialSentinel, routeSentinel, externalSentinel, correlationSentinel], };}