Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
10.0 kB · 203 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204import { z } from "zod";import { REVIEW_RESPONSE_OUTPUT_CONTRACT, reviewResponseOutputSchema, reviewResponseSummary,} from "../agents/output-contracts.js";import { modelAdapterIdentitySchema } from "../adapters/model-adapters.js";import type { JsonObject } from "../core/json.js";
export const REVIEW_PROMPT_EVENT_TYPE = "stream.thought.source.review.prompt";export const REVIEW_ITEM_EVENT_TYPE = "stream.thought.review.item.created";export const REVIEW_DECISION_EVENT_TYPE = "stream.thought.review.decision";export const REVIEW_RESPONSE_EVENT_TYPE = "stream.thought.derived.review.response";
export const REVIEW_PROMPT_SCHEMA_VERSION = 1;export const REVIEW_ITEM_SCHEMA_VERSION = 1;export const REVIEW_DECISION_SCHEMA_VERSION = 1;
export const REVIEW_RUNTIME = "thoughtstream-review-v1";
export const reviewCriterionSchema = z.object({ id: z.string().min(1).max(200).regex(/^[a-z0-9][a-z0-9._-]*$/), version: z.number().int().positive(), label: z.string().min(1).max(200), instructions: z.string().min(1).max(10_000), reasonCodes: z.array(z.string().min(1).max(100).regex(/^[a-z0-9][a-z0-9._-]*$/)).max(30), responseTags: z.array(z.string().min(1).max(100).regex(/^[a-z0-9][a-z0-9._-]*$/)).max(30),}).strict().superRefine((value, context) => { addDuplicateIssues(value.reasonCodes, "reasonCodes", context); addDuplicateIssues(value.responseTags, "responseTags", context);});
export const reviewPromptPayloadSchema = z.object({ campaign: z.object({ id: z.string().min(1).max(200).regex(/^[a-z0-9][a-z0-9._-]*$/), version: z.number().int().positive(), label: z.string().min(1).max(200), }).strict(), prompt: z.string().min(1).max(64_000), evidence: z.string().min(1).max(128_000).optional(), criterion: reviewCriterionSchema, candidateAgentIds: z.array(z.string().min(1).max(200)).min(2).max(8), externalExportEligible: z.boolean(),}).strict().superRefine((value, context) => { addDuplicateIssues(value.candidateAgentIds, "candidateAgentIds", context);}) as unknown as z.ZodType<JsonObject>;
const contractIdentitySchema = z.object({ id: z.string().min(1), version: z.number().int().positive(), sha256: z.string().regex(/^[a-f0-9]{64}$/),}).strict();
export const reviewResponseEventPayloadSchema = z.object({ runId: z.string().min(1), executionKey: z.string().min(1), inputEventId: z.string().min(1), inputSourceSequence: z.number().int().positive(), summary: z.string().min(1).max(2_000), outputContract: contractIdentitySchema, structuredOutput: reviewResponseOutputSchema, model: z.record(z.string(), z.unknown()).optional(), executionAdapterRevision: z.string().min(1).max(500).optional(), modelAdapter: modelAdapterIdentitySchema.optional(), adapterCatalogDigest: z.string().regex(/^[a-f0-9]{64}$/).optional(), adapterCatalogGeneration: z.number().int().positive().optional(),}).strict().superRefine((value, context) => { const expected = REVIEW_RESPONSE_OUTPUT_CONTRACT.identity; if ( value.outputContract.id !== expected.id || value.outputContract.version !== expected.version || value.outputContract.sha256 !== expected.sha256 ) { context.addIssue({ code: "custom", path: ["outputContract"], message: "Review response must name the canonical output contract" }); } const expectedSummary = reviewResponseSummary(value.structuredOutput.response); if (value.summary !== expectedSummary) { context.addIssue({ code: "custom", path: ["summary"], message: "Review response summary must be the canonical bounded preview" }); }}) as unknown as z.ZodType<JsonObject>;
export const reviewItemPayloadSchema = z.object({ promptEventId: z.string().min(1), candidateRunIds: z.tuple([z.string().min(1), z.string().min(1)]), candidateOutputEventIds: z.tuple([z.string().min(1), z.string().min(1)]), candidateRunReceiptDigests: z.tuple([ z.string().regex(/^sha256:[a-f0-9]{64}$/), z.string().regex(/^sha256:[a-f0-9]{64}$/), ]), displayOrder: z.tuple([z.string().min(1), z.string().min(1)]), criterion: z.object({ id: z.string().min(1).max(200), version: z.number().int().positive(), }).strict(), outputContract: contractIdentitySchema,}).strict().superRefine((value, context) => { if (value.candidateRunIds[0] === value.candidateRunIds[1]) { context.addIssue({ code: "custom", path: ["candidateRunIds"], message: "Review candidates must be distinct" }); } if (value.candidateOutputEventIds[0] === value.candidateOutputEventIds[1]) { context.addIssue({ code: "custom", path: ["candidateOutputEventIds"], message: "Review outputs must be distinct" }); } if ( new Set(value.displayOrder).size !== 2 || !value.candidateRunIds.includes(value.displayOrder[0]) || !value.candidateRunIds.includes(value.displayOrder[1]) ) { context.addIssue({ code: "custom", path: ["displayOrder"], message: "Display order must be a permutation of candidate runs" }); }}) as unknown as z.ZodType<JsonObject>;
export const reviewDispositionSchema = z.enum([ "prefer", "tie", "correct", "underdetermined", "malformed", "skip",]);
export type ReviewDisposition = z.infer<typeof reviewDispositionSchema>;
export const reviewDecisionPayloadSchema = z.object({ reviewItemEventId: z.string().min(1), promptEventId: z.string().min(1), disposition: reviewDispositionSchema, preferredRunId: z.string().min(1).optional(), preferenceStrength: z.enum(["slight", "strong"]).optional(), replacementResponse: z.string().min(1).max(60_000).optional(), confidence: z.enum(["low", "medium", "high"]).optional(), reasonCodes: z.array(z.string().min(1).max(100)).max(2), responseTags: z.array(z.string().min(1).max(100)).max(2), notes: z.string().min(1).max(4_000).optional(), trainingEligible: z.boolean(), externalExportEligible: z.boolean(), submissionId: z.string().min(16).max(200).regex(/^[A-Za-z0-9._~-]+$/), supersedesDecisionEventId: z.string().min(1).optional(),}).strict().superRefine((value, context) => { addDuplicateIssues(value.reasonCodes, "reasonCodes", context); addDuplicateIssues(value.responseTags, "responseTags", context); if (value.disposition === "prefer") { if (!value.preferredRunId) context.addIssue({ code: "custom", path: ["preferredRunId"], message: "Preference requires one candidate" }); if (!value.preferenceStrength) context.addIssue({ code: "custom", path: ["preferenceStrength"], message: "Preference requires strength" }); } else if (value.preferredRunId || value.preferenceStrength) { context.addIssue({ code: "custom", path: ["preferredRunId"], message: "Only preference may identify a winning run" }); } if (value.disposition === "correct") { if (!value.replacementResponse) context.addIssue({ code: "custom", path: ["replacementResponse"], message: "Correction requires a replacement response" }); } else if (value.replacementResponse) { context.addIssue({ code: "custom", path: ["replacementResponse"], message: "Only correction may carry a replacement response" }); } if (value.trainingEligible && value.disposition !== "prefer" && value.disposition !== "correct") { context.addIssue({ code: "custom", path: ["trainingEligible"], message: "Only preference or correction may be used for training" }); } if (value.externalExportEligible && !value.trainingEligible) { context.addIssue({ code: "custom", path: ["externalExportEligible"], message: "External export requires training eligibility" }); }}) as unknown as z.ZodType<JsonObject>;
export const reviewDecisionSubmissionSchema = z.object({ disposition: reviewDispositionSchema, preferredCandidate: z.enum(["A", "B"]).optional(), preferenceStrength: z.enum(["slight", "strong"]).optional(), replacementResponse: z.string().min(1).max(60_000).optional(), confidence: z.enum(["low", "medium", "high"]).optional(), reasonCodes: z.array(z.string().min(1).max(100)).max(2).default([]), responseTags: z.array(z.string().min(1).max(100)).max(2).default([]), notes: z.string().min(1).max(4_000).optional(), trainingEligible: z.boolean(), submissionId: z.string().min(16).max(200).regex(/^[A-Za-z0-9._~-]+$/), supersedesDecisionEventId: z.string().min(1).optional(),}).strict().superRefine((value, context) => { addDuplicateIssues(value.reasonCodes, "reasonCodes", context); addDuplicateIssues(value.responseTags, "responseTags", context); if (value.disposition === "prefer") { if (!value.preferredCandidate) context.addIssue({ code: "custom", path: ["preferredCandidate"], message: "Preference requires one candidate" }); if (!value.preferenceStrength) context.addIssue({ code: "custom", path: ["preferenceStrength"], message: "Preference requires strength" }); } else if (value.preferredCandidate || value.preferenceStrength) { context.addIssue({ code: "custom", path: ["preferredCandidate"], message: "Only preference may identify a winning candidate" }); } if (value.disposition === "correct") { if (!value.replacementResponse) context.addIssue({ code: "custom", path: ["replacementResponse"], message: "Correction requires a replacement response" }); } else if (value.replacementResponse) { context.addIssue({ code: "custom", path: ["replacementResponse"], message: "Only correction may carry a replacement response" }); } if (value.trainingEligible && value.disposition !== "prefer" && value.disposition !== "correct") { context.addIssue({ code: "custom", path: ["trainingEligible"], message: "Only preference or correction may be used for training" }); }});
function addDuplicateIssues( values: string[], path: string, context: z.core.$RefinementCtx<unknown>,): void { const seen = new Set<string>(); for (let index = 0; index < values.length; index += 1) { if (seen.has(values[index]!)) { context.addIssue({ code: "custom", path: [path, index], message: "Values must be unique" }); } seen.add(values[index]!); }}