Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
32 kB · 788 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789import { z } from "zod";import { CONCEPTUALIZATION_OUTPUT_CONTRACT, conceptualizationOutputSchema, createOutputContractRegistry, PUBLIC_KNOWLEDGE_RECOMMENDATION_OUTPUT_CONTRACT, publicKnowledgeRecommendationOutputSchema,} from "../agents/output-contracts.js";import { sha256, type JsonObject, type JsonValue } from "../core/json.js";import { modelAdapterIdentitySchema } from "../adapters/model-adapters.js";import { OPERATIONAL_INCIDENT_EVENT_TYPE, operationalIncidentPayloadSchema } from "../incidents/types.js";import { AGENT_PROPOSAL_SCHEMA_VERSION, CORRECTION_PROPOSAL_EVENT_TYPE, MEMORY_MATERIALIZATION_FAILED_EVENT_TYPE, MEMORY_MATERIALIZED_EVENT_TYPE, MEMORY_PROPOSAL_EVENT_TYPE, PROPOSAL_DECISION_EVENT_TYPE, correctionProposalPayloadSchema as agentCorrectionProposalPayloadSchema, memoryMaterializationFailedPayloadSchema, memoryMaterializedPayloadSchema, memoryProposalPayloadSchema, proposalDecisionPayloadSchema,} from "../agent-proposals/contracts.js";import { artifactPayloadSchema, artifactRequestPayloadSchema } from "../artifacts/schema.js";import { ARTIFACT_EVENT_TYPE, ARTIFACT_REQUEST_EVENT_TYPE, ARTIFACT_REQUEST_SCHEMA_VERSION, ARTIFACT_SCHEMA_VERSION,} from "../artifacts/types.js";import { REVIEW_DECISION_EVENT_TYPE, REVIEW_DECISION_SCHEMA_VERSION, REVIEW_ITEM_EVENT_TYPE, REVIEW_ITEM_SCHEMA_VERSION, REVIEW_PROMPT_EVENT_TYPE, REVIEW_PROMPT_SCHEMA_VERSION, REVIEW_RESPONSE_EVENT_TYPE, reviewDecisionPayloadSchema, reviewItemPayloadSchema, reviewPromptPayloadSchema, reviewResponseEventPayloadSchema,} from "../review/types.js";import { PUBLIC_KNOWLEDGE_PROPOSED_DIFF_EVENT_TYPE, PUBLIC_KNOWLEDGE_PROPOSED_DIFF_SCHEMA_VERSION, publicKnowledgeProposedDiffPayloadSchema,} from "../public-knowledge/types.js";import { FOCUS_DECLARATION_PROPOSED_EVENT_TYPE, FOCUS_DECLARATION_SCHEMA_VERSION, focusDeclarationProposalPayloadSchema,} from "../focuses/types.js";import { CONVERSATION_COMPACTION_EVENT_TYPE, CONVERSATION_COMPACTION_SCHEMA_VERSION, conversationCompactionEventPayloadSchema,} from "../agents/conversation-compaction.js";import { AGENT_MESSAGE_RESPONSE_EVENT_TYPE, AGENT_MESSAGE_SCHEMA_VERSION, AGENT_MESSAGE_SOURCE_EVENT_TYPE, agentMessageResponsePayloadSchema, agentMessageSourcePayloadSchema,} from "../agents/agent-messages.js";import { COURSE_QUESTION_EVENT_TYPE, COURSE_QUESTION_SCHEMA_VERSION, courseQuestionPayloadSchema,} from "../courses/questions.js";import { CONTEXT_SELECTED_EVENT_TYPE as WORKBENCH_CONTEXT_SELECTED_EVENT_TYPE, DOCUMENT_CREATED_EVENT_TYPE as WORKBENCH_DOCUMENT_CREATED_EVENT_TYPE, DOCUMENT_VERSION_EVENT_TYPE as WORKBENCH_DOCUMENT_VERSION_EVENT_TYPE, PROPOSAL_DECISION_EVENT_TYPE as WORKBENCH_PROPOSAL_DECISION_EVENT_TYPE, PROPOSAL_PROPOSED_EVENT_TYPE as WORKBENCH_PROPOSAL_PROPOSED_EVENT_TYPE, PROPOSAL_REQUESTED_EVENT_TYPE as WORKBENCH_PROPOSAL_REQUESTED_EVENT_TYPE, WORKBENCH_SCHEMA_VERSION, contextSelectedPayloadSchema as workbenchContextSelectedPayloadSchema, documentCreatedPayloadSchema as workbenchDocumentCreatedPayloadSchema, documentVersionPayloadSchema as workbenchDocumentVersionPayloadSchema, proposalDecisionPayloadSchema as workbenchProposalDecisionPayloadSchema, proposalProposedPayloadSchema as workbenchProposalProposedPayloadSchema, proposalRequestedPayloadSchema as workbenchProposalRequestedPayloadSchema,} from "../workbench/contracts.js";import type { ThoughtEvent } from "./types.js";import { X_ACTIVITY_SOURCE_EVENT_TYPE, X_ACTIVITY_SOURCE_SCHEMA_VERSION, xActivityPayloadSchema,} from "../connectors/x-contract.js";
export interface RegisteredEventType { type: string; schemaVersion: number; description: string; payload: z.ZodType<JsonObject>; minimumPrivacy?: ThoughtEvent["privacy"] | undefined;}
export class EventRegistry { private readonly types = new Map<string, RegisteredEventType>();
register(definition: RegisteredEventType): this { const key = registryKey(definition.type, definition.schemaVersion); if (this.types.has(key)) throw new Error(`Event type already registered: ${key}`); this.types.set(key, definition); return this; }
validate(type: string, schemaVersion: number, payload: JsonObject): JsonObject { const definition = this.types.get(registryKey(type, schemaVersion)); if (!definition) throw new Error(`Unknown event type/version: ${type}@${schemaVersion}`); return definition.payload.parse(payload); }
validateEvent( type: string, schemaVersion: number, privacy: ThoughtEvent["privacy"], payload: JsonObject, ): JsonObject { const definition = this.types.get(registryKey(type, schemaVersion)); if (!definition) throw new Error(`Unknown event type/version: ${type}@${schemaVersion}`); if (definition.minimumPrivacy && privacyRank(privacy) < privacyRank(definition.minimumPrivacy)) { throw new Error(`${type}@${schemaVersion} requires ${definition.minimumPrivacy} or stricter privacy`); } return definition.payload.parse(payload); }
has(type: string, schemaVersion: number): boolean { return this.types.has(registryKey(type, schemaVersion)); }
list(): RegisteredEventType[] { return [...this.types.values()].sort((left, right) => left.type.localeCompare(right.type) || left.schemaVersion - right.schemaVersion, ); }}
const jsonValueSchema: z.ZodType<JsonValue> = z.lazy(() => z.union([ z.string(), z.number(), z.boolean(), z.null(), z.array(jsonValueSchema), z.record(z.string(), jsonValueSchema), ]),);
const objectPayload = z.record(z.string(), jsonValueSchema) as z.ZodType<JsonObject>;
const filePayload = z.object({ documentId: z.string().min(1), path: z.string().min(1), versionId: z.string().optional(), previousVersionId: z.string().optional(), previousPath: z.string().optional(), sha256: z.string().optional(), previousSha256: z.string().optional(), contentType: z.string(), sizeBytes: z.number().int().nonnegative(), mtimeMs: z.number().nonnegative(), diff: z.string().optional(), identityConfidence: z.enum(["explicit", "inode", "content", "path", "uncertain"]),}) as unknown as z.ZodType<JsonObject>;
const runPayload = z.object({ runId: z.string().min(1), agentId: z.string().min(1), agentVersion: z.number().int().positive(), inputEventIds: z.array(z.string()), attempt: z.number().int().positive(), status: z.string(),}).passthrough() as unknown as z.ZodType<JsonObject>;
const sha256Schema = z.string().regex(/^[a-f0-9]{64}$/);const telegramCorrectionPayload = z.object({ accountId: z.string().min(1).max(100), accountUsername: z.string().min(1).max(200).optional(), updateId: z.string().min(1).max(100), chatId: z.string().min(1).max(100), messageId: z.string().min(1).max(100), senderId: z.string().min(1).max(100), chatType: z.literal("direct"), occurredAt: z.iso.datetime(), commandVersion: z.literal("telegram-correct-v1"), replacementText: z.string().min(1).max(2_000).optional(), replacementChars: z.number().int().nonnegative(), replacementSha256: sha256Schema, replyToMessageId: z.string().min(1).max(100).optional(), resolutionStatus: z.enum([ "resolved", "invalid-command", "replacement-too-long", "missing-reply-target", "unknown-delivery", "ambiguous-delivery", "ambiguous-run", "missing-run", "incomplete-run", "ambiguous-output", "missing-output", "invalid-lineage", ]), deliveryReceiptEventId: z.string().min(1).optional(), runId: z.string().min(1).optional(), outputEventId: z.string().min(1).optional(), sourceRootEventId: z.string().min(1).optional(), transport: z.literal("telegram-bot-api-webhook"),}).strict().superRefine((value, context) => { if (!["invalid-command", "replacement-too-long"].includes(value.resolutionStatus) && value.replacementText === undefined) { context.addIssue({ code: "custom", path: ["replacementText"], message: "Admitted correction must retain its exact bounded replacement" }); } if (value.replacementText !== undefined && value.replacementText.length !== value.replacementChars) { context.addIssue({ code: "custom", path: ["replacementChars"], message: "Correction replacement length does not match" }); } if (value.replacementText !== undefined && sha256(value.replacementText) !== value.replacementSha256) { context.addIssue({ code: "custom", path: ["replacementSha256"], message: "Correction replacement hash does not match" }); } if (value.resolutionStatus === "resolved" && ( !value.replyToMessageId || !value.deliveryReceiptEventId || !value.runId || !value.outputEventId || !value.sourceRootEventId )) { context.addIssue({ code: "custom", path: ["resolutionStatus"], message: "Resolved correction requires complete target lineage" }); }}) as unknown as z.ZodType<JsonObject>;const batchMemberReferenceSchema = z.object({ eventId: z.string().min(1), source: z.string().min(1).max(200), sourceSequence: z.number().int().positive(), type: z.string().min(1).max(300), schemaVersion: z.number().int().positive(), privacy: z.enum(["public-source", "private", "sensitive"]), occurredAt: z.iso.datetime(), observedAt: z.iso.datetime(), payloadHash: sha256Schema,}).strict();const batchPayloadSchema = z.object({ declaration: z.object({ id: z.string().min(1).max(200), version: z.number().int().positive(), fingerprint: sha256Schema, }).strict(), flushReason: z.enum(["quiet-window", "max-age", "max-items"]), firstOccurredAt: z.iso.datetime(), lastOccurredAt: z.iso.datetime(), members: z.array(batchMemberReferenceSchema).min(1).max(1_000),}).strict().superRefine((value, context) => { const lastSequenceBySource = new Map<string, number>(); for (const [index, current] of value.members.entries()) { const previousSequence = lastSequenceBySource.get(current.source); if (previousSequence !== undefined && previousSequence >= current.sourceSequence) { context.addIssue({ code: "custom", path: ["members", index], message: "Batch members must preserve one source's increasing sequence" }); } lastSequenceBySource.set(current.source, current.sourceSequence); }}) as unknown as z.ZodType<JsonObject>;const outputContractIdentitySchema = z.object({ id: z.string().min(1), version: z.number().int().positive(), sha256: sha256Schema,}).strict();const repairPolicyIdentitySchema = z.object({ id: z.string().min(1), version: z.number().int().positive(), sha256: sha256Schema,}).strict();const validationIssueSchema = z.object({ code: z.string().min(1).max(100), path: z.array(z.union([z.string().max(200), z.number().int().nonnegative()])).max(12),}).strict();const repairModelIdentitySchema = z.object({ provider: z.string().min(1).max(200), id: z.string().min(1).max(500), checkpointRevision: z.string().min(1).max(500).optional(), adapterRevision: z.string().min(1).max(500).optional(), executionAdapterRevision: z.string().min(1).max(500).optional(), modelAdapter: modelAdapterIdentitySchema.optional(), adapterCatalogDigest: sha256Schema.optional(), adapterCatalogGeneration: z.number().int().positive().optional(),}).strict();const repairFailureSchema = z.object({ code: z.enum(["invalid-final-output", "semantic-output-invalid"]), stage: z.literal("final-output-validation"), reason: z.enum([ "invalid-json", "expected-one-text-part", "empty-final-text", "final-text-too-large", "output-contract-invalid", "semantic-output-invalid", ]), assistantMessages: z.number().int().nonnegative(), textParts: z.number().int().nonnegative(), textChars: z.number().int().nonnegative(), textSha256: sha256Schema.optional(), thinkingParts: z.number().int().nonnegative(), thinkingChars: z.number().int().nonnegative(), thinkingRedacted: z.literal(true), toolCallParts: z.number().int().nonnegative(), otherParts: z.number().int().nonnegative(), stopReason: z.enum(["stop", "length", "toolUse", "error", "aborted"]).optional(), validationIssues: z.array(validationIssueSchema).max(20).optional(), semanticValidationId: z.string().min(1).max(200).optional(),}).strict();const repairRequestPayload = z.object({ policy: repairPolicyIdentitySchema, originalRunId: z.string().min(1), originalAgentId: z.string().min(1), originalAgentVersion: z.number().int().positive(), originalTriggerEventId: z.string().min(1), originalFailedEventId: z.string().min(1), sourceRootEventId: z.string().min(1), declarationFingerprint: sha256Schema, promptRef: z.string().min(1), promptSha256: sha256Schema, outputContract: outputContractIdentitySchema, model: repairModelIdentitySchema, failure: repairFailureSchema, evidenceSha256: sha256Schema,}).strict() as unknown as z.ZodType<JsonObject>;const correctionProposalPayload = z.object({ originalRunId: z.string().min(1), repairRequestEventId: z.string().min(1), repairRunId: z.string().min(1), originalTriggerEventId: z.string().min(1), sourceRootEventId: z.string().min(1), outputContract: outputContractIdentitySchema, structuredOutput: objectPayload, originalModel: repairModelIdentitySchema, repairModel: repairModelIdentitySchema, executionAdapterRevision: z.string().min(1).max(500).optional(), modelAdapter: modelAdapterIdentitySchema.optional(), adapterCatalogDigest: sha256Schema.optional(), adapterCatalogGeneration: z.number().int().positive().optional(),}).strict().superRefine((value, context) => { try { createOutputContractRegistry().canonicalize(value.outputContract, value.structuredOutput); } catch { context.addIssue({ code: "custom", path: ["structuredOutput"], message: "Correction proposal must satisfy its registered output contract", }); }}) as unknown as z.ZodType<JsonObject>;
const conceptualizationGraphPayload = 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), confidence: z.number().min(0).max(1), outputContract: outputContractIdentitySchema, structuredOutput: conceptualizationOutputSchema, model: objectPayload.optional(), executionAdapterRevision: z.string().min(1).max(500).optional(), modelAdapter: modelAdapterIdentitySchema.optional(), adapterCatalogDigest: sha256Schema.optional(), adapterCatalogGeneration: z.number().int().positive().optional(),}).strict().superRefine((value, context) => { const expectedContract = CONCEPTUALIZATION_OUTPUT_CONTRACT.identity; if ( value.outputContract.id !== expectedContract.id || value.outputContract.version !== expectedContract.version || value.outputContract.sha256 !== expectedContract.sha256 ) { context.addIssue({ code: "custom", path: ["outputContract"], message: "Concept graph must name the canonical conceptualization contract identity", }); } else { try { createOutputContractRegistry().canonicalize(value.outputContract, value.structuredOutput); } catch { context.addIssue({ code: "custom", path: ["structuredOutput"], message: "Concept graph must satisfy its canonical conceptualization contract", }); } } if (value.summary !== value.structuredOutput.summary) { context.addIssue({ code: "custom", path: ["summary"], message: "Summary must equal canonical structured output" }); } if (value.confidence !== value.structuredOutput.confidence) { context.addIssue({ code: "custom", path: ["confidence"], message: "Confidence must equal canonical structured output" }); }}) as unknown as z.ZodType<JsonObject>;
const publicKnowledgeRecommendationPayload = 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(1_600), confidence: z.enum(["low", "medium", "high"]), outputContract: outputContractIdentitySchema, structuredOutput: publicKnowledgeRecommendationOutputSchema, model: objectPayload.optional(), executionAdapterRevision: z.string().min(1).max(500).optional(),}).strict().superRefine((value, context) => { const expected = PUBLIC_KNOWLEDGE_RECOMMENDATION_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: "Recommendation must name the canonical Public Knowledge contract" }); } if (value.summary !== value.structuredOutput.summary) { context.addIssue({ code: "custom", path: ["summary"], message: "Summary must equal canonical structured output" }); } if (value.confidence !== value.structuredOutput.confidence) { context.addIssue({ code: "custom", path: ["confidence"], message: "Confidence must equal canonical structured output" }); }}) as unknown as z.ZodType<JsonObject>;
const lettaConversationBindingPayload = z.object({ id: z.string().min(1), declarationId: z.string().min(1), declarationVersion: z.number().int().positive(), agentId: z.string().min(1), source: z.string().min(1), scopeType: z.literal("document"), scopeKey: z.string().min(1), conversationId: z.string().min(1), remoteMarker: z.string().regex(/^thoughtstream:coil:[a-f0-9]{64}$/), currentPath: z.string().min(1), createdAt: z.iso.datetime(), updatedAt: z.iso.datetime(),}).strict() as unknown as z.ZodType<JsonObject>;
const legacyJudgmentPayload = z.object({ runId: z.string().min(1), outputEventId: z.string().min(1), kind: z.enum(["accept", "reject", "correct", "prefer"]), criterion: z.string().min(1), criterionVersion: z.number().int().positive(), exportEligible: z.boolean(), comparedRunId: z.string().min(1).optional(), comparedOutputEventId: z.string().min(1).optional(), replacementOutput: objectPayload.optional(), notes: z.string().max(10_000).optional(), feedbackSourceEventId: z.string().min(1).optional(), deliveryReceiptEventId: z.string().min(1).optional(), supersedesJudgmentEventId: z.string().min(1).optional(),}).strict().superRefine((value, context) => { if (value.kind === "correct" && !value.replacementOutput) { context.addIssue({ code: "custom", path: ["replacementOutput"], message: "Corrections require a replacement output" }); } if (value.kind === "prefer" && (!value.comparedRunId || !value.comparedOutputEventId)) { context.addIssue({ code: "custom", path: ["comparedRunId"], message: "Preferences require a compared run and output" }); }}) as unknown as z.ZodType<JsonObject>;
const judgmentPayload = z.object({ runId: z.string().min(1), outputEventId: z.string().min(1), kind: z.enum(["accept", "reject", "correct", "prefer"]), criterion: z.string().min(1), criterionVersion: z.number().int().positive(), qualityEligible: z.boolean(), externalExportEligible: z.boolean(), comparedRunId: z.string().min(1).optional(), comparedOutputEventId: z.string().min(1).optional(), replacementOutput: objectPayload.optional(), notes: z.string().max(10_000).optional(), feedbackSourceEventId: z.string().min(1).optional(), deliveryReceiptEventId: z.string().min(1).optional(), supersedesJudgmentEventId: z.string().min(1).optional(),}).strict().superRefine((value, context) => { if (value.kind === "correct" && !value.replacementOutput) { context.addIssue({ code: "custom", path: ["replacementOutput"], message: "Corrections require a replacement output" }); } if (value.kind === "prefer" && (!value.comparedRunId || !value.comparedOutputEventId)) { context.addIssue({ code: "custom", path: ["comparedRunId"], message: "Preferences require a compared run and output" }); } if (value.externalExportEligible && !value.qualityEligible) { context.addIssue({ code: "custom", path: ["externalExportEligible"], message: "External export eligibility requires quality eligibility", }); }}) as unknown as z.ZodType<JsonObject>;
const judgmentRetractionPayload = z.object({ runId: z.string().min(1), outputEventId: z.string().min(1), criterion: z.string().min(1), criterionVersion: z.number().int().positive(), retractedJudgmentEventId: z.string().min(1), feedbackSourceEventId: z.string().min(1), deliveryReceiptEventId: z.string().min(1),}) as unknown as z.ZodType<JsonObject>;
export function createDefaultRegistry(): EventRegistry { const registry = new EventRegistry(); for (const type of [ "stream.thought.source.file.added", "stream.thought.source.file.changed", "stream.thought.source.file.renamed", "stream.thought.source.file.deleted", ]) { registry.register({ type, schemaVersion: 1, description: "Filesystem document observation", payload: filePayload }); }
for (const type of [ "stream.thought.agent.run.started", "stream.thought.agent.run.trace", "stream.thought.agent.run.completed", "stream.thought.agent.run.failed", "stream.thought.agent.run.blocked", "stream.thought.agent.run.skipped", "stream.thought.agent.run.abandoned", ]) { registry.register({ type, schemaVersion: 1, description: "Agent lifecycle evidence", payload: runPayload }); }
registry.register({ type: "stream.thought.agent.repair.requested", schemaVersion: 1, description: "Deterministic request for one sandboxed output-correction proposal", payload: repairRequestPayload, }); registry.register({ type: "stream.thought.derived.output.correction.proposed", schemaVersion: 1, description: "Inert contract-valid output correction proposal", payload: correctionProposalPayload, }); registry.register({ type: MEMORY_PROPOSAL_EVENT_TYPE, schemaVersion: AGENT_PROPOSAL_SCHEMA_VERSION, description: "Inert snapshot-bound agent memory-change proposal", payload: memoryProposalPayloadSchema, minimumPrivacy: "sensitive", }); registry.register({ type: CORRECTION_PROPOSAL_EVENT_TYPE, schemaVersion: AGENT_PROPOSAL_SCHEMA_VERSION, description: "Inert target-bound agent self-correction proposal", payload: agentCorrectionProposalPayloadSchema, minimumPrivacy: "sensitive", }); registry.register({ type: PROPOSAL_DECISION_EVENT_TYPE, schemaVersion: AGENT_PROPOSAL_SCHEMA_VERSION, description: "Append-only human decision over one exact agent proposal", payload: proposalDecisionPayloadSchema as unknown as z.ZodType<JsonObject>, minimumPrivacy: "sensitive", }); registry.register({ type: MEMORY_MATERIALIZED_EVENT_TYPE, schemaVersion: AGENT_PROPOSAL_SCHEMA_VERSION, description: "Verified stale-checked materialization of an accepted Stream memory proposal", payload: memoryMaterializedPayloadSchema, minimumPrivacy: "sensitive", }); registry.register({ type: MEMORY_MATERIALIZATION_FAILED_EVENT_TYPE, schemaVersion: AGENT_PROPOSAL_SCHEMA_VERSION, description: "Content-dark failed materialization of an accepted Stream memory proposal", payload: memoryMaterializationFailedPayloadSchema, minimumPrivacy: "sensitive", });
registry.register({ type: "stream.thought.judgment.training-example", schemaVersion: 1, description: "Legacy explicit judgment with one export-eligibility field", payload: legacyJudgmentPayload, }); registry.register({ type: "stream.thought.judgment.training-example", schemaVersion: 2, description: "Explicit judgment with separate quality and external-export eligibility", payload: judgmentPayload, }); registry.register({ type: "stream.thought.judgment.training-example.retracted", schemaVersion: 1, description: "Append-only retraction of an earlier training judgment", payload: judgmentRetractionPayload, });
registry.register({ type: REVIEW_PROMPT_EVENT_TYPE, schemaVersion: REVIEW_PROMPT_SCHEMA_VERSION, description: "Complete evidence-bearing prompt for one private review campaign", payload: reviewPromptPayloadSchema, }); registry.register({ type: REVIEW_RESPONSE_EVENT_TYPE, schemaVersion: 1, description: "Canonical candidate response generated for one review prompt", payload: reviewResponseEventPayloadSchema, }); registry.register({ type: REVIEW_ITEM_EVENT_TYPE, schemaVersion: REVIEW_ITEM_SCHEMA_VERSION, description: "Immutable blinded pair of exact completed review candidate runs", payload: reviewItemPayloadSchema, }); registry.register({ type: REVIEW_DECISION_EVENT_TYPE, schemaVersion: REVIEW_DECISION_SCHEMA_VERSION, description: "Append-only human judgeability, preference, correction, tie, or skip decision", payload: reviewDecisionPayloadSchema, });
registry.register({ type: "stream.thought.derived.event.batch", schemaVersion: 1, description: "Deterministic replayable batch of strong canonical event references", payload: batchPayloadSchema, });
registry.register({ type: "stream.thought.derived.concept.graph", schemaVersion: 1, description: "Validated private concept graph derived from one bounded source event", payload: conceptualizationGraphPayload, minimumPrivacy: "private", });
registry.register({ type: CONVERSATION_COMPACTION_EVENT_TYPE, schemaVersion: CONVERSATION_COMPACTION_SCHEMA_VERSION, description: "Private recursive conversation-history compaction boundary over a frozen canonical prefix", payload: conversationCompactionEventPayloadSchema, minimumPrivacy: "sensitive", });
registry.register({ type: "stream.thought.derived.public-knowledge.recommendation", schemaVersion: 1, description: "Validated private editorial recommendation for one exact Coil document version", payload: publicKnowledgeRecommendationPayload, minimumPrivacy: "sensitive", });
registry.register({ type: PUBLIC_KNOWLEDGE_PROPOSED_DIFF_EVENT_TYPE, schemaVersion: PUBLIC_KNOWLEDGE_PROPOSED_DIFF_SCHEMA_VERSION, description: "Inert context-bound private Public Knowledge diff proposal for one exact Coil document version", payload: publicKnowledgeProposedDiffPayloadSchema, minimumPrivacy: "sensitive", });
registry.register({ type: FOCUS_DECLARATION_PROPOSED_EVENT_TYPE, schemaVersion: FOCUS_DECLARATION_SCHEMA_VERSION, description: "Inert versioned focus-declaration proposal", payload: focusDeclarationProposalPayloadSchema as unknown as z.ZodType<JsonObject>, minimumPrivacy: "sensitive", });
registry.register({ type: "stream.thought.runtime.letta-conversation.binding", schemaVersion: 1, description: "Immutable evidence for one recoverable file-scoped Letta conversation binding", payload: lettaConversationBindingPayload, minimumPrivacy: "sensitive", });
registry.register({ type: OPERATIONAL_INCIDENT_EVENT_TYPE, schemaVersion: 1, description: "Content-dark operational incident or recovery evidence", payload: operationalIncidentPayloadSchema as z.ZodType<JsonObject>, }); registry.register({ type: AGENT_MESSAGE_SOURCE_EVENT_TYPE, schemaVersion: AGENT_MESSAGE_SCHEMA_VERSION, description: "Private typed message from one admitted agent to another", payload: agentMessageSourcePayloadSchema as unknown as z.ZodType<JsonObject>, minimumPrivacy: "sensitive", }); registry.register({ type: AGENT_MESSAGE_RESPONSE_EVENT_TYPE, schemaVersion: AGENT_MESSAGE_SCHEMA_VERSION, description: "Private typed response to an admitted agent message", payload: agentMessageResponsePayloadSchema as unknown as z.ZodType<JsonObject>, minimumPrivacy: "sensitive", }); registry.register({ type: COURSE_QUESTION_EVENT_TYPE, schemaVersion: COURSE_QUESTION_SCHEMA_VERSION, description: "Private authenticated question about one exact course lesson revision", payload: courseQuestionPayloadSchema as unknown as z.ZodType<JsonObject>, minimumPrivacy: "sensitive", }); registry.register({ type: "stream.thought.source.telegram.correction", schemaVersion: 1, description: "Target-bound private Telegram correction feedback", payload: telegramCorrectionPayload, minimumPrivacy: "sensitive", }); registry.register({ type: X_ACTIVITY_SOURCE_EVENT_TYPE, schemaVersion: X_ACTIVITY_SOURCE_SCHEMA_VERSION, description: "Signature-verified allowlisted X Activity webhook observation", payload: xActivityPayloadSchema, }); for (const [type, payload, description] of [ [WORKBENCH_DOCUMENT_CREATED_EVENT_TYPE, workbenchDocumentCreatedPayloadSchema, "Operator-owned working document created from one origin event"], [WORKBENCH_DOCUMENT_VERSION_EVENT_TYPE, workbenchDocumentVersionPayloadSchema, "Immutable working-document version appended through the trusted workbench path"], [WORKBENCH_CONTEXT_SELECTED_EVENT_TYPE, workbenchContextSelectedPayloadSchema, "Exact operator-selected evidence snapshot for one working document"], [WORKBENCH_PROPOSAL_REQUESTED_EVENT_TYPE, workbenchProposalRequestedPayloadSchema, "Operator request for one runner proposal over one exact base version"], [WORKBENCH_PROPOSAL_PROPOSED_EVENT_TYPE, workbenchProposalProposedPayloadSchema, "Inert runner-proposed replacement for one exact working-document base version"], [WORKBENCH_PROPOSAL_DECISION_EVENT_TYPE, workbenchProposalDecisionPayloadSchema, "Append-only human accept/reject decision over one workbench proposal"], ] as const) { registry.register({ type, schemaVersion: WORKBENCH_SCHEMA_VERSION, description, payload: payload as unknown as z.ZodType<JsonObject>, minimumPrivacy: "sensitive", }); } registry.register({ type: ARTIFACT_REQUEST_EVENT_TYPE, schemaVersion: ARTIFACT_REQUEST_SCHEMA_VERSION, description: "Append-only private request to materialize one workspace file", payload: artifactRequestPayloadSchema, minimumPrivacy: "private", }); registry.register({ type: ARTIFACT_EVENT_TYPE, schemaVersion: ARTIFACT_SCHEMA_VERSION, description: "Immutable private metadata reference to a content-addressed durable artifact", payload: artifactPayloadSchema, minimumPrivacy: "private", });
for (const type of [ "stream.thought.derived.topics", "stream.thought.derived.entities", "stream.thought.derived.message.observation", "stream.thought.derived.post.candidate", "stream.thought.derived.retrieval.recommendation", "stream.thought.derived.document.structure", "stream.thought.derived.document.read", "stream.thought.derived.charter.reevaluation", "stream.thought.source.rss.item", "stream.thought.source.atproto.commit", "stream.thought.source.email.observed", "stream.thought.source.telegram.message", "stream.thought.source.telegram.nonconversation", "stream.thought.source.telegram.reaction", "stream.thought.runtime.notice", "stream.thought.dispatcher.activated", "stream.thought.connector.poll.started", "stream.thought.connector.poll.completed", "stream.thought.connector.ingest.started", "stream.thought.connector.ingest.completed", "stream.thought.connector.subscription.started", "stream.thought.connector.subscription.connected", "stream.thought.connector.subscription.stopped", "stream.thought.connector.cursor.advanced", "stream.thought.connector.failed", "stream.thought.connector.recovered", "stream.thought.action.telegram.send.started", "stream.thought.action.telegram.send.delivered", "stream.thought.action.telegram.send.failed", ]) { registry.register({ type, schemaVersion: 1, description: "Registered thought stream payload", payload: objectPayload }); } return registry;}
function registryKey(type: string, schemaVersion: number): string { return `${type}@${schemaVersion}`;}
function privacyRank(privacy: ThoughtEvent["privacy"]): number { return privacy === "public-source" ? 0 : privacy === "private" ? 1 : 2;}