Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
22 kB · 358 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359import fs from "node:fs/promises";import { afterEach, describe, expect, test } from "vitest";import { sha256 } from "../src/core/json.js";import type { ThoughtEvent } from "../src/events/types.js";import { CONTEXT_SELECTED_EVENT_TYPE, DOCUMENT_CREATED_EVENT_TYPE, DOCUMENT_VERSION_EVENT_TYPE, PROPOSAL_DECISION_EVENT_TYPE, PROPOSAL_PROPOSED_EVENT_TYPE, PROPOSAL_REQUESTED_EVENT_TYPE, WORKBENCH_CONTEXT_SOURCE, WORKBENCH_DOCUMENT_SOURCE, decodeDocumentContent, encodeDocumentContent, plainExcerpt,} from "../src/workbench/contracts.js";import { FIXTURE_RUNNER_ID, fixtureDeterministicRunner } from "../src/workbench/runners.js";import { StaleBaseError, WorkbenchConflictError, WorkbenchRunnerError, WorkbenchWorkflow, workingDocumentId,} from "../src/workbench/workflow.js";import { JazzThoughtStore } from "../src/jazz/store.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 })));});
async function openStore(appId?: string): Promise<{ root: string; store: JazzThoughtStore }> { const projectRoot = await temporaryProject("thoughtstream-workbench-"); roots.push(projectRoot); const store = appId ? await JazzThoughtStore.open({ projectRoot, appId, runtimeRevision: "test" }) : await testStore(projectRoot); stores.push(store); return { root: projectRoot, store };}
async function seedOrigin(store: JazzThoughtStore, suffix = "1"): Promise<ThoughtEvent> { return (await store.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:fixture", sourceKind: "rss", externalId: `item-${suffix}`, idempotencyKey: `rss:item-${suffix}`, occurredAt: `2026-07-14T00:0${suffix}:00.000Z`, actor: "rss:fixture", correlationId: "poll-1", privacy: "public-source", payload: { title: `Fixture item ${suffix}`, summary: `Synthetic summary ${suffix} with <b>markup</b>` }, })).event;}
async function eventsOfType(store: JazzThoughtStore, type: string): Promise<ThoughtEvent[]> { return store.listEvents({ types: [type], source: WORKBENCH_DOCUMENT_SOURCE });}
async function judgments(store: JazzThoughtStore): Promise<ThoughtEvent[]> { return store.listEvents({ types: ["stream.thought.judgment.training-example"] });}
describe("workbench document workflow", () => { test("frontmatter encoding is self-describing and reversible", () => { const content = encodeDocumentContent('A "quoted" title', "# Body\n\ntext"); expect(content.startsWith("---\ntitle: ")).toBe(true); expect(decodeDocumentContent(content)).toEqual({ title: 'A "quoted" title', body: "# Body\n\ntext" }); expect(decodeDocumentContent("no frontmatter")).toEqual({ title: "Untitled", body: "no frontmatter" }); expect(plainExcerpt({ title: "T", text: "x".repeat(5_000) }).length).toBe(2_000); expect(plainExcerpt({ nested: { deep: true } })).toBe('{"nested":{"deep":true}}'); });
test("creates a working document whose origin reference survives store reopen", async () => { const { root, store } = await openStore(); const origin = await seedOrigin(store); const workflow = new WorkbenchWorkflow(store, { actor: "operator:test" }); const head = await workflow.createWorkingDocument({ originEventId: origin.id, title: "Notes on fixture item", body: "# Notes\n\nFirst draft.\n", requestId: "req-create-0001", }); expect(head.documentId).toBe(workingDocumentId(origin.id, "req-create-0001")); const replay = await workflow.createWorkingDocument({ originEventId: origin.id, title: "Notes on fixture item", body: "# Notes\n\nFirst draft.\n", requestId: "req-create-0001", }); expect(replay.versionId).toBe(head.versionId); await expect(workflow.createWorkingDocument({ originEventId: origin.id, title: "Different", body: "content", requestId: "req-create-0001", })).rejects.toBeInstanceOf(WorkbenchConflictError); await expect(workflow.createWorkingDocument({ originEventId: "evt_missing", title: "x", body: "y", requestId: "req-create-0002", })).rejects.toThrow("Origin event was not found");
const created = await eventsOfType(store, DOCUMENT_CREATED_EVENT_TYPE); expect(created).toHaveLength(1); expect(created[0]!.rootEventId).toBe(origin.rootEventId); expect(created[0]!.parentEventId).toBe(origin.id); expect(created[0]!.privacy).toBe("sensitive"); expect(await eventsOfType(store, DOCUMENT_VERSION_EVENT_TYPE)).toHaveLength(1);
await store.close(); stores.splice(stores.indexOf(store), 1); const reopened = await testStore(root); stores.push(reopened); const detail = await new WorkbenchWorkflow(reopened).getDocument(head.documentId); expect(detail?.origin?.eventId).toBe(origin.id); expect(detail?.head.title).toBe("Notes on fixture item"); expect(detail?.head.body).toBe("# Notes\n\nFirst draft.\n"); expect(detail?.head.versionId).toBe(head.versionId); const version = await reopened.getDocumentVersion(head.versionId); expect(version?.source).toBe(WORKBENCH_DOCUMENT_SOURCE); expect(version?.path).toBe(`documents/${encodeURIComponent(head.documentId)}.md`); expect(version?.contentType).toBe("text/markdown"); const list = await new WorkbenchWorkflow(reopened).listDocuments(); expect(list).toHaveLength(1); expect(list[0]).toMatchObject({ documentId: head.documentId, originEventId: origin.id, versionCount: 1, proposalCount: 0 }); });
test("operator edits require the exact head and are idempotent per request", async () => { const { store } = await openStore(); const origin = await seedOrigin(store); const workflow = new WorkbenchWorkflow(store); const head = await workflow.createWorkingDocument({ originEventId: origin.id, title: "T", body: "one\n", requestId: "req-create-0001" }); const edited = await workflow.saveOperatorEdit({ documentId: head.documentId, baseVersionId: head.versionId, title: "T2", body: "two\n", requestId: "req-edit-00001" }); expect(edited.versionId).not.toBe(head.versionId); const stale = workflow.saveOperatorEdit({ documentId: head.documentId, baseVersionId: head.versionId, title: "T3", body: "three\n", requestId: "req-edit-00002" }); await expect(stale).rejects.toBeInstanceOf(StaleBaseError); await expect(stale).rejects.toMatchObject({ headVersionId: edited.versionId }); const replay = await workflow.saveOperatorEdit({ documentId: head.documentId, baseVersionId: head.versionId, title: "T2", body: "two\n", requestId: "req-edit-00001" }); expect(replay.versionId).toBe(edited.versionId); const noop = await workflow.saveOperatorEdit({ documentId: head.documentId, baseVersionId: edited.versionId, title: "T2", body: "two\n", requestId: "req-edit-00003" }); expect(noop.versionId).toBe(edited.versionId); expect(await eventsOfType(store, DOCUMENT_VERSION_EVENT_TYPE)).toHaveLength(2); const detail = await workflow.getDocument(head.documentId); expect(detail?.head.title).toBe("T2"); expect(detail?.versions.map((version) => version.reason)).toEqual(["created", "operator-edit"]); });
test("selects exact context and stores a bounded plain-text snapshot", async () => { const { store } = await openStore(); const origin = await seedOrigin(store, "1"); const other = await seedOrigin(store, "2"); const workflow = new WorkbenchWorkflow(store); const head = await workflow.createWorkingDocument({ originEventId: origin.id, title: "T", body: "body\n", requestId: "req-create-0001" }); await expect(workflow.selectContext({ documentId: head.documentId, eventIds: ["evt_missing"], versionIds: [], requestId: "req-select-001" })) .rejects.toThrow("Selected event was not found"); const selection = await workflow.selectContext({ documentId: head.documentId, eventIds: [other.id, origin.id], versionIds: [head.versionId], requestId: "req-select-002", }); expect(selection.baseVersionId).toBe(head.versionId); expect(selection.selectedEvents.map((event) => event.eventId)).toEqual([other.id, origin.id]); expect(selection.selectedEvents[0]).toMatchObject({ type: other.type, source: other.source, payloadHash: other.payloadHash }); expect(selection.selectedEvents[0]!.excerpt).toContain("Synthetic summary 2 with <b>markup</b>"); expect(selection.selectedVersions).toEqual([{ documentId: head.documentId, versionId: head.versionId, sha256: head.sha256, path: `documents/${encodeURIComponent(head.documentId)}.md` }]); const snapshot = await store.getDocumentVersion(selection.snapshotVersionId); expect(snapshot?.source).toBe(WORKBENCH_CONTEXT_SOURCE); expect(snapshot?.contentType).toBe("application/json"); expect(snapshot?.sha256).toBe(sha256(snapshot!.content)); expect(JSON.parse(snapshot!.content)).toMatchObject({ documentId: head.documentId, baseVersionId: head.versionId }); const replay = await workflow.selectContext({ documentId: head.documentId, eventIds: [other.id, origin.id], versionIds: [head.versionId], requestId: "req-select-002" }); expect(replay.eventId).toBe(selection.eventId); expect(await eventsOfType(store, CONTEXT_SELECTED_EVENT_TYPE)).toHaveLength(1); const candidates = await workflow.listCandidates(head.documentId); expect(candidates?.events.some((event) => event.eventId === origin.id && event.origin)).toBe(true); expect(candidates?.events.some((event) => event.eventId === other.id)).toBe(true); expect(candidates?.versions).toEqual([{ versionId: head.versionId, sha256: head.sha256, reason: "created", createdAt: expect.any(String), head: true }]); });
test("deterministic fixture proposal, diff, accept once, duplicate submission, judgment lineage", async () => { const appId = "thoughtstream-workbench-attach-test"; const { root, store } = await openStore(appId); const origin = await seedOrigin(store); const workflow = new WorkbenchWorkflow(store, { actor: "operator:test" }); const head = await workflow.createWorkingDocument({ originEventId: origin.id, title: "Doc", body: "# Doc\n\nIntro.\n", requestId: "req-create-0001" }); const selection = await workflow.selectContext({ documentId: head.documentId, eventIds: [origin.id], versionIds: [head.versionId], requestId: "req-select-001" }); await expect(workflow.requestProposal({ documentId: head.documentId, selectionId: selection.selectionId, runnerId: "missing", requestId: "req-prop-00001" })) .rejects.toThrow("Unknown proposal runner"); const first = await workflow.requestProposal({ documentId: head.documentId, selectionId: selection.selectionId, runnerId: FIXTURE_RUNNER_ID, requestId: "req-prop-00001" }); const replay = await workflow.requestProposal({ documentId: head.documentId, selectionId: selection.selectionId, runnerId: FIXTURE_RUNNER_ID, requestId: "req-prop-00001" }); expect(replay.proposalEventId).toBe(first.proposalEventId); expect(await eventsOfType(store, PROPOSAL_REQUESTED_EVENT_TYPE)).toHaveLength(1); expect(await eventsOfType(store, PROPOSAL_PROPOSED_EVENT_TYPE)).toHaveLength(1); const proposal = (await store.getEvent(first.proposalEventId))!; expect(proposal.parentEventId).toBe(first.requestedEventId); expect(proposal.privacy).toBe("sensitive"); expect(proposal.payload).toMatchObject({ proposalState: "runner-proposed", operation: "replace-document", publicationEligible: false, reason: "fixture: append sources section", evidenceEventIds: [origin.id], target: { documentId: head.documentId, baseVersionId: head.versionId, baseSha256: head.sha256 }, }); const run = (await store.getRun(first.runId))!; expect(run.status).toBe("completed"); expect(run.triggerEventId).toBe(first.requestedEventId); expect(run.outputEventIds).toEqual([first.proposalEventId]); expect(run.provider).toBe("fixture"); expect(run.contextManifest.contextSnapshot).toMatchObject({ id: selection.snapshotVersionId, storage: "jazz-document-version" });
const direct = await fixtureDeterministicRunner.propose({ documentId: head.documentId, baseVersionId: head.versionId, baseSha256: head.sha256, title: "Doc", baseText: "# Doc\n\nIntro.\n", snapshot: { documentId: head.documentId, baseVersionId: head.versionId, selectedEvents: selection.selectedEvents, selectedVersions: selection.selectedVersions }, maxProposedChars: 64_000, }); expect(direct.proposedText).toBe(proposal.payload.proposedText); expect(direct.proposedText).toContain("## Sources"); expect(direct.proposedText).toContain("Summary: this document cites 2 selected sources.");
const detailBefore = (await workflow.getDocument(head.documentId))!; expect(detailBefore.proposals).toHaveLength(1); expect(detailBefore.proposals[0]).toMatchObject({ status: "pending", stale: false, headVersionId: head.versionId, inference: "none" }); expect(detailBefore.proposals[0]!.diff).toContain("+## Sources"); expect(detailBefore.proposals[0]!.diff).toContain(`base ${head.versionId}`);
const accepted = await workflow.decideProposal({ proposalEventId: first.proposalEventId, disposition: "accept", submissionId: "sub-accept-0001" }); expect(accepted.replayed).toBe(false); expect(accepted.resultVersionId).toBeDefined(); expect(accepted.judgmentEventId).toBeDefined(); const duplicate = await workflow.decideProposal({ proposalEventId: first.proposalEventId, disposition: "accept", submissionId: "sub-accept-0001" }); expect(duplicate).toMatchObject({ decisionEventId: accepted.decisionEventId, resultVersionId: accepted.resultVersionId, judgmentEventId: accepted.judgmentEventId, replayed: true }); await expect(workflow.decideProposal({ proposalEventId: first.proposalEventId, disposition: "reject", submissionId: "sub-accept-0001" })) .rejects.toThrow("submission id conflicts"); await expect(workflow.decideProposal({ proposalEventId: first.proposalEventId, disposition: "reject", submissionId: "sub-reject-0002" })) .rejects.toThrow("already has a human decision");
expect(await eventsOfType(store, DOCUMENT_VERSION_EVENT_TYPE)).toHaveLength(2); expect(await eventsOfType(store, PROPOSAL_DECISION_EVENT_TYPE)).toHaveLength(1); const allJudgments = await judgments(store); expect(allJudgments).toHaveLength(1); expect(allJudgments[0]!.payload).toMatchObject({ runId: first.runId, outputEventId: first.proposalEventId, kind: "accept", criterion: "workbench-document-proposal", qualityEligible: false, externalExportEligible: false, feedbackSourceEventId: accepted.decisionEventId, }); expect(allJudgments[0]!.parentEventId).toBe(accepted.decisionEventId); expect(allJudgments[0]!.rootEventId).toBe(origin.rootEventId);
const detail = (await workflow.getDocument(head.documentId))!; expect(detail.head.versionId).toBe(accepted.resultVersionId); expect(detail.head.title).toBe("Doc"); expect(detail.head.body).toBe(direct.proposedText); expect(detail.proposals[0]).toMatchObject({ status: "accepted", decision: { disposition: "accept", resultVersionId: accepted.resultVersionId, judgmentEventId: accepted.judgmentEventId } }); expect(detail.versions.at(-1)).toMatchObject({ reason: "proposal-accepted", proposalEventId: first.proposalEventId, decisionEventId: accepted.decisionEventId, baseVersionId: head.versionId });
// A second independent client over the same project root sees the accepted head. const second = await JazzThoughtStore.open({ projectRoot: root, appId, runtimeRevision: "test" }); stores.push(second); const seen = await new WorkbenchWorkflow(second).getDocument(head.documentId); expect(seen?.head.versionId).toBe(accepted.resultVersionId); expect(seen?.head.body).toBe(direct.proposedText); expect((await second.listCurrentDocuments(WORKBENCH_DOCUMENT_SOURCE))[0]?.versionId).toBe(accepted.resultVersionId); });
test("reject records a decision and judgment and leaves the document unchanged", async () => { const { store } = await openStore(); const origin = await seedOrigin(store); const workflow = new WorkbenchWorkflow(store); const head = await workflow.createWorkingDocument({ originEventId: origin.id, title: "Doc", body: "body\n", requestId: "req-create-0001" }); const selection = await workflow.selectContext({ documentId: head.documentId, eventIds: [origin.id], versionIds: [], requestId: "req-select-001" }); const proposal = await workflow.requestProposal({ documentId: head.documentId, selectionId: selection.selectionId, runnerId: FIXTURE_RUNNER_ID, requestId: "req-prop-00001" }); const rejected = await workflow.decideProposal({ proposalEventId: proposal.proposalEventId, disposition: "reject", submissionId: "sub-reject-0001" }); expect(rejected.resultVersionId).toBeUndefined(); expect(rejected.judgmentEventId).toBeDefined(); const detail = (await workflow.getDocument(head.documentId))!; expect(detail.head.versionId).toBe(head.versionId); expect(detail.head.body).toBe("body\n"); expect(detail.versions).toHaveLength(1); expect(detail.proposals[0]).toMatchObject({ status: "rejected", decision: { disposition: "reject" } }); expect((await judgments(store))[0]!.payload).toMatchObject({ kind: "reject", qualityEligible: false, externalExportEligible: false }); const replay = await workflow.decideProposal({ proposalEventId: proposal.proposalEventId, disposition: "reject", submissionId: "sub-reject-0001" }); expect(replay).toMatchObject({ decisionEventId: rejected.decisionEventId, replayed: true }); expect(await eventsOfType(store, PROPOSAL_DECISION_EVENT_TYPE)).toHaveLength(1); expect(await judgments(store)).toHaveLength(1); });
test("an operator edit makes an older proposal stale; accept fails closed and the edit survives", async () => { const { store } = await openStore(); const origin = await seedOrigin(store); const workflow = new WorkbenchWorkflow(store); const head = await workflow.createWorkingDocument({ originEventId: origin.id, title: "Doc", body: "body\n", requestId: "req-create-0001" }); const selection = await workflow.selectContext({ documentId: head.documentId, eventIds: [origin.id], versionIds: [], requestId: "req-select-001" }); const proposal = await workflow.requestProposal({ documentId: head.documentId, selectionId: selection.selectionId, runnerId: FIXTURE_RUNNER_ID, requestId: "req-prop-00001" }); const edited = await workflow.saveOperatorEdit({ documentId: head.documentId, baseVersionId: head.versionId, title: "Doc", body: "newer operator work\n", requestId: "req-edit-00001" }); const stale = workflow.decideProposal({ proposalEventId: proposal.proposalEventId, disposition: "accept", submissionId: "sub-accept-0001" }); await expect(stale).rejects.toBeInstanceOf(StaleBaseError); await expect(stale).rejects.toMatchObject({ headVersionId: edited.versionId }); const detail = (await workflow.getDocument(head.documentId))!; expect(detail.head.versionId).toBe(edited.versionId); expect(detail.head.body).toBe("newer operator work\n"); expect(detail.proposals[0]).toMatchObject({ status: "stale", stale: true, headVersionId: edited.versionId }); expect(detail.proposals[0]!.decision).toBeUndefined(); expect(await eventsOfType(store, PROPOSAL_DECISION_EVENT_TYPE)).toHaveLength(0); expect(await eventsOfType(store, DOCUMENT_VERSION_EVENT_TYPE)).toHaveLength(2); expect(await judgments(store)).toHaveLength(0); // Rejecting a stale proposal remains allowed. const rejected = await workflow.decideProposal({ proposalEventId: proposal.proposalEventId, disposition: "reject", submissionId: "sub-reject-0001" }); expect(rejected.disposition).toBe("reject"); expect((await workflow.getDocument(head.documentId))!.head.versionId).toBe(edited.versionId); });
test("runner failure records a failed run and no proposal", async () => { const { store } = await openStore(); const origin = await seedOrigin(store); const failing = new Map([["failing", { id: "failing", revision: "1", label: "Failing fixture", inference: "none" as const, async propose() { throw new Error("synthetic runner failure"); }, }]]); const workflow = new WorkbenchWorkflow(store, { runners: failing }); const head = await workflow.createWorkingDocument({ originEventId: origin.id, title: "Doc", body: "body\n", requestId: "req-create-0001" }); const selection = await workflow.selectContext({ documentId: head.documentId, eventIds: [origin.id], versionIds: [], requestId: "req-select-001" }); const attempt = workflow.requestProposal({ documentId: head.documentId, selectionId: selection.selectionId, runnerId: "failing", requestId: "req-prop-00001" }); await expect(attempt).rejects.toBeInstanceOf(WorkbenchRunnerError); await expect(attempt).rejects.toThrow("synthetic runner failure"); expect(await eventsOfType(store, PROPOSAL_REQUESTED_EVENT_TYPE)).toHaveLength(1); expect(await eventsOfType(store, PROPOSAL_PROPOSED_EVENT_TYPE)).toHaveLength(0); const runs = await store.listRuns(); expect(runs).toHaveLength(1); expect(runs[0]).toMatchObject({ status: "failed", errorText: "runner-failed: synthetic runner failure" }); await expect(workflow.requestProposal({ documentId: head.documentId, selectionId: selection.selectionId, runnerId: "failing", requestId: "req-prop-00001" })) .rejects.toBeInstanceOf(WorkbenchRunnerError); expect(await store.listRuns()).toHaveLength(1); });});