Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
14 kB · 327 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328import fs from "node:fs/promises";import path from "node:path";import { afterEach, describe, expect, test } from "vitest";import { loadAgentDeclarations } from "../src/agents/declarations.js";import { ThoughtAgentRuntime } from "../src/agents/runtime.js";import type { AgentRunner } from "../src/agents/types.js";import { canonicalJson } from "../src/core/json.js";import { MACHINE_PROJECT_EULER_REVISION, POST_TRAINING_COURSE, postTrainingLessonContext,} from "../src/courses/post-training.js";import { COURSE_QUESTION_EVENT_TYPE, COURSE_QUESTION_SOURCE, COURSE_TUTOR_AGENT_ID, CourseQuestionConflictError, CourseRevisionConflictError, appendPostTrainingCourseQuestion, parseCourseQuestionEvent,} from "../src/courses/questions.js";import { COURSE_CHAT_NONCE_HEADER, COURSE_CHAT_SIGNATURE_HEADER, COURSE_CHAT_TIMESTAMP_HEADER, signCourseChatRequest,} from "../src/courses/web-capability.js";import { startInspectorServer } from "../src/web/inspector.js";import { eventIdFor, type JazzThoughtStore } from "../src/jazz/store.js";import { temporaryProject, testDeclarationEnvironment, testStore } from "./helpers.js";
const roots: string[] = [];const stores: JazzThoughtStore[] = [];const servers: import("node:http").Server[] = [];
afterEach(async () => { await Promise.all(servers.splice(0).map((server) => new Promise<void>((resolve) => { server.closeAllConnections(); server.close(() => resolve()); }))); 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("post-training course questions", () => { test("binds eight lessons and model context to one content revision", () => { expect(POST_TRAINING_COURSE.revision).toMatch(/^[a-f0-9]{64}$/); expect(POST_TRAINING_COURSE.lessons).toHaveLength(8); expect(POST_TRAINING_COURSE.lessons.map((lesson) => lesson.number)).toEqual([1, 2, 3, 4, 5, 6, 7, 8]); expect(POST_TRAINING_COURSE.lessons.map((lesson) => lesson.workshop.kind)).toEqual([ "capability-contract", "machine-gate", "data-split", "environment", "loss-mask", "preference", "reward", "factory", ]); const lessonIds = new Set<string>(); for (const lesson of POST_TRAINING_COURSE.lessons) { expect(lessonIds.has(lesson.id)).toBe(false); lessonIds.add(lesson.id); expect(lesson.objective.length).toBeGreaterThan(20); expect(lesson.sections.length).toBeGreaterThan(0); expect(new Set(lesson.sections.map((section) => section.id)).size).toBe(lesson.sections.length); expect(lesson.references.length).toBeGreaterThan(0); expect(lesson.references.every((reference) => reference.url.startsWith("https://"))).toBe(true); expect(postTrainingLessonContext(lesson.id).length).toBeLessThanOrEqual(48_000); } const machineLesson = POST_TRAINING_COURSE.lessons.find((lesson) => lesson.id === "machine-project-euler")!; expect(machineLesson.references[0]?.url).toContain(`/blob/${MACHINE_PROJECT_EULER_REVISION}/`); expect(JSON.stringify(machineLesson)).toContain("17 of 24"); expect(JSON.stringify(machineLesson)).toContain("0 of 6 in the full Pi harness with zero tool calls"); expect(JSON.stringify(machineLesson)).toContain("Project Euler problem 2 remains locked"); const context = postTrainingLessonContext("machine-project-euler", "promotion-gate"); expect(context).toContain("at least six of eight trajectories in every language for two consecutive rounds"); expect(context).toContain(`Course revision: ${POST_TRAINING_COURSE.revision}`); expect(context).not.toContain("The target is tool-shaped problem solving"); });
test("appends one exact sensitive web event and rejects stale or divergent reuse", async () => { const { store } = await fixture(); const input = { requestId: "question-0001", courseRevision: POST_TRAINING_COURSE.revision, lessonId: "capability-hello-world", sectionId: "smallest-loop", question: "Why does exact output alone fail to prove tool use?", }; const first = await appendPostTrainingCourseQuestion(store, input, { occurredAt: "2026-08-22T08:00:00.000Z", }); const second = await appendPostTrainingCourseQuestion(store, input, { occurredAt: "2026-08-22T08:01:00.000Z", }); expect(first.inserted).toBe(true); expect(second.inserted).toBe(false); expect(second.event.id).toBe(first.event.id); expect(first.event).toMatchObject({ type: COURSE_QUESTION_EVENT_TYPE, source: COURSE_QUESTION_SOURCE, sourceKind: "web", actor: "operator:cameron", privacy: "sensitive", rootEventId: first.event.id, }); expect(parseCourseQuestionEvent(first.event)).toMatchObject({ requestId: input.requestId, lessonId: input.lessonId, sectionId: input.sectionId, question: input.question, interface: "thought-stream-inspector", }); await expect(appendPostTrainingCourseQuestion(store, { ...input, question: "Different" })) .rejects.toBeInstanceOf(CourseQuestionConflictError); await expect(appendPostTrainingCourseQuestion(store, { ...input, requestId: "question-0002", courseRevision: "0".repeat(64) })) .rejects.toBeInstanceOf(CourseRevisionConflictError); const serialized = canonicalJson(first.event.payload); expect(serialized).not.toContain("credential"); expect(serialized).not.toContain("provider"); });
test("serves course content and accepts only a signed question mutation", async () => { const { store } = await fixture(); const capability = Buffer.alloc(32, 41); const server = await startInspectorServer(store, { port: 0, courseChatCapability: capability }); servers.push(server); const base = serverBase(server); const courseResponse = await fetch(`${base}/api/courses/post-training`); expect(courseResponse.status).toBe(200); expect(await courseResponse.json()).toMatchObject({ id: "post-training-model-factory", revision: POST_TRAINING_COURSE.revision, lessons: expect.arrayContaining([expect.objectContaining({ id: "model-factory" })]), }); const body = Buffer.from(JSON.stringify({ requestId: "question-web-0001", courseRevision: POST_TRAINING_COURSE.revision, lessonId: "model-factory", question: "Why are rollout workers separate from learners?", })); expect((await fetch(`${base}/api/courses/post-training/questions`, { method: "POST", headers: { "content-type": "application/json" }, body, })).status).toBe(403); const path = "/api/courses/post-training/questions"; const signed = signCourseChatRequest(capability, { method: "POST", path, body }); const accepted = await fetch(`${base}${path}`, { method: "POST", headers: { "content-type": "application/json", [COURSE_CHAT_TIMESTAMP_HEADER]: signed.timestamp, [COURSE_CHAT_NONCE_HEADER]: signed.nonce, [COURSE_CHAT_SIGNATURE_HEADER]: signed.signature, }, body, }); expect(accepted.status).toBe(202); const receipt = await accepted.json() as { eventId: string; statusPath: string }; expect(receipt.statusPath).toBe(`api/courses/post-training/questions/${encodeURIComponent(receipt.eventId)}`); expect(await (await fetch(`${base}/${receipt.statusPath}`)).json()).toEqual({ status: "pending", eventId: receipt.eventId, }); expect((await fetch(`${base}${path}`, { method: "POST", headers: { "content-type": "application/json", [COURSE_CHAT_TIMESTAMP_HEADER]: signed.timestamp, [COURSE_CHAT_NONCE_HEADER]: signed.nonce, [COURSE_CHAT_SIGNATURE_HEADER]: signed.signature, }, body, })).status).toBe(403); });
test("returns an answer only through exact completed tutor lineage", async () => { const { store } = await fixture(); const question = await appendPostTrainingCourseQuestion(store, { requestId: "question-answer-0001", courseRevision: POST_TRAINING_COURSE.revision, lessonId: "supervised-fine-tuning", question: "Why does teacher forcing hide rollout errors?", }); const answer = "Teacher forcing supplies the correct previous token, so the model never has to recover from its own earlier mistake during the supervised update."; const runId = "run_course_tutor_0001"; const outputId = eventIdFor(`agent:${COURSE_TUTOR_AGENT_ID}`, `${runId}:output`); const run = { id: runId, executionKey: "execution_course_tutor_0001", triggerEventId: question.event.id, agentId: COURSE_TUTOR_AGENT_ID, agentVersion: 1, status: "completed" as const, inputEventIds: [question.event.id], outputEventIds: [outputId], attempt: 1, provider: "tinker", model: "fixture-model", privacy: "sensitive" as const, promptHash: "fixture-prompt", contextManifest: {}, result: { summary: answer, tags: ["post-training-course"], importance: "normal" as const, confidence: 0.95 }, createdAt: "2026-08-22T08:00:01.000Z", completedAt: "2026-08-22T08:00:02.000Z", updatedAt: "2026-08-22T08:00:02.000Z", }; await store.upsertRun(run); const output = await store.appendEvent({ type: "stream.thought.derived.message.observation", schemaVersion: 1, source: `agent:${COURSE_TUTOR_AGENT_ID}`, sourceKind: "agent", externalId: runId, idempotencyKey: `${runId}:output`, occurredAt: run.completedAt, actor: COURSE_TUTOR_AGENT_ID, rootEventId: question.event.id, parentEventId: question.event.id, correlationId: question.event.correlationId, privacy: "sensitive", traceId: runId, payload: { runId, executionKey: run.executionKey, inputEventId: question.event.id, inputSourceSequence: question.event.sourceSequence, summary: answer, tags: ["post-training-course"], importance: "normal", confidence: 0.95, outputContract: { id: "stream.thought.output.observation", version: 1, sha256: "f".repeat(64) }, structuredOutput: { summary: answer, tags: ["post-training-course"], importance: "normal", confidence: 0.95 }, }, }); expect(output.event.id).toBe(outputId); const server = await startInspectorServer(store, { port: 0 }); servers.push(server); expect(await (await fetch(`${serverBase(server)}/api/courses/post-training/questions/${question.event.id}`)).json()).toEqual({ status: "completed", eventId: question.event.id, runId, outputEventId: outputId, answer, }); });
test("runs the declared tutor from a live course subscription and closes exact answer lineage", async () => { const { store } = await fixture(); const declaration = (await loadAgentDeclarations( path.join(process.cwd(), "agents"), testDeclarationEnvironment, )).find((candidate) => candidate.id === COURSE_TUTOR_AGENT_ID)!; let modelContext = ""; const answer = "Teacher forcing supplies the correct previous token, so rollout evaluation must test recovery from the model's own earlier actions."; const runner: AgentRunner = { mode: "pi", run: async ({ context }) => { modelContext = context.text; return { summary: answer, tags: ["post-training-course", "sft"], importance: "normal", confidence: 0.94, model: { provider: "tinker", id: declaration.model! }, usage: { inputTokens: 480, outputTokens: 36 }, }; }, }; const consumers = await new ThoughtAgentRuntime(store, [runner], { reconcileIntervalMs: 20 }) .startConsumers([declaration]); const question = await appendPostTrainingCourseQuestion(store, { requestId: "question-runtime-0001", courseRevision: POST_TRAINING_COURSE.revision, lessonId: "supervised-fine-tuning", sectionId: "teacher-forcing", question: "What does teacher forcing fail to demonstrate?", }); await waitFor(async () => (await store.listRuns()).some((run) => ( run.agentId === COURSE_TUTOR_AGENT_ID && run.status === "completed" ))); await consumers.stop();
expect(modelContext).toContain("What does teacher forcing fail to demonstrate?"); expect(modelContext).toContain("Teacher forcing hides rollout errors"); expect(modelContext).not.toContain("The target is tool-shaped problem solving"); const run = (await store.listRuns()).find((candidate) => candidate.agentId === COURSE_TUTOR_AGENT_ID)!; expect(run).toMatchObject({ triggerEventId: question.event.id, status: "completed", result: { summary: answer }, }); const server = await startInspectorServer(store, { port: 0 }); servers.push(server); expect(await (await fetch(`${serverBase(server)}/api/courses/post-training/questions/${question.event.id}`)).json()).toMatchObject({ status: "completed", eventId: question.event.id, runId: run.id, answer, }); });});
async function fixture(): Promise<{ store: JazzThoughtStore }> { const root = await temporaryProject("thoughtstream-course-question-"); roots.push(root); const store = await testStore(root); stores.push(store); return { store };}
function serverBase(server: import("node:http").Server): string { const address = server.address(); if (!address || typeof address === "string") throw new Error("Missing server address"); return `http://127.0.0.1:${address.port}`;}
async function waitFor(predicate: () => Promise<boolean>, timeoutMs = 5_000): Promise<void> { const deadline = Date.now() + timeoutMs; while (Date.now() < deadline) { if (await predicate()) return; await new Promise((resolve) => setTimeout(resolve, 20)); } throw new Error("Timed out waiting for course tutor settlement");}