import { 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; } export async function activeJudgments(store: JazzThoughtStore): Promise { 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(); 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 { 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 { 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 { 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 { 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 { 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); }