import 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; 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>, 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>, 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], }; }