import { 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 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; minimumPrivacy?: ThoughtEvent["privacy"] | undefined; } export class EventRegistry { private readonly types = new Map(); 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 = 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; 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; 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; 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; 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(); 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; 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; 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; 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; 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; 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; 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; 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; 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; 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, 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, 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, }); 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, 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, 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, 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, }); 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; }