Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
12 kB · 292 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293import { canonicalStructuredOutput, createOutputContractRegistry, outputContractIdentityJson, parseOutputContractIdentity, type OutputContractIdentity,} from "../agents/output-contracts.js";import { CORRECTION_PROPOSAL_EVENT_TYPE, REPAIR_REQUEST_EVENT_TYPE } from "../agents/repairs.js";import { canonicalJson, type JsonObject } from "../core/json.js";import { stableKey } from "../core/ids.js";import type { ThoughtEvent } from "../events/types.js";import type { JazzThoughtStore } from "../jazz/store.js";import type { AgentRun, Projection } from "../store/types.js";import { joinPrivacy, runPrivacy } from "../security/privacy.js";
export const EFFECTIVE_OUTPUT_PROJECTION_VERSION = 2;
export interface ActiveJudgmentSet { active: ThoughtEvent[]; inactiveIds: Set<string>;}
export async function activeJudgments(store: JazzThoughtStore): Promise<ActiveJudgmentSet> { const judgments = await store.listEvents({ types: ["stream.thought.judgment.training-example"] }); const retractions = await store.listEvents({ types: ["stream.thought.judgment.training-example.retracted"] }); const inactiveIds = new Set<string>(); for (const judgment of judgments) { const superseded = optionalString(judgment.payload.supersedesJudgmentEventId); if (superseded) inactiveIds.add(superseded); } for (const retraction of retractions) { const retracted = optionalString(retraction.payload.retractedJudgmentEventId); if (retracted) inactiveIds.add(retracted); } return { active: judgments.filter((judgment) => !inactiveIds.has(judgment.id)), inactiveIds, };}
export async function originalRunIdForRun(store: JazzThoughtStore, run: AgentRun): Promise<string> { if (run.contextManifest.agentRole !== "repair") return run.id; const proposalId = run.outputEventIds[0]; if (proposalId) { const proposal = await store.getEvent(proposalId); const originalRunId = proposal && optionalString(proposal.payload.originalRunId); if (originalRunId) return originalRunId; } const request = await store.getEvent(run.triggerEventId); if (request?.type !== REPAIR_REQUEST_EVENT_TYPE) throw new Error(`Repair run ${run.id} has no valid repair request`); const originalRunId = optionalString(request.payload.originalRunId); if (!originalRunId) throw new Error(`Repair request ${request.id} has no original run id`); return originalRunId;}
export async function rebuildEffectiveOutputForRun( store: JazzThoughtStore, runId: string,): Promise<Projection> { const run = await store.getRun(runId); if (!run) throw new Error(`Run not found: ${runId}`); return rebuildEffectiveOutput(store, await originalRunIdForRun(store, run));}
export async function rebuildEffectiveOutput( store: JazzThoughtStore, originalRunId: string,): Promise<Projection> { const originalRun = await store.getRun(originalRunId); if (!originalRun) throw new Error(`Original run not found: ${originalRunId}`); if (originalRun.contextManifest.agentRole === "repair") { throw new Error(`Effective-output projection requires a non-repair original run: ${originalRunId}`); } const originalSource = await requireEvent(store, originalRun.triggerEventId); const contract = parseOutputContractIdentity(originalRun.contextManifest.outputContract); const registry = createOutputContractRegistry(); registry.resolve(contract); const judgmentSet = await activeJudgments(store); const originalOutput = await validatedOriginalOutput(store, originalRun, contract); const directCorrections: Array<{ judgment: ThoughtEvent; output: JsonObject }> = []; if (originalOutput) { for (const judgment of judgmentSet.active.filter((event) => ( event.payload.runId === originalRun.id && event.payload.outputEventId === originalOutput.event.id && event.payload.kind === "correct" ))) { const candidate = objectField(judgment.payload.replacementOutput); if (!candidate) continue; try { directCorrections.push({ judgment, output: canonicalStructuredOutput(registry, contract, candidate), }); } catch { // Corrupt historical direct corrections remain inert. } } } directCorrections.sort((left, right) => compareEvents(left.judgment, right.judgment)); const selectedDirectCorrection = directCorrections.at(-1); const proposals = (await store.listEvents({ types: [CORRECTION_PROPOSAL_EVENT_TYPE] })) .filter((event) => event.payload.originalRunId === originalRunId); const acceptedRepairs: Array<{ judgment: ThoughtEvent; proposal: ThoughtEvent; repairRun: AgentRun; output: JsonObject; }> = [];
for (const proposal of proposals) { const repairRunId = optionalString(proposal.payload.repairRunId); if (!repairRunId) continue; const repairRun = await store.getRun(repairRunId); if (!repairRun || repairRun.status !== "completed" || !repairRun.outputEventIds.includes(proposal.id)) continue; if (!sameContract(proposal.payload.outputContract, contract)) continue; const judgments = judgmentSet.active.filter((judgment) => ( judgment.payload.runId === repairRun.id && (judgment.payload.kind === "accept" || judgment.payload.kind === "correct") )); for (const judgment of judgments) { const candidate = judgment.payload.kind === "correct" ? objectField(judgment.payload.replacementOutput) : objectField(proposal.payload.structuredOutput); if (!candidate) continue; try { acceptedRepairs.push({ judgment, proposal, repairRun, output: canonicalStructuredOutput(registry, contract, candidate), }); } catch { // A corrupt historical judgment cannot become effective merely because it exists. } } } acceptedRepairs.sort((left, right) => compareEvents(left.judgment, right.judgment)); const selectedRepair = acceptedRepairs.at(-1);
let payload: JsonObject; let lastEventId: string; if (selectedDirectCorrection && originalOutput) { payload = { originalRunId, status: "corrected", sourceRootEventId: originalSource.rootEventId, originalModel: executionProvenance(originalRun), privacy: joinPrivacy( originalSource.privacy, runPrivacy(originalRun), originalOutput.event.privacy, selectedDirectCorrection.judgment.privacy, ), outputContract: outputContractIdentityJson(contract), structuredOutput: selectedDirectCorrection.output, outputEventId: originalOutput.event.id, judgmentEventId: selectedDirectCorrection.judgment.id, ...(optionalString(selectedDirectCorrection.judgment.payload.feedbackSourceEventId) ? { feedbackSourceEventId: optionalString(selectedDirectCorrection.judgment.payload.feedbackSourceEventId)!, } : {}), authority: "correct", }; lastEventId = selectedDirectCorrection.judgment.id; } else if (selectedRepair) { payload = { originalRunId, status: "repair", sourceRootEventId: originalSource.rootEventId, originalModel: executionProvenance(originalRun), repairModel: executionProvenance(selectedRepair.repairRun), privacy: joinPrivacy( originalSource.privacy, runPrivacy(originalRun), selectedRepair.proposal.privacy, runPrivacy(selectedRepair.repairRun), selectedRepair.judgment.privacy, ), outputContract: outputContractIdentityJson(contract), structuredOutput: selectedRepair.output, outputEventId: selectedRepair.proposal.id, repairRunId: selectedRepair.repairRun.id, repairRequestEventId: optionalString(selectedRepair.proposal.payload.repairRequestEventId) ?? selectedRepair.repairRun.triggerEventId, judgmentEventId: selectedRepair.judgment.id, authority: selectedRepair.judgment.payload.kind === "correct" ? "correct" : "accept", }; lastEventId = selectedRepair.judgment.id; } else { if (originalOutput) { payload = { originalRunId, status: "original", sourceRootEventId: originalSource.rootEventId, originalModel: executionProvenance(originalRun), privacy: joinPrivacy(originalSource.privacy, runPrivacy(originalRun), originalOutput.event.privacy), outputContract: outputContractIdentityJson(contract), structuredOutput: originalOutput.output, outputEventId: originalOutput.event.id, }; lastEventId = originalOutput.event.id; } else { const failed = (await store.listEvents({ types: ["stream.thought.agent.run.failed"] })) .filter((event) => event.payload.runId === originalRunId) .at(-1); payload = { originalRunId, status: "unresolved", sourceRootEventId: originalSource.rootEventId, originalModel: executionProvenance(originalRun), privacy: joinPrivacy(originalSource.privacy, runPrivacy(originalRun), failed?.privacy), outputContract: outputContractIdentityJson(contract), }; lastEventId = failed?.id ?? originalSource.id; } }
const projection: Projection = { id: stableKey("effective-output", originalRunId), payload, lastEventId, projectionVersion: EFFECTIVE_OUTPUT_PROJECTION_VERSION, updatedAt: new Date().toISOString(), }; await store.upsertProjection(projection); return projection;}
function executionProvenance(run: AgentRun): JsonObject { return { provider: run.provider, id: run.model, ...(run.checkpointRevision ? { checkpointRevision: run.checkpointRevision } : {}), ...(run.executionAdapterRevision ? { executionAdapterRevision: run.executionAdapterRevision } : {}), ...(run.modelAdapter ? { modelAdapter: run.modelAdapter as unknown as JsonObject } : {}), ...(run.adapterCatalogDigest ? { adapterCatalogDigest: run.adapterCatalogDigest } : {}), ...(run.adapterCatalogGeneration ? { adapterCatalogGeneration: run.adapterCatalogGeneration } : {}), };}
export async function rebuildAllEffectiveOutputs(store: JazzThoughtStore): Promise<Projection[]> { const runs = await store.listRuns(); const originals = runs.filter((run) => run.contextManifest.agentRole !== "repair"); const projections: Projection[] = []; for (const run of originals) { try { projections.push(await rebuildEffectiveOutput(store, run.id)); } catch { // Runs predating contract-bound manifests are outside this projection version. } } return projections;}
async function validatedOriginalOutput( store: JazzThoughtStore, run: AgentRun, contract: OutputContractIdentity,): Promise<{ event: ThoughtEvent; output: JsonObject } | undefined> { if (run.status !== "completed" || run.outputEventIds.length !== 1) return undefined; const event = await store.getEvent(run.outputEventIds[0]!); if (!event || !sameContract(event.payload.outputContract, contract)) return undefined; const structuredOutput = objectField(event.payload.structuredOutput); if (!structuredOutput) return undefined; const output = canonicalStructuredOutput(createOutputContractRegistry(), contract, structuredOutput); return { event, output };}
function sameContract(value: unknown, expected: OutputContractIdentity): boolean { try { return canonicalJson(outputContractIdentityJson(parseOutputContractIdentity(value))) === canonicalJson(outputContractIdentityJson(expected)); } catch { return false; }}
async function requireEvent(store: JazzThoughtStore, id: string): Promise<ThoughtEvent> { const event = await store.getEvent(id); if (!event) throw new Error(`Event not found: ${id}`); return event;}
function objectField(value: unknown): JsonObject | undefined { return value && typeof value === "object" && !Array.isArray(value) ? value as JsonObject : undefined;}
function optionalString(value: unknown): string | undefined { return typeof value === "string" && value.length > 0 ? value : undefined;}
function compareEvents(left: ThoughtEvent, right: ThoughtEvent): number { return left.observedAt.localeCompare(right.observedAt) || left.id.localeCompare(right.id);}