Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
46 kB · 1096 lines
TypeScript
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097import { createTwoFilesPatch } from "diff";import { canonicalJson, hashJson, parseJsonObject, sha256, type JsonObject } from "../core/json.js";import { stableKey } from "../core/ids.js";import type { CurrentDocument, DocumentVersion, EventCandidate, ThoughtEvent } from "../events/types.js";import { eventIdFor, type JazzThoughtStore } from "../jazz/store.js";import { buildRecentRootObservations } from "../projections/activity.js";import type { AgentRun } from "../store/types.js";import { recordJudgment } from "../training/judgments.js";import { CONTEXT_SELECTED_EVENT_TYPE, DOCUMENT_CREATED_EVENT_TYPE, DOCUMENT_VERSION_EVENT_TYPE, MAX_DOCUMENT_BODY_CHARS, MAX_DOCUMENT_TITLE_CHARS, MAX_EVIDENCE_IDS, MAX_PROPOSED_TEXT_CHARS, PROPOSAL_DECISION_EVENT_TYPE, PROPOSAL_PROPOSED_EVENT_TYPE, PROPOSAL_REQUESTED_EVENT_TYPE, WORKBENCH_CONTEXT_SOURCE, WORKBENCH_DOCUMENT_CONTENT_TYPE, WORKBENCH_DOCUMENT_SOURCE, WORKBENCH_JUDGMENT_CRITERION, WORKBENCH_JUDGMENT_CRITERION_VERSION, WORKBENCH_RUNTIME, WORKBENCH_SCHEMA_VERSION, contextSelectedPayloadSchema, decodeDocumentContent, documentCreatedPayloadSchema, documentVersionPayloadSchema, encodeDocumentContent, plainExcerpt, proposalDecisionPayloadSchema, proposalProposedPayloadSchema, proposalRequestedPayloadSchema, workbenchContextSnapshotSchema, type ContextSelectedPayload, type DocumentVersionPayload, type ProposalProposedPayload, type VersionReason, type WorkbenchContextSnapshot,} from "./contracts.js";import { createDefaultWorkbenchRunnerRegistry, type WorkbenchProposalRequest, type WorkbenchRunnerRegistry,} from "./runners.js";
const JUDGMENT_EVENT_TYPE = "stream.thought.judgment.training-example";const MAX_DIFF_CHARS = 200_000;const MAX_CANDIDATE_EVENTS = 50;
export class StaleBaseError extends Error { constructor(readonly headVersionId: string, message = "The document changed since this base version") { super(message); this.name = "StaleBaseError"; }}
export class WorkbenchNotFoundError extends Error { constructor(message: string) { super(message); this.name = "WorkbenchNotFoundError"; }}
export class WorkbenchConflictError extends Error { constructor(message: string) { super(message); this.name = "WorkbenchConflictError"; }}
export class WorkbenchRunnerError extends Error { constructor(readonly runId: string, message: string) { super(message); this.name = "WorkbenchRunnerError"; }}
export interface WorkbenchWorkflowOptions { actor?: string | undefined; runners?: WorkbenchRunnerRegistry | undefined;}
export interface WorkingDocumentHead { documentId: string; title: string; body: string; versionId: string; sha256: string; sizeBytes: number; updatedAt: string;}
export interface WorkbenchVersionSummary { eventId: string; versionId: string; baseVersionId: string | null; sha256: string; sizeBytes: number; reason: VersionReason; proposalEventId?: string; decisionEventId?: string; createdAt: string;}
export interface WorkbenchSelectionSummary { eventId: string; selectionId: string; snapshotVersionId: string; snapshotSha256: string; baseVersionId: string; createdAt: string; selectedEvents: WorkbenchContextSnapshot["selectedEvents"]; selectedVersions: WorkbenchContextSnapshot["selectedVersions"];}
export interface WorkbenchProposalSummary { eventId: string; requestedEventId: string; selectionId: string; snapshotVersionId: string; runId: string | undefined; runStatus: string | undefined; runnerId: string; runnerRevision: string; inference: "none" | "model" | "unknown"; reason: string; evidenceEventIds: string[]; baseVersionId: string; baseSha256: string; proposedText: string; proposedTextSha256: string; proposedTextChars: number; diff: string; createdAt: string; stale: boolean; headVersionId: string; status: "pending" | "accepted" | "rejected" | "stale"; decision?: { eventId: string; disposition: "accept" | "reject"; submissionId: string; resultVersionId?: string; judgmentEventId?: string; createdAt: string; };}
export interface WorkbenchDocumentDetail { head: WorkingDocumentHead; origin: { eventId: string; type: string; source: string; occurredAt: string; excerpt: string } | undefined; createdEventId: string; versions: WorkbenchVersionSummary[]; selections: WorkbenchSelectionSummary[]; proposals: WorkbenchProposalSummary[]; runners: Array<{ id: string; revision: string; label: string; inference: "none" | "model" }>;}
export interface WorkbenchDocumentListItem { documentId: string; title: string; originEventId: string; headVersionId: string; sha256: string; updatedAt: string; versionCount: number; proposalCount: number; pendingProposalCount: number;}
export interface WorkbenchCandidates { documentId: string; events: Array<{ eventId: string; type: string; source: string; occurredAt: string; title: string; excerpt: string; origin: boolean }>; versions: Array<{ versionId: string; sha256: string; reason: VersionReason; createdAt: string; head: boolean }>; bounded: boolean;}
export function workingDocumentId(originEventId: string, requestId: string): string { return stableKey("working-document", originEventId, requestId);}
export function workingDocumentVersionId(documentId: string, contentSha256: string): string { return stableKey("version", WORKBENCH_DOCUMENT_SOURCE, documentId, contentSha256);}
export function workbenchSelectionId(documentId: string, requestId: string): string { return stableKey("workbench-selection", documentId, requestId);}
export function workbenchRunId(documentId: string, requestId: string): string { return stableKey("workbench-run", documentId, requestId);}
function headRowId(documentId: string): string { return stableKey("document", WORKBENCH_DOCUMENT_SOURCE, documentId);}
function createdEventKey(documentId: string): string { return stableKey("workbench-document-created", documentId);}
function versionEventKey(documentId: string, reason: VersionReason, requestId: string): string { return stableKey("workbench-document-version", documentId, reason, requestId);}
function selectedEventKey(selectionId: string): string { return stableKey("workbench-context-selected", selectionId);}
function requestedEventKey(documentId: string, requestId: string): string { return stableKey("workbench-proposal-requested", documentId, requestId);}
function proposedEventKey(documentId: string, requestId: string): string { return stableKey("workbench-proposal-proposed", documentId, requestId);}
function decisionEventKey(proposalEventId: string): string { return stableKey("workbench-proposal-decision", proposalEventId);}
const documentQueues = new Map<string, Promise<void>>();
async function withDocumentLock<T>(documentId: string, operation: () => Promise<T>): Promise<T> { const previous = documentQueues.get(documentId) ?? Promise.resolve(); let release!: () => void; const current = new Promise<void>((resolve) => { release = resolve; }); const queued = previous.then(() => current); documentQueues.set(documentId, queued); await previous; try { return await operation(); } finally { release(); if (documentQueues.get(documentId) === queued) documentQueues.delete(documentId); }}
/** * Trusted mutation path for working documents. Every method appends * append-only events through the existing store, keeps the `documents` row as * the current-head projection, and never touches a file on disk. */export class WorkbenchWorkflow { private readonly actor: string; private readonly runners: WorkbenchRunnerRegistry;
constructor(private readonly store: JazzThoughtStore, options: WorkbenchWorkflowOptions = {}) { this.actor = options.actor ?? "operator:local"; this.runners = options.runners ?? createDefaultWorkbenchRunnerRegistry(); }
listRunners(): WorkbenchDocumentDetail["runners"] { return [...this.runners.values()].map((runner) => ({ id: runner.id, revision: runner.revision, label: runner.label, inference: runner.inference, })); }
async createWorkingDocument(input: { originEventId: string; title: string; body: string; requestId: string; }): Promise<WorkingDocumentHead> { const origin = await this.store.getEvent(input.originEventId); if (!origin) throw new WorkbenchNotFoundError("Origin event was not found"); const title = validateTitle(input.title); const body = validateBody(input.body); const documentId = workingDocumentId(origin.id, input.requestId); return withDocumentLock(documentId, async () => { const content = encodeDocumentContent(title, body); const contentSha256 = sha256(content); const versionId = workingDocumentVersionId(documentId, contentSha256); const existing = await this.getHeadRow(documentId); if (existing) { if (existing.sha256 !== contentSha256) { throw new WorkbenchConflictError("Request id was reused with different document content"); } return this.headFromRow(existing); } const createdAt = new Date().toISOString(); await this.appendVersionRow(documentId, versionId, content, contentSha256, createdAt); const created = await this.store.appendEvent(this.candidate({ type: DOCUMENT_CREATED_EVENT_TYPE, idempotencyKey: createdEventKey(documentId), rootEventId: origin.rootEventId, parentEventId: origin.id, correlationId: documentId, occurredAt: createdAt, payload: documentCreatedPayloadSchema.parse({ documentId, title, originEventId: origin.id, versionId, sha256: contentSha256, requestId: input.requestId, }), })); await this.store.appendEvent(this.candidate({ type: DOCUMENT_VERSION_EVENT_TYPE, idempotencyKey: versionEventKey(documentId, "created", input.requestId), rootEventId: origin.rootEventId, parentEventId: created.event.id, correlationId: documentId, occurredAt: createdAt, payload: versionPayload({ documentId, versionId, baseVersionId: null, sha256: contentSha256, sizeBytes: Buffer.byteLength(content), reason: "created", requestId: input.requestId, }), })); const head = await this.upsertHead(documentId, versionId, contentSha256, Buffer.byteLength(content), createdAt); return this.headFromRow(head); }); }
async saveOperatorEdit(input: { documentId: string; baseVersionId: string; title: string; body: string; requestId: string; }): Promise<WorkingDocumentHead> { const title = validateTitle(input.title); const body = validateBody(input.body); return withDocumentLock(input.documentId, async () => { const lineage = await this.requireLineage(input.documentId); const replayed = await this.store.getEvent(eventIdFor( WORKBENCH_DOCUMENT_SOURCE, versionEventKey(input.documentId, "operator-edit", input.requestId), )); if (replayed) { // A retry after an uncertain outcome: settle the head if the append // landed but the projection did not, then return the current head. const payload = documentVersionPayloadSchema.parse(replayed.payload); return this.headFromRow(await this.settleHead(input.documentId, payload)); } const head = await this.requireHeadRow(input.documentId); if (input.baseVersionId !== head.versionId) throw new StaleBaseError(head.versionId); const content = encodeDocumentContent(title, body); const contentSha256 = sha256(content); const versionId = workingDocumentVersionId(input.documentId, contentSha256); if (versionId === head.versionId) return this.headFromRow(head); const createdAt = new Date().toISOString(); await this.appendVersionRow(input.documentId, versionId, content, contentSha256, createdAt); await this.store.appendEvent(this.candidate({ type: DOCUMENT_VERSION_EVENT_TYPE, idempotencyKey: versionEventKey(input.documentId, "operator-edit", input.requestId), rootEventId: lineage.created.rootEventId, parentEventId: lineage.created.id, correlationId: input.documentId, occurredAt: createdAt, payload: versionPayload({ documentId: input.documentId, versionId, baseVersionId: head.versionId, sha256: contentSha256, sizeBytes: Buffer.byteLength(content), reason: "operator-edit", requestId: input.requestId, }), })); const updated = await this.upsertHead(input.documentId, versionId, contentSha256, Buffer.byteLength(content), createdAt); return this.headFromRow(updated); }); }
async selectContext(input: { documentId: string; eventIds: string[]; versionIds: string[]; requestId: string; }): Promise<WorkbenchSelectionSummary> { const lineage = await this.requireLineage(input.documentId); const head = await this.requireHeadRow(input.documentId); const selectionId = workbenchSelectionId(input.documentId, input.requestId); const existing = await this.store.getEvent(eventIdFor(WORKBENCH_DOCUMENT_SOURCE, selectedEventKey(selectionId))); if (existing) return this.selectionSummary(existing); const eventIds = [...new Set(input.eventIds)]; const versionIds = [...new Set(input.versionIds)]; const events = await this.store.getEvents(eventIds); const eventsById = new Map(events.map((event) => [event.id, event])); for (const id of eventIds) if (!eventsById.has(id)) throw new WorkbenchNotFoundError(`Selected event was not found: ${id}`); const versions: DocumentVersion[] = []; for (const id of versionIds) { const version = await this.store.getDocumentVersion(id); if (!version) throw new WorkbenchNotFoundError(`Selected version was not found: ${id}`); versions.push(version); } const snapshot: WorkbenchContextSnapshot = { documentId: input.documentId, baseVersionId: head.versionId, selectedEvents: eventIds.map((id) => { const event = eventsById.get(id)!; return { eventId: event.id, type: event.type, source: event.source, occurredAt: event.occurredAt, payloadHash: event.payloadHash, excerpt: plainExcerpt(event.payload), }; }), selectedVersions: versions.map((version) => ({ documentId: version.documentId, versionId: version.id, sha256: version.sha256, path: version.path, })), }; const content = canonicalJson(workbenchContextSnapshotSchema.parse(snapshot) as unknown as JsonObject); // Deterministic over the exact base version and the sorted selected ids; // the id list is hashed so the identity stays bounded for any selection size. const sortedIds = [...eventIds, ...versionIds].sort(); const snapshotVersionId = stableKey( "workbench-context-snapshot", input.documentId, sha256(canonicalJson({ baseVersionId: head.versionId, ids: sortedIds })), ); const createdAt = new Date().toISOString(); await this.store.appendDocumentVersion({ id: snapshotVersionId, source: WORKBENCH_CONTEXT_SOURCE, documentId: snapshotVersionId, path: `workbench-context/${encodeURIComponent(input.documentId)}/${encodeURIComponent(selectionId)}.json`, contentType: "application/json", sha256: sha256(content), content, sizeBytes: Buffer.byteLength(content), mtimeMs: Date.parse(createdAt), createdAt, }); const stored = await this.store.getDocumentVersion(snapshotVersionId); if (!stored) throw new Error("Context snapshot insertion left no durable evidence"); if (stored.sha256 !== sha256(content)) throw new WorkbenchConflictError("Context snapshot identity collides with different content"); const selected = await this.store.appendEvent(this.candidate({ type: CONTEXT_SELECTED_EVENT_TYPE, idempotencyKey: selectedEventKey(selectionId), rootEventId: lineage.created.rootEventId, parentEventId: lineage.created.id, correlationId: input.documentId, occurredAt: createdAt, payload: contextSelectedPayloadSchema.parse({ documentId: input.documentId, selectionId, snapshotVersionId, selectedEventIds: eventIds, selectedVersionIds: versionIds, requestId: input.requestId, }) as unknown as JsonObject, })); return this.selectionSummary(selected.event); }
async requestProposal(input: { documentId: string; selectionId: string; runnerId: string; requestId: string; }): Promise<{ proposalEventId: string; requestedEventId: string; runId: string }> { const runner = this.runners.get(input.runnerId); if (!runner) throw new WorkbenchNotFoundError(`Unknown proposal runner: ${input.runnerId}`); return withDocumentLock(input.documentId, async () => { const lineage = await this.requireLineage(input.documentId); const runId = workbenchRunId(input.documentId, input.requestId); const requestedEventId = eventIdFor(WORKBENCH_DOCUMENT_SOURCE, requestedEventKey(input.documentId, input.requestId)); const proposalEventId = eventIdFor(WORKBENCH_DOCUMENT_SOURCE, proposedEventKey(input.documentId, input.requestId)); const replayed = await this.store.getEvent(proposalEventId); if (replayed) return { proposalEventId, requestedEventId, runId }; const priorRun = await this.store.getRun(runId); if (priorRun?.status === "failed") { throw new WorkbenchRunnerError(runId, priorRun.errorText || "The proposal runner failed"); } const selected = await this.store.getEvent(eventIdFor(WORKBENCH_DOCUMENT_SOURCE, selectedEventKey(input.selectionId))); if (!selected || selected.type !== CONTEXT_SELECTED_EVENT_TYPE) throw new WorkbenchNotFoundError("Context selection was not found"); const selection = contextSelectedPayloadSchema.parse(selected.payload); if (selection.documentId !== input.documentId) throw new WorkbenchNotFoundError("Context selection belongs to another document"); const snapshotVersion = await this.store.getDocumentVersion(selection.snapshotVersionId); if (!snapshotVersion) throw new WorkbenchNotFoundError("Context snapshot was not found"); const snapshot = workbenchContextSnapshotSchema.parse(parseJsonObject(snapshotVersion.content)) as WorkbenchContextSnapshot; const head = await this.requireHeadRow(input.documentId); const headVersion = await this.requireVersion(head.versionId); const decoded = decodeDocumentContent(headVersion.content); const createdAt = new Date().toISOString(); const requested = await this.store.appendEvent(this.candidate({ type: PROPOSAL_REQUESTED_EVENT_TYPE, idempotencyKey: requestedEventKey(input.documentId, input.requestId), rootEventId: lineage.created.rootEventId, parentEventId: selected.id, correlationId: input.documentId, occurredAt: createdAt, payload: proposalRequestedPayloadSchema.parse({ documentId: input.documentId, baseVersionId: head.versionId, baseSha256: head.sha256, selectionId: input.selectionId, snapshotVersionId: selection.snapshotVersionId, runnerId: runner.id, requestId: input.requestId, }), })); const request: WorkbenchProposalRequest = { documentId: input.documentId, baseVersionId: head.versionId, baseSha256: head.sha256, title: decoded.title, baseText: decoded.body, snapshot, maxProposedChars: MAX_PROPOSED_TEXT_CHARS, }; const run: AgentRun = { id: runId, executionKey: runId, triggerEventId: requested.event.id, agentId: runner.id, agentVersion: 1, status: "running", inputEventIds: [requested.event.id, ...selection.selectedEventIds], outputEventIds: [], attempt: 1, provider: "fixture", model: "deterministic-workbench@1", privacy: "sensitive", promptHash: hashJson(request as unknown as JsonObject), contextManifest: { contextStrategy: "workbench-selection", inference: runner.inference, documentId: input.documentId, baseVersionId: head.versionId, baseSha256: head.sha256, selectionId: input.selectionId, runnerId: runner.id, runnerRevision: runner.revision, tools: [], proposals: [], externalActions: false, privacy: "sensitive", contextSnapshot: { id: selection.snapshotVersionId, storage: "jazz-document-version", sha256: snapshotVersion.sha256, }, }, createdAt, startedAt: createdAt, updatedAt: createdAt, }; await this.store.upsertRun(run); let payload: ProposalProposedPayload; try { const result = await runner.propose(request); const admitted = new Set(selection.selectedEventIds); const evidenceEventIds = [...new Set(result.evidenceEventIds)].slice(0, MAX_EVIDENCE_IDS); for (const id of evidenceEventIds) { if (!admitted.has(id)) throw new Error("Runner cited an event outside the selected context"); } payload = proposalProposedPayloadSchema.parse({ proposalState: "runner-proposed", proposer: { runnerId: runner.id, runnerRevision: runner.revision, requestId: input.requestId, runId, contextSnapshotId: selection.snapshotVersionId, }, target: { documentId: input.documentId, baseVersionId: head.versionId, baseSha256: head.sha256 }, operation: "replace-document", proposedText: result.proposedText, proposedTextChars: result.proposedText.length, proposedTextSha256: sha256(result.proposedText), reason: result.reason, evidenceEventIds, publicationEligible: false, }); } catch (error) { const completedAt = new Date().toISOString(); await this.store.upsertRun({ ...run, status: "failed", errorText: classifyRunnerFailure(error), completedAt, updatedAt: completedAt, }); throw new WorkbenchRunnerError(runId, classifyRunnerFailure(error)); } const proposed = await this.store.appendEvent(this.candidate({ type: PROPOSAL_PROPOSED_EVENT_TYPE, idempotencyKey: proposedEventKey(input.documentId, input.requestId), rootEventId: lineage.created.rootEventId, parentEventId: requested.event.id, correlationId: input.documentId, occurredAt: createdAt, payload: payload as unknown as JsonObject, })); const completedAt = new Date().toISOString(); await this.store.upsertRun({ ...run, status: "completed", outputEventIds: [proposed.event.id], result: { summary: payload.reason, proposedTextSha256: payload.proposedTextSha256, proposedTextChars: payload.proposedTextChars, evidenceCount: payload.evidenceEventIds.length, }, completedAt, updatedAt: completedAt, }); return { proposalEventId: proposed.event.id, requestedEventId: requested.event.id, runId }; }); }
async decideProposal(input: { proposalEventId: string; disposition: "accept" | "reject"; submissionId: string; }): Promise<{ decisionEventId: string; disposition: "accept" | "reject"; resultVersionId?: string; judgmentEventId?: string; replayed: boolean }> { const proposal = await this.store.getEvent(input.proposalEventId); if (!proposal || proposal.type !== PROPOSAL_PROPOSED_EVENT_TYPE) throw new WorkbenchNotFoundError("Proposal was not found"); const proposalPayload = proposalProposedPayloadSchema.parse(proposal.payload); const documentId = proposalPayload.target.documentId; return withDocumentLock(documentId, async () => { const base = await this.requireVersion(proposalPayload.target.baseVersionId); if (base.sha256 !== proposalPayload.target.baseSha256) throw new WorkbenchConflictError("Proposal base version evidence is inconsistent"); const baseTitle = decodeDocumentContent(base.content).title; const resultContent = encodeDocumentContent(baseTitle, proposalPayload.proposedText); const resultSha256 = sha256(resultContent); const resultVersionId = workingDocumentVersionId(documentId, resultSha256); const decisionEventId = eventIdFor(WORKBENCH_DOCUMENT_SOURCE, decisionEventKey(proposal.id)); const payload = proposalDecisionPayloadSchema.parse({ proposalEventId: proposal.id, disposition: input.disposition, submissionId: input.submissionId, authority: "human", ...(input.disposition === "accept" ? { resultVersionId } : {}), }); const decisions = (await this.store.listChildEvents(proposal.id)).filter((event) => event.type === PROPOSAL_DECISION_EVENT_TYPE); const sameSubmission = decisions.find((decision) => decision.payload.submissionId === input.submissionId); if (sameSubmission) { if (canonicalJson(sameSubmission.payload) !== canonicalJson(payload as unknown as JsonObject)) { throw new WorkbenchConflictError("Proposal decision submission id conflicts with existing content"); } // Replay: finish any settlement the earlier attempt did not reach; apply nothing twice. if (payload.disposition === "accept") { const versionEvent = await this.store.getEvent(eventIdFor( WORKBENCH_DOCUMENT_SOURCE, versionEventKey(documentId, "proposal-accepted", input.submissionId), )); if (versionEvent) await this.settleHead(documentId, documentVersionPayloadSchema.parse(versionEvent.payload)); } const judgment = await this.recordDecisionJudgment(proposalPayload, sameSubmission, payload.disposition); return { decisionEventId: sameSubmission.id, disposition: payload.disposition, ...(payload.resultVersionId ? { resultVersionId: payload.resultVersionId } : {}), ...(judgment ? { judgmentEventId: judgment.id } : {}), replayed: true, }; } if (decisions.length > 0) throw new WorkbenchConflictError("Proposal already has a human decision"); const head = await this.requireHeadRow(documentId); const occurredAt = new Date().toISOString(); if (payload.disposition === "accept") { if (head.versionId !== proposalPayload.target.baseVersionId) throw new StaleBaseError(head.versionId); // Write order for crash safety: version row, version event, decision event, head upsert, judgment. await this.appendVersionRow(documentId, resultVersionId, resultContent, resultSha256, occurredAt); await this.store.appendEvent(this.candidate({ type: DOCUMENT_VERSION_EVENT_TYPE, idempotencyKey: versionEventKey(documentId, "proposal-accepted", input.submissionId), rootEventId: proposal.rootEventId, parentEventId: proposal.id, correlationId: documentId, occurredAt, payload: versionPayload({ documentId, versionId: resultVersionId, baseVersionId: proposalPayload.target.baseVersionId, sha256: resultSha256, sizeBytes: Buffer.byteLength(resultContent), reason: "proposal-accepted", proposalEventId: proposal.id, decisionEventId, requestId: input.submissionId, }), })); } const decision = await this.store.appendEvent(this.candidate({ type: PROPOSAL_DECISION_EVENT_TYPE, idempotencyKey: decisionEventKey(proposal.id), rootEventId: proposal.rootEventId, parentEventId: proposal.id, correlationId: documentId, occurredAt, payload: payload as unknown as JsonObject, })); if (decision.event.id !== decisionEventId) throw new Error("Decision event identity diverged from its deterministic id"); if (payload.disposition === "accept") { await this.upsertHead(documentId, resultVersionId, resultSha256, Buffer.byteLength(resultContent), occurredAt); } const judgment = await this.recordDecisionJudgment(proposalPayload, decision.event, payload.disposition); return { decisionEventId: decision.event.id, disposition: payload.disposition, ...(payload.resultVersionId ? { resultVersionId: payload.resultVersionId } : {}), ...(judgment ? { judgmentEventId: judgment.id } : {}), replayed: false, }; }); }
async listDocuments(): Promise<WorkbenchDocumentListItem[]> { const heads = (await this.store.listCurrentDocuments(WORKBENCH_DOCUMENT_SOURCE)).filter((row) => !row.deleted); if (heads.length === 0) return []; const events = await this.workbenchEvents(); const items: WorkbenchDocumentListItem[] = []; for (const row of heads) { const documentId = row.documentId; const created = events.find((event) => event.type === DOCUMENT_CREATED_EVENT_TYPE && event.payload.documentId === documentId); if (!created) continue; const version = await this.store.getDocumentVersion(row.versionId); const title = version ? decodeDocumentContent(version.content).title : String(created.payload.title); const proposals = events.filter((event) => event.type === PROPOSAL_PROPOSED_EVENT_TYPE && proposalDocumentId(event) === documentId); const decided = new Set(events .filter((event) => event.type === PROPOSAL_DECISION_EVENT_TYPE) .map((event) => String(event.payload.proposalEventId))); items.push({ documentId, title, originEventId: String(created.payload.originEventId), headVersionId: row.versionId, sha256: row.sha256, updatedAt: row.updatedAt, versionCount: events.filter((event) => event.type === DOCUMENT_VERSION_EVENT_TYPE && event.payload.documentId === documentId).length, proposalCount: proposals.length, pendingProposalCount: proposals.filter((event) => !decided.has(event.id)).length, }); } return items.sort((left, right) => right.updatedAt.localeCompare(left.updatedAt) || left.documentId.localeCompare(right.documentId)); }
async getDocument(documentId: string): Promise<WorkbenchDocumentDetail | undefined> { const row = await this.getHeadRow(documentId); if (!row || row.deleted) return undefined; const events = (await this.workbenchEvents()).filter((event) => eventDocumentId(event) === documentId); const created = events.find((event) => event.type === DOCUMENT_CREATED_EVENT_TYPE); if (!created) return undefined; const head = this.headFromRow(row); const headVersion = await this.store.getDocumentVersion(row.versionId); if (headVersion) { const decoded = decodeDocumentContent(headVersion.content); head.title = decoded.title; head.body = decoded.body; } const origin = await this.store.getEvent(String(created.payload.originEventId)); const versions = events .filter((event) => event.type === DOCUMENT_VERSION_EVENT_TYPE) .map((event) => { const payload = documentVersionPayloadSchema.parse(event.payload); return { eventId: event.id, versionId: payload.versionId, baseVersionId: payload.baseVersionId, sha256: payload.sha256, sizeBytes: payload.sizeBytes, reason: payload.reason, ...(payload.proposalEventId ? { proposalEventId: payload.proposalEventId } : {}), ...(payload.decisionEventId ? { decisionEventId: payload.decisionEventId } : {}), createdAt: event.observedAt, } satisfies WorkbenchVersionSummary; }); const selections: WorkbenchSelectionSummary[] = []; for (const event of events.filter((candidate) => candidate.type === CONTEXT_SELECTED_EVENT_TYPE)) { selections.push(await this.selectionSummary(event)); } const decisions = events.filter((event) => event.type === PROPOSAL_DECISION_EVENT_TYPE); const proposals: WorkbenchProposalSummary[] = []; for (const event of events.filter((candidate) => candidate.type === PROPOSAL_PROPOSED_EVENT_TYPE)) { const payload = proposalProposedPayloadSchema.parse(event.payload); const requested = events.find((candidate) => candidate.id === event.parentEventId && candidate.type === PROPOSAL_REQUESTED_EVENT_TYPE); const requestedPayload = requested ? proposalRequestedPayloadSchema.parse(requested.payload) : undefined; const decision = decisions.find((candidate) => candidate.payload.proposalEventId === event.id); const decisionPayload = decision ? proposalDecisionPayloadSchema.parse(decision.payload) : undefined; const judgment = decision ? (await this.store.listChildEvents(decision.id)).find((candidate) => candidate.type === JUDGMENT_EVENT_TYPE) : undefined; const run = payload.proposer.runId ? await this.store.getRun(payload.proposer.runId) : undefined; const base = await this.store.getDocumentVersion(payload.target.baseVersionId); const baseBody = base ? decodeDocumentContent(base.content).body : ""; const stale = row.versionId !== payload.target.baseVersionId; const status: WorkbenchProposalSummary["status"] = decisionPayload ? decisionPayload.disposition === "accept" ? "accepted" : "rejected" : stale ? "stale" : "pending"; const inference = run?.contextManifest.inference === "none" || run?.contextManifest.inference === "model" ? run.contextManifest.inference : "unknown"; proposals.push({ eventId: event.id, requestedEventId: requested?.id ?? "", selectionId: requestedPayload?.selectionId ?? "", snapshotVersionId: payload.proposer.contextSnapshotId, runId: payload.proposer.runId, runStatus: run?.status, runnerId: payload.proposer.runnerId, runnerRevision: payload.proposer.runnerRevision, inference, reason: payload.reason, evidenceEventIds: payload.evidenceEventIds, baseVersionId: payload.target.baseVersionId, baseSha256: payload.target.baseSha256, proposedText: payload.proposedText, proposedTextSha256: payload.proposedTextSha256, proposedTextChars: payload.proposedTextChars, diff: boundedUnifiedDiff(documentId, payload.target.baseVersionId, event.id, baseBody, payload.proposedText), createdAt: event.observedAt, stale, headVersionId: row.versionId, status, ...(decision && decisionPayload ? { decision: { eventId: decision.id, disposition: decisionPayload.disposition, submissionId: decisionPayload.submissionId, ...(decisionPayload.resultVersionId ? { resultVersionId: decisionPayload.resultVersionId } : {}), ...(judgment ? { judgmentEventId: judgment.id } : {}), createdAt: decision.observedAt, }, } : {}), }); } return { head, origin: origin ? { eventId: origin.id, type: origin.type, source: origin.source, occurredAt: origin.occurredAt, excerpt: plainExcerpt(origin.payload, 240), } : undefined, createdEventId: created.id, versions, selections, proposals, runners: this.listRunners(), }; }
/** Bounded selectable context: recent root observations plus this document's own versions. */ async listCandidates(documentId: string): Promise<WorkbenchCandidates | undefined> { const row = await this.getHeadRow(documentId); if (!row || row.deleted) return undefined; const events = (await this.workbenchEvents()).filter((event) => eventDocumentId(event) === documentId); const created = events.find((event) => event.type === DOCUMENT_CREATED_EVENT_TYPE); if (!created) return undefined; const originEventId = String(created.payload.originEventId); const observations = await buildRecentRootObservations(this.store, MAX_CANDIDATE_EVENTS); const candidates: WorkbenchCandidates["events"] = observations.items.map((item) => ({ eventId: item.id, type: item.type, source: item.source, occurredAt: item.occurredAt, title: item.presentation?.title ?? item.summary, excerpt: (item.presentation?.body || item.presentation?.objectLabel || item.summary).slice(0, 200), origin: item.id === originEventId, })); if (!candidates.some((candidate) => candidate.origin)) { const origin = await this.store.getEvent(originEventId); if (origin) { candidates.unshift({ eventId: origin.id, type: origin.type, source: origin.source, occurredAt: origin.occurredAt, title: origin.type.split(".").at(-1) ?? origin.type, excerpt: plainExcerpt(origin.payload, 200), origin: true, }); } } const versions = events .filter((event) => event.type === DOCUMENT_VERSION_EVENT_TYPE) .map((event) => { const payload = documentVersionPayloadSchema.parse(event.payload); return { versionId: payload.versionId, sha256: payload.sha256, reason: payload.reason, createdAt: event.observedAt, head: payload.versionId === row.versionId, }; }); return { documentId, events: candidates, versions, bounded: observations.rootWindowComplete === false || observations.items.length >= MAX_CANDIDATE_EVENTS }; }
private async recordDecisionJudgment( proposal: ProposalProposedPayload, decision: ThoughtEvent, disposition: "accept" | "reject", ): Promise<ThoughtEvent | undefined> { if (!proposal.proposer.runId) return undefined; return recordJudgment(this.store, { runId: proposal.proposer.runId, kind: disposition, criterion: WORKBENCH_JUDGMENT_CRITERION, criterionVersion: WORKBENCH_JUDGMENT_CRITERION_VERSION, qualityEligible: false, externalExportEligible: false, notes: disposition === "accept" ? "Operator accepted a workbench document proposal" : "Operator rejected a workbench document proposal", actor: this.actor, source: "judgment:workbench-documents", feedbackSourceEventId: decision.id, }); }
private async selectionSummary(event: ThoughtEvent): Promise<WorkbenchSelectionSummary> { const payload = contextSelectedPayloadSchema.parse(event.payload) as ContextSelectedPayload; const snapshotVersion = await this.store.getDocumentVersion(payload.snapshotVersionId); const snapshot = snapshotVersion ? workbenchContextSnapshotSchema.parse(parseJsonObject(snapshotVersion.content)) as WorkbenchContextSnapshot : undefined; return { eventId: event.id, selectionId: payload.selectionId, snapshotVersionId: payload.snapshotVersionId, snapshotSha256: snapshotVersion?.sha256 ?? "", baseVersionId: snapshot?.baseVersionId ?? "", createdAt: event.observedAt, selectedEvents: snapshot?.selectedEvents ?? [], selectedVersions: snapshot?.selectedVersions ?? [], }; }
private async workbenchEvents(): Promise<ThoughtEvent[]> { return this.store.listEvents({ source: WORKBENCH_DOCUMENT_SOURCE }); }
private async requireLineage(documentId: string): Promise<{ created: ThoughtEvent }> { const created = await this.store.getEvent(eventIdFor(WORKBENCH_DOCUMENT_SOURCE, createdEventKey(documentId))); if (!created || created.type !== DOCUMENT_CREATED_EVENT_TYPE) throw new WorkbenchNotFoundError("Working document was not found"); return { created }; }
private async getHeadRow(documentId: string): Promise<CurrentDocument | undefined> { const id = headRowId(documentId); return (await this.store.listCurrentDocuments(WORKBENCH_DOCUMENT_SOURCE)).find((row) => row.id === id); }
private async requireHeadRow(documentId: string): Promise<CurrentDocument> { const row = await this.getHeadRow(documentId); if (!row || row.deleted) throw new WorkbenchNotFoundError("Working document was not found"); return row; }
private async requireVersion(versionId: string): Promise<DocumentVersion> { const version = await this.store.getDocumentVersion(versionId); if (!version) throw new WorkbenchNotFoundError(`Document version was not found: ${versionId}`); return version; }
private headFromRow(row: CurrentDocument): WorkingDocumentHead { return { documentId: row.documentId, title: "", body: "", versionId: row.versionId, sha256: row.sha256, sizeBytes: row.sizeBytes, updatedAt: row.updatedAt, }; }
private async appendVersionRow(documentId: string, versionId: string, content: string, contentSha256: string, createdAt: string): Promise<void> { await this.store.appendDocumentVersion({ id: versionId, source: WORKBENCH_DOCUMENT_SOURCE, documentId, path: `documents/${encodeURIComponent(documentId)}.md`, contentType: WORKBENCH_DOCUMENT_CONTENT_TYPE, sha256: contentSha256, content, sizeBytes: Buffer.byteLength(content), mtimeMs: Date.parse(createdAt), createdAt, }); const stored = await this.store.getDocumentVersion(versionId); if (!stored || stored.sha256 !== contentSha256) throw new WorkbenchConflictError("Version identity collides with different content"); }
private async upsertHead(documentId: string, versionId: string, contentSha256: string, sizeBytes: number, updatedAt: string): Promise<CurrentDocument> { const row: CurrentDocument = { id: headRowId(documentId), source: WORKBENCH_DOCUMENT_SOURCE, documentId, path: `documents/${encodeURIComponent(documentId)}.md`, versionId, sha256: contentSha256, contentType: WORKBENCH_DOCUMENT_CONTENT_TYPE, sizeBytes, mtimeMs: Date.parse(updatedAt), deleted: false, updatedAt, }; await this.store.upsertCurrentDocument(row); return row; }
/** Advance the head only when it still points at the version the event was based on. */ private async settleHead(documentId: string, payload: DocumentVersionPayload): Promise<CurrentDocument> { const head = await this.requireHeadRow(documentId); if (head.versionId === payload.baseVersionId && head.versionId !== payload.versionId) { return this.upsertHead(documentId, payload.versionId, payload.sha256, payload.sizeBytes, new Date().toISOString()); } return head; }
private candidate(input: { type: string; idempotencyKey: string; rootEventId: string; parentEventId: string; correlationId: string; occurredAt: string; payload: JsonObject; }): EventCandidate { return { type: input.type, schemaVersion: WORKBENCH_SCHEMA_VERSION, source: WORKBENCH_DOCUMENT_SOURCE, sourceKind: "system", externalId: input.idempotencyKey, idempotencyKey: input.idempotencyKey, occurredAt: input.occurredAt, actor: this.actor, rootEventId: input.rootEventId, parentEventId: input.parentEventId, correlationId: input.correlationId, privacy: "sensitive", payload: input.payload, createdByRuntime: WORKBENCH_RUNTIME, }; }}
function versionPayload(input: DocumentVersionPayload): JsonObject { return documentVersionPayloadSchema.parse(input) as unknown as JsonObject;}
function validateTitle(title: string): string { const trimmed = title.replace(/[\r\n]+/g, " ").trim(); if (!trimmed) throw new WorkbenchConflictError("Title is required"); if (trimmed.length > MAX_DOCUMENT_TITLE_CHARS) throw new WorkbenchConflictError("Title is too long"); return trimmed;}
function validateBody(body: string): string { const normalized = body.replace(/\r\n?/g, "\n"); if (normalized.length > MAX_DOCUMENT_BODY_CHARS) throw new WorkbenchConflictError("Document body is too long"); return normalized;}
function classifyRunnerFailure(error: unknown): string { const message = error instanceof Error ? error.message : String(error); return `runner-failed: ${message.slice(0, 500)}`;}
function proposalDocumentId(event: ThoughtEvent): string | undefined { const target = event.payload.target; return target && typeof target === "object" && !Array.isArray(target) && typeof target.documentId === "string" ? target.documentId : undefined;}
function eventDocumentId(event: ThoughtEvent): string | undefined { if (event.type === PROPOSAL_PROPOSED_EVENT_TYPE) return proposalDocumentId(event); if (event.type === PROPOSAL_DECISION_EVENT_TYPE) return event.correlationId; return typeof event.payload.documentId === "string" ? event.payload.documentId : undefined;}
export function boundedUnifiedDiff(documentId: string, baseVersionId: string, proposalEventId: string, base: string, proposed: string): string { const name = `documents/${encodeURIComponent(documentId)}.md`; const patch = createTwoFilesPatch(name, name, base, proposed, `base ${baseVersionId}`, `proposal ${proposalEventId}`, { context: 3 }); if (patch.length <= MAX_DIFF_CHARS) return patch; return `${patch.slice(0, MAX_DIFF_CHARS)}\n... diff truncated at ${MAX_DIFF_CHARS} characters ...\n`;}