Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
6.1 kB · 162 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163import { z } from "zod";import { canonicalJson, sha256, type JsonObject } from "../core/json.js";import type { EventCandidate, ThoughtEvent } from "../events/types.js";import type { JazzThoughtStore } from "../jazz/store.js";import { observationOutputSchema } from "./output-contracts.js";
export const AGENT_MESSAGE_SOURCE_EVENT_TYPE = "stream.thought.source.agent.message";export const AGENT_MESSAGE_RESPONSE_EVENT_TYPE = "stream.thought.derived.agent.message";export const AGENT_MESSAGE_SCHEMA_VERSION = 1;export const CO_AGENT_ID = "co";export const STREAM_AGENT_ENDPOINT_ID = "stream-agent-conversation";export const AGENT_MESSAGE_TEXT_MAX_CHARS = 8_000;
const agentIdSchema = z.string().min(1).max(100).regex(/^[a-z][a-z0-9]*(?:-[a-z0-9]+)*$/);const routeIdSchema = z.string().min(1).max(200).regex(/^[A-Za-z0-9][A-Za-z0-9._:-]*$/);const sha256Schema = z.string().regex(/^[a-f0-9]{64}$/);
export const agentMessageSourcePayloadSchema = z.object({ messageId: routeIdSchema, threadId: routeIdSchema, senderAgentId: agentIdSchema, recipientAgentId: agentIdSchema, text: z.string().min(1).max(AGENT_MESSAGE_TEXT_MAX_CHARS),}).strict();
export type AgentMessageSourcePayload = z.infer<typeof agentMessageSourcePayloadSchema>;
const outputContractIdentitySchema = z.object({ id: z.string().min(1), version: z.number().int().positive(), sha256: sha256Schema,}).strict();
const modelSchema = z.object({ provider: z.string().min(1).max(100), id: z.string().min(1).max(500), revision: z.string().min(1).max(500).optional(),}).strict();
export const agentMessageResponsePayloadSchema = z.object({ messageId: z.string().min(1).max(500), inReplyToMessageId: routeIdSchema, threadId: routeIdSchema, senderAgentId: agentIdSchema, recipientAgentId: agentIdSchema, runId: z.string().min(1), executionKey: z.string().min(1), inputEventId: z.string().min(1), inputSourceSequence: z.number().int().positive(), agentVersion: z.number().int().positive(), declarationFingerprint: sha256Schema, summary: z.string().min(1).max(2_000), tags: z.array(z.string().min(1).max(100)).max(20), importance: z.enum(["low", "normal", "high"]), confidence: z.number().min(0).max(1), outputContract: outputContractIdentitySchema, structuredOutput: observationOutputSchema, model: modelSchema.optional(),}).strict();
export type AgentMessageResponsePayload = z.infer<typeof agentMessageResponsePayloadSchema>;
export class AgentMessageNotAdmitted extends Error { readonly code = "agent-message-not-admitted"; readonly evidence: JsonObject;
constructor(event: ThoughtEvent) { super("Agent message is outside the endpoint route"); this.name = "AgentMessageNotAdmitted"; this.evidence = { eventType: event.type, source: event.source, sourceKind: event.sourceKind, privacy: event.privacy, }; }}
export function assertAgentMessageRoute( event: ThoughtEvent, recipientAgentId = STREAM_AGENT_ENDPOINT_ID,): AgentMessageSourcePayload { if (event.type !== AGENT_MESSAGE_SOURCE_EVENT_TYPE || event.schemaVersion !== AGENT_MESSAGE_SCHEMA_VERSION || event.sourceKind !== "agent" || event.privacy !== "sensitive") { throw new Error("Agent message event type, source kind, or privacy is invalid"); } const payload = agentMessageSourcePayloadSchema.parse(event.payload); if (payload.recipientAgentId !== recipientAgentId || event.source !== `agent-message:${payload.senderAgentId}` || event.actor !== `agent:${payload.senderAgentId}` || event.externalId !== payload.messageId || event.correlationId !== payload.threadId) { throw new Error("Agent message route identity is inconsistent"); } return payload;}
export async function appendCoAgentMessage( store: JazzThoughtStore, input: { messageId: string; threadId: string; text: string; occurredAt?: string | undefined; },) { const payload = agentMessageSourcePayloadSchema.parse({ messageId: input.messageId, threadId: input.threadId, senderAgentId: CO_AGENT_ID, recipientAgentId: STREAM_AGENT_ENDPOINT_ID, text: input.text, }); const idempotencyKey = [ "agent-message-v1", payload.senderAgentId, payload.recipientAgentId, payload.threadId, payload.messageId, ].join(":"); const source = `agent-message:${CO_AGENT_ID}`; const eventId = `evt_${sha256(`${source}\u0000${idempotencyKey}`).slice(0, 48)}`; const existing = await store.getEvent(eventId); if (existing) { const actual = assertAgentMessageRoute(existing); if (canonicalJson(actual as unknown as JsonObject) !== canonicalJson(payload as unknown as JsonObject)) { throw new Error("Existing agent message identity has divergent content"); } return { event: existing, inserted: false }; } const occurredAt = input.occurredAt ?? new Date().toISOString(); const candidate: EventCandidate = { type: AGENT_MESSAGE_SOURCE_EVENT_TYPE, schemaVersion: AGENT_MESSAGE_SCHEMA_VERSION, source, sourceKind: "agent", externalId: payload.messageId, idempotencyKey, occurredAt, actor: `agent:${CO_AGENT_ID}`, correlationId: payload.threadId, privacy: "sensitive", payload: payload as unknown as JsonObject, }; const appended = await store.appendEvent(candidate); const actual = assertAgentMessageRoute(appended.event); if (canonicalJson(actual as unknown as JsonObject) !== canonicalJson(payload as unknown as JsonObject)) { throw new Error("Existing agent message identity has divergent content"); } return appended;}
export function parseAgentMessageResponseEvent(event: ThoughtEvent): AgentMessageResponsePayload | undefined { if (event.type !== AGENT_MESSAGE_RESPONSE_EVENT_TYPE || event.schemaVersion !== AGENT_MESSAGE_SCHEMA_VERSION || event.sourceKind !== "agent" || event.privacy !== "sensitive") return undefined; const parsed = agentMessageResponsePayloadSchema.safeParse(event.payload); return parsed.success ? parsed.data : undefined;}