Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
16 kB · 383 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384import { canonicalJson, sha256, type JsonObject } from "../core/json.js";import type { ThoughtEvent } from "../events/types.js";import type { JazzThoughtStore } from "../jazz/store.js";import type { AgentRun } from "../store/types.js";import { joinPrivacy, runPrivacy } from "../security/privacy.js";import { agentRole, declarationFingerprint } from "./declarations.js";import { createOutputContractRegistry, OBSERVATION_OUTPUT_CONTRACT_ID, outputContractForDeclaration, outputContractIdentityJson, parseOutputContractIdentity, type OutputContractRegistry, type SanitizedOutputIssue,} from "./output-contracts.js";import type { ThoughtAgentDeclaration } from "./types.js";
export const REPAIR_REQUEST_EVENT_TYPE = "stream.thought.agent.repair.requested";export const CORRECTION_PROPOSAL_EVENT_TYPE = "stream.thought.derived.output.correction.proposed";export const REPAIR_EVENT_SOURCE = "repair:output-contract";export const REPAIR_POLICY_ID = "stream.thought.agent.output-repair";export const REPAIR_POLICY_VERSION = 1;
export interface RepairPolicyIdentity { id: string; version: number; sha256: string;}
export interface RepairPolicy { identity: RepairPolicyIdentity; allowlistedSemanticValidationIds: ReadonlySet<string>; definition: JsonObject;}
export interface RepairReconciliationResult { examined: number; eligible: number; inserted: number; unchanged: number; requestEventIds: string[];}
const basePolicyDefinition = { id: REPAIR_POLICY_ID, version: REPAIR_POLICY_VERSION, eligibleReasons: [ "invalid-json", "expected-one-text-part", "empty-final-text", "final-text-too-large", "output-contract-invalid", ], allowlistedSemanticValidationIds: [], requireTerminalFailure: true, requireCompleteImmutableEvidence: true, excludeRepairOrigin: true, maxRequestsPerRunAndPolicy: 1,} satisfies JsonObject;
export const DEFAULT_REPAIR_POLICY: RepairPolicy = createRepairPolicy();
export function createRepairPolicy(allowlistedSemanticValidationIds: string[] = []): RepairPolicy { const definition: JsonObject = { ...basePolicyDefinition, allowlistedSemanticValidationIds: [...new Set(allowlistedSemanticValidationIds)].sort(), }; return { identity: { id: REPAIR_POLICY_ID, version: REPAIR_POLICY_VERSION, sha256: sha256(canonicalJson(definition)), }, allowlistedSemanticValidationIds: new Set(allowlistedSemanticValidationIds), definition, };}
export class RepairRequestCoordinator { private readonly outputContracts: OutputContractRegistry;
constructor( private readonly store: JazzThoughtStore, private readonly policy: RepairPolicy = DEFAULT_REPAIR_POLICY, outputContracts?: OutputContractRegistry, ) { this.outputContracts = outputContracts ?? createOutputContractRegistry(); }
async reconcile(declarations: ThoughtAgentDeclaration[]): Promise<RepairReconciliationResult> { const declarationsByVersion = new Map(declarations.map((declaration) => [ `${declaration.id}@${declaration.version}`, declaration, ])); const failedRuns = await this.store.listRuns({ status: "failed" }); const result: RepairReconciliationResult = { examined: failedRuns.length, eligible: 0, inserted: 0, unchanged: 0, requestEventIds: [], }; for (const run of failedRuns) { const declaration = declarationsByVersion.get(`${run.agentId}@${run.agentVersion}`); if (!declaration) continue; const request = await this.appendForRun(run, declaration); if (!request) continue; result.eligible += 1; result.requestEventIds.push(request.event.id); if (request.inserted) result.inserted += 1; else result.unchanged += 1; } return result; }
async appendForRun( run: AgentRun, declaration: ThoughtAgentDeclaration, ): Promise<{ event: ThoughtEvent; inserted: boolean } | undefined> { const evidence = await repairEvidence(this.store, run, declaration, this.policy, this.outputContracts); if (!evidence) return undefined; const occurredAt = new Date().toISOString(); const payload: JsonObject = { policy: policyIdentityJson(this.policy.identity), originalRunId: run.id, originalAgentId: declaration.id, originalAgentVersion: declaration.version, originalTriggerEventId: evidence.source.id, originalFailedEventId: evidence.failed.id, sourceRootEventId: evidence.source.rootEventId, declarationFingerprint: evidence.declarationFingerprint, promptRef: declaration.promptRef, promptSha256: run.promptHash, outputContract: outputContractIdentityJson(evidence.outputContract), model: evidence.model, failure: evidence.failure, evidenceSha256: evidence.evidenceSha256, }; return this.store.appendEvent({ type: REPAIR_REQUEST_EVENT_TYPE, schemaVersion: 1, source: REPAIR_EVENT_SOURCE, sourceKind: "agent", externalId: `${run.id}:${this.policy.identity.id}@${this.policy.identity.version}`, idempotencyKey: `${run.id}:${this.policy.identity.id}@${this.policy.identity.version}:${this.policy.identity.sha256}`, occurredAt, actor: "repair-coordinator", rootEventId: evidence.source.rootEventId, parentEventId: evidence.failed.id, correlationId: run.id, privacy: joinPrivacy(evidence.source.privacy, evidence.failed.privacy, runPrivacy(run)), payload, traceId: run.id, }); }}
interface EligibleRepairEvidence { source: ThoughtEvent; failed: ThoughtEvent; declarationFingerprint: string; outputContract: ReturnType<typeof outputContractForDeclaration>; model: JsonObject; failure: JsonObject; evidenceSha256: string;}
async function repairEvidence( store: JazzThoughtStore, run: AgentRun, declaration: ThoughtAgentDeclaration, policy: RepairPolicy, outputContracts: OutputContractRegistry,): Promise<EligibleRepairEvidence | undefined> { if (run.status !== "failed" || run.outputEventIds.length !== 0) return undefined; if (run.agentId !== declaration.id || run.agentVersion !== declaration.version) return undefined; if (agentRole(declaration) !== "standard") return undefined; if (run.inputEventIds.length !== 1 || run.inputEventIds[0] !== run.triggerEventId) return undefined;
const source = await store.getEvent(run.triggerEventId); if (!source || source.type === REPAIR_REQUEST_EVENT_TYPE || source.rootEventId.length === 0) return undefined; const expectedPromptHash = sha256(declaration.systemPrompt); if (run.promptHash !== expectedPromptHash) return undefined; const expectedFingerprint = declaration.declarationFingerprint ?? declarationFingerprint(declaration); if (stringField(run.contextManifest.declarationFingerprint) !== expectedFingerprint) return undefined; if (stringField(run.contextManifest.promptRef) !== declaration.promptRef) return undefined; if (stringField(run.contextManifest.promptRevision) !== expectedPromptHash) return undefined; if (stringField(run.contextManifest.agentRole) !== "standard") return undefined;
const outputContract = outputContractForDeclaration(declaration); if (outputContract.id !== OBSERVATION_OUTPUT_CONTRACT_ID) return undefined; try { outputContracts.resolve(outputContract); const manifestContract = parseOutputContractIdentity(run.contextManifest.outputContract); if (canonicalJson(outputContractIdentityJson(manifestContract)) !== canonicalJson(outputContractIdentityJson(outputContract))) return undefined; } catch { return undefined; }
const diagnostic = objectField(run.result, "failureDiagnostic"); const failure = eligibleFailureDiagnostic(diagnostic, outputContract, policy); if (!failure) return undefined; const diagnosticProvider = stringField(diagnostic?.provider); const diagnosticModel = stringField(diagnostic?.model); if ((diagnosticProvider && diagnosticProvider !== run.provider) || (diagnosticModel && diagnosticModel !== run.model)) return undefined; const checkpointRevision = stringField(diagnostic?.checkpointRevision); const model: JsonObject = { provider: run.provider, id: run.model, ...(checkpointRevision ? { checkpointRevision } : {}), ...(run.adapterRevision ? { adapterRevision: run.adapterRevision } : {}), ...(run.executionAdapterRevision ? { executionAdapterRevision: run.executionAdapterRevision } : {}), ...(run.modelAdapter ? { modelAdapter: run.modelAdapter as unknown as JsonObject } : {}), ...(run.adapterCatalogDigest ? { adapterCatalogDigest: run.adapterCatalogDigest } : {}), ...(run.adapterCatalogGeneration ? { adapterCatalogGeneration: run.adapterCatalogGeneration } : {}), }; const failedEvents = (await store.listEvents({ types: ["stream.thought.agent.run.failed"] })) .filter((event) => event.payload.runId === run.id); if (failedEvents.length !== 1) return undefined; const failed = failedEvents[0]!; const failedDiagnostic = objectField(failed.payload, "failureDiagnostic"); if (!failedDiagnostic || canonicalJson(failedDiagnostic) !== canonicalJson(diagnostic!)) return undefined; if ( failed.source !== `agent:${run.agentId}` || failed.actor !== run.agentId || failed.traceId !== run.id || failed.rootEventId !== source.rootEventId || failed.parentEventId !== source.id || failed.privacy !== joinPrivacy(source.privacy, runPrivacy(run)) || failed.payload.runId !== run.id || failed.payload.executionKey !== run.executionKey || failed.payload.agentId !== run.agentId || failed.payload.agentVersion !== run.agentVersion || failed.payload.agentRole !== "standard" || failed.payload.declarationFingerprint !== expectedFingerprint || failed.payload.status !== "failed" || failed.payload.attempt !== run.attempt || canonicalJson(failed.payload.inputEventIds as JsonObject) !== canonicalJson(run.inputEventIds as unknown as JsonObject) || canonicalJson(failed.payload.outputContract as JsonObject) !== canonicalJson(outputContractIdentityJson(outputContract)) ) return undefined;
const evidenceBody: JsonObject = { policy: policyIdentityJson(policy.identity), originalRunId: run.id, originalAgentId: run.agentId, originalAgentVersion: run.agentVersion, originalTriggerEventId: source.id, originalFailedEventId: failed.id, sourceRootEventId: source.rootEventId, declarationFingerprint: expectedFingerprint, promptRef: declaration.promptRef, promptSha256: expectedPromptHash, outputContract: outputContractIdentityJson(outputContract), model, failure, }; return { source, failed, declarationFingerprint: expectedFingerprint, outputContract, model, failure, evidenceSha256: sha256(canonicalJson(evidenceBody)), };}
function eligibleFailureDiagnostic( diagnostic: JsonObject | undefined, outputContract: ReturnType<typeof outputContractForDeclaration>, policy: RepairPolicy,): JsonObject | undefined { if (!diagnostic) return undefined; const code = stringField(diagnostic.code); const stage = stringField(diagnostic.stage); const reason = stringField(diagnostic.reason); if (stage !== "final-output-validation") return undefined; const contractValue = diagnostic.outputContract; try { const diagnosticContract = parseOutputContractIdentity(contractValue); if (canonicalJson(outputContractIdentityJson(diagnosticContract)) !== canonicalJson(outputContractIdentityJson(outputContract))) return undefined; } catch { return undefined; }
const assistantMessages = nonnegativeInteger(diagnostic.assistantMessages); const textParts = nonnegativeInteger(diagnostic.textParts); const textChars = nonnegativeInteger(diagnostic.textChars); const thinkingParts = nonnegativeInteger(diagnostic.thinkingParts); const thinkingChars = nonnegativeInteger(diagnostic.thinkingChars); const toolCallParts = nonnegativeInteger(diagnostic.toolCallParts); const otherParts = nonnegativeInteger(diagnostic.otherParts); if ([assistantMessages, textParts, textChars, thinkingParts, thinkingChars, toolCallParts, otherParts].some((value) => value === undefined)) return undefined; if (diagnostic.thinkingRedacted !== true) return undefined;
const base: JsonObject = { code: code ?? "", stage, reason: reason ?? "", assistantMessages: assistantMessages!, textParts: textParts!, textChars: textChars!, thinkingParts: thinkingParts!, thinkingChars: thinkingChars!, thinkingRedacted: true, toolCallParts: toolCallParts!, otherParts: otherParts!, ...(safeStopReason(diagnostic.stopReason) ? { stopReason: safeStopReason(diagnostic.stopReason)! } : {}), };
const textSha256 = stringField(diagnostic.textSha256); const requiresTextHash = reason === "invalid-json" || reason === "final-text-too-large" || reason === "output-contract-invalid" || reason === "semantic-output-invalid"; if (requiresTextHash && (!textSha256 || !/^[a-f0-9]{64}$/.test(textSha256))) return undefined; if (textSha256) base.textSha256 = textSha256;
if (code === "invalid-final-output" && [ "invalid-json", "expected-one-text-part", "empty-final-text", "final-text-too-large", "output-contract-invalid", ].includes(reason ?? "")) { if (reason === "invalid-json" && (textParts !== 1 || textChars === 0)) return undefined; if (reason === "empty-final-text" && (textParts !== 1 || textChars !== 0)) return undefined; if (reason === "output-contract-invalid") { const issues = sanitizedIssues(diagnostic.validationIssues); if (issues.length === 0) return undefined; base.validationIssues = issues.map((issue) => ({ code: issue.code, path: issue.path })); } return base; }
if (code === "semantic-output-invalid" && reason === "semantic-output-invalid") { const semanticValidationId = stringField(diagnostic.semanticValidationId); if (!semanticValidationId || !policy.allowlistedSemanticValidationIds.has(semanticValidationId)) return undefined; base.semanticValidationId = semanticValidationId; const issues = sanitizedIssues(diagnostic.validationIssues); if (issues.length > 0) base.validationIssues = issues.map((issue) => ({ code: issue.code, path: issue.path })); return base; } return undefined;}
function sanitizedIssues(value: unknown): SanitizedOutputIssue[] { if (!Array.isArray(value) || value.length > 20) return []; const issues: SanitizedOutputIssue[] = []; for (const entry of value) { if (!entry || typeof entry !== "object" || Array.isArray(entry)) return []; const record = entry as Record<string, unknown>; if (typeof record.code !== "string" || !record.code || record.code.length > 100 || !Array.isArray(record.path) || record.path.length > 12) return []; const path: Array<string | number> = []; for (const part of record.path) { if (typeof part === "number" && Number.isSafeInteger(part) && part >= 0) path.push(part); else if (typeof part === "string" && part.length <= 200) path.push(part); else return []; } issues.push({ code: record.code, path }); } return issues;}
function policyIdentityJson(identity: RepairPolicyIdentity): JsonObject { return { id: identity.id, version: identity.version, sha256: identity.sha256 };}
function objectField(value: JsonObject | undefined, key: string): JsonObject | undefined { const field = value?.[key]; return field && typeof field === "object" && !Array.isArray(field) ? field as JsonObject : undefined;}
function stringField(value: unknown): string | undefined { return typeof value === "string" && value.length > 0 ? value : undefined;}
function nonnegativeInteger(value: unknown): number | undefined { return Number.isSafeInteger(value) && Number(value) >= 0 ? Number(value) : undefined;}
function safeStopReason(value: unknown): string | undefined { return typeof value === "string" && ["stop", "length", "toolUse", "error", "aborted"].includes(value) ? value : undefined;}