Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
43 kB · 862 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863import fs from "node:fs/promises";import path from "node:path";import YAML from "yaml";import { z } from "zod";import { canonicalJson, sha256, type JsonObject } from "../core/json.js";import { createDefaultRegistry } from "../events/registry.js";import { CONCEPTUALIZATION_OUTPUT_CONTRACT_ID, CONVERSATION_COMPACTION_OUTPUT_CONTRACT_ID, createOutputContractRegistry, OBSERVATION_OUTPUT_CONTRACT_ID, OBSERVATION_OUTPUT_CONTRACT_VERSION, PUBLIC_KNOWLEDGE_PROPOSED_DIFF_OUTPUT_CONTRACT_ID, PUBLIC_KNOWLEDGE_PROPOSED_DIFF_OUTPUT_CONTRACT_VERSION,} from "./output-contracts.js";import { CONVERSATION_COMPACTION_EVENT_TYPE } from "./conversation-compaction.js";import { AGENT_MESSAGE_RESPONSE_EVENT_TYPE, AGENT_MESSAGE_SOURCE_EVENT_TYPE, CO_AGENT_ID, STREAM_AGENT_ENDPOINT_ID,} from "./agent-messages.js";import { PUBLIC_KNOWLEDGE_PROPOSED_DIFF_EVENT_TYPE } from "../public-knowledge/types.js";import { BUILTIN_PROVIDER_PROFILE_IDS, providerKindForProfile } from "./provider-profiles.js";import { AGENT_TOOL_NAMES } from "./tools.js";import { PROPOSAL_DECLARATION_NAMES } from "./proposals.js";import { loadAdapterCatalog, privateCheckpointFor, type LoadedAdapterCatalog } from "../adapters/model-adapters.js";import type { ThoughtAgentDeclaration } from "./types.js";
const privateDeclarationCheckpoints = new WeakMap<ThoughtAgentDeclaration, string>();
const budgetCounterSchema = z.number().int().positive().max(1_000_000_000);const budgetCostSchema = z.number().int().positive().max(1_000_000_000_000);const budgetLimitFields = { maxCalls: z.number().int().positive().max(1_000_000), maxInputTokens: budgetCounterSchema, maxOutputTokens: budgetCounterSchema, maxCostMicrousd: budgetCostSchema.optional(),};const budgetLimitSchema = z.discriminatedUnion("window", [ z.object({ window: z.literal("rolling"), durationMs: z.number().int().min(1_000).max(30 * 86_400_000), ...budgetLimitFields, }).strict(), z.object({ window: z.literal("hour"), ...budgetLimitFields }).strict(), z.object({ window: z.literal("day"), ...budgetLimitFields }).strict(),]);const inferenceAccountingSchema = z.object({ leaseMs: z.number().int().min(10_000).max(3_600_000), reservation: z.object({ inputTokens: budgetCounterSchema, outputTokens: budgetCounterSchema, costMicrousd: budgetCostSchema.optional(), }).strict(), limits: z.array(budgetLimitSchema).min(1).max(8), onExhaustion: z.enum(["advance", "defer"]).default("advance"),}).strict().superRefine((value, context) => { const keys = new Set<string>(); for (let index = 0; index < value.limits.length; index += 1) { const limit = value.limits[index]!; const key = limit.window === "rolling" ? `rolling:${limit.durationMs}` : limit.window; if (keys.has(key)) context.addIssue({ code: "custom", path: ["limits", index], message: `Duplicate accounting window ${key}` }); keys.add(key); if (limit.maxCostMicrousd !== undefined && value.reservation.costMicrousd === undefined) { context.addIssue({ code: "custom", path: ["limits", index, "maxCostMicrousd"], message: "Window cost limit requires a cost reservation", }); } else if (limit.maxCostMicrousd !== undefined && value.reservation.costMicrousd !== undefined && limit.maxCostMicrousd < value.reservation.costMicrousd) { context.addIssue({ code: "custom", path: ["limits", index, "maxCostMicrousd"], message: "Window cost limit must admit one reservation" }); } if (limit.maxInputTokens < value.reservation.inputTokens) { context.addIssue({ code: "custom", path: ["limits", index, "maxInputTokens"], message: "Window input-token limit must admit one reservation" }); } if (limit.maxOutputTokens < value.reservation.outputTokens) { context.addIssue({ code: "custom", path: ["limits", index, "maxOutputTokens"], message: "Window output-token limit must admit one reservation" }); } }});
const runnerLimits = { maxOutputTokens: z.number().int().positive().max(32_000).default(2_000), timeoutMs: z.number().int().positive().max(600_000).default(60_000),};
const conversationCompactionPolicySchema = z.discriminatedUnion("mode", [ z.object({ mode: z.literal("consume"), agentId: z.string().regex(/^[a-z0-9]+(?:-[a-z0-9]+)*$/), }).strict(), z.object({ mode: z.literal("produce"), targetAgentId: z.string().regex(/^[a-z0-9]+(?:-[a-z0-9]+)*$/), triggerInputChars: z.number().int().min(16).max(1_000_000), retainInputChars: z.number().int().positive().max(999_999), }).strict().superRefine((value, context) => { if (value.retainInputChars >= value.triggerInputChars) { context.addIssue({ code: "custom", path: ["retainInputChars"], message: "Compaction must retain fewer input characters than its trigger" }); } }),]);
const declarationIdSchema = z.string().min(1).regex(/^[a-z0-9][a-z0-9-]*$/);const contextDocumentPathSchema = z.string().min(1).max(1_000).superRefine((value, context) => { if (path.posix.isAbsolute(value) || value.includes("\\") || value.split("/").some((part) => part === ".." || part === "." || part === "")) { context.addIssue({ code: "custom", message: "Subscribed document paths must be normalized relative POSIX paths" }); }});
const contextDocumentsSchema = z.object({ maxChars: z.number().int().min(1_024).max(500_000), subscriptions: z.array(z.object({ source: z.string().min(1).max(200), paths: z.array(contextDocumentPathSchema).min(1).max(100), required: z.boolean().default(true), }).strict()).min(1).max(32),}).strict().superRefine((value, context) => { const seen = new Set<string>(); for (let subscriptionIndex = 0; subscriptionIndex < value.subscriptions.length; subscriptionIndex += 1) { const subscription = value.subscriptions[subscriptionIndex]!; for (let pathIndex = 0; pathIndex < subscription.paths.length; pathIndex += 1) { const key = `${subscription.source}\u0000${subscription.paths[pathIndex]}`; if (seen.has(key)) { context.addIssue({ code: "custom", path: ["subscriptions", subscriptionIndex, "paths", pathIndex], message: "Duplicate subscribed document", }); } seen.add(key); } }});
const deterministicRunnerSchema = z.object({ kind: z.literal("deterministic"), model: z.string().min(1).max(500).optional(), ...runnerLimits,}).strict();
const piRunnerSchema = z.object({ kind: z.literal("pi"), profile: z.enum(BUILTIN_PROVIDER_PROFILE_IDS).optional(), tier: z.enum(["triage-small", "reasoning-small", "escalation"]).optional(), model: z.string().min(1).max(500).optional(), thinkingLevel: z.enum(["off", "minimal", "low", "medium", "high", "xhigh", "max"]).optional(), adapter: z.object({ id: z.string().min(1), version: z.number().int().positive() }).strict().optional(), outputMode: z.enum(["strict-json", "conversation-text", "compaction-text"]).default("strict-json"), ...runnerLimits,}).strict();
const lettaAgentRunnerSchema = z.object({ kind: z.literal("letta-agent-sdk"), backend: z.enum(["cloud", "local"]), agentIdEnv: z.string().regex(/^[A-Z_][A-Z0-9_]*$/, "agentIdEnv must be an environment variable name"), conversation: z.enum(["main", "per-document"]).default("main"), responseMode: z.enum(["strict-json", "conversation-text"]).default("strict-json"), batchResponseMode: z.literal("strict-json").optional(), outputOnly: z.boolean().default(false), proposalTool: z.literal("public-knowledge-diff").optional(), permissionMode: z.enum(["standard", "acceptEdits", "unrestricted", "strict"]).default("unrestricted"), skillSources: z.array(z.enum(["bundled", "global", "agent", "project"])).max(4).optional(), memoryDirEnv: z.string().regex(/^[A-Z_][A-Z0-9_]*$/, "memoryDirEnv must be an environment variable name").optional(), reasoningEffort: z.enum(["none", "minimal", "low", "medium", "high", "xhigh"]).optional(), dreaming: z.object({ trigger: z.enum(["off", "step-count", "compaction-event"]).default("off"), stepCount: z.number().int().positive().max(10_000).optional(), }).strict().default({ trigger: "off" }), sandbox: z.object({ ttlMinutes: z.number().int().min(1).max(60).default(5), terminateOnClose: z.boolean().default(false), }).strict().default({ ttlMinutes: 5, terminateOnClose: false }), model: z.string().min(1).max(500).optional(), ...runnerLimits,}).strict().superRefine((value, context) => { if (value.dreaming.trigger === "step-count" && value.dreaming.stepCount === undefined) { context.addIssue({ code: "custom", path: ["dreaming", "stepCount"], message: "Step-count dreaming requires stepCount" }); } if (value.dreaming.trigger !== "step-count" && value.dreaming.stepCount !== undefined) { context.addIssue({ code: "custom", path: ["dreaming", "stepCount"], message: "stepCount is valid only for step-count dreaming" }); } if (value.backend === "local") { if (!value.memoryDirEnv) context.addIssue({ code: "custom", path: ["memoryDirEnv"], message: "Local Agent SDK sessions require a memory-root environment reference" }); if (!value.outputOnly) context.addIssue({ code: "custom", path: ["outputOnly"], message: "Local Agent SDK sessions must be explicitly output-only" }); if (value.permissionMode !== "strict") context.addIssue({ code: "custom", path: ["permissionMode"], message: "Local Agent SDK sessions require strict permission mode" }); if ((value.skillSources ?? []).length > 0) context.addIssue({ code: "custom", path: ["skillSources"], message: "Local Agent SDK sessions must disable skills" }); if (value.conversation !== "per-document") context.addIssue({ code: "custom", path: ["conversation"], message: "The local-memory Agent SDK profile is limited to per-document conversations" }); if (value.responseMode !== "strict-json") context.addIssue({ code: "custom", path: ["responseMode"], message: "Per-document local Agent SDK sessions require strict JSON" }); if (value.dreaming.trigger !== "off") context.addIssue({ code: "custom", path: ["dreaming"], message: "Per-document local Agent SDK sessions disable dreaming" }); } else if (value.memoryDirEnv) { context.addIssue({ code: "custom", path: ["memoryDirEnv"], message: "Cloud Agent SDK sessions cannot select a host memory root" }); }});
const declarationFileSchema = z.object({ id: declarationIdSchema, version: z.number().int().positive(), name: z.string().min(1).optional(), description: z.string().min(1), role: z.enum(["standard", "repair", "compactor"]).default("standard"), outputContract: z.object({ id: z.string().min(1), version: z.number().int().positive(), }).strict().default({ id: OBSERVATION_OUTPUT_CONTRACT_ID, version: OBSERVATION_OUTPUT_CONTRACT_VERSION, }), subscribe: z.object({ types: z.array(z.string().min(1)).min(1), sources: z.array(z.string().min(1)).min(1).default(["*"]), privacy: z.array(z.enum(["public-source", "private", "sensitive"])).min(1), replay: z.enum(["beginning", "now"]).default("beginning"), }).strict(), privacyFloor: z.enum(["public-source", "private", "sensitive"]).default("public-source"), context: z.object({ maxEvents: z.number().int().positive().max(1_000).default(1), maxChars: z.number().int().positive().max(1_000_000).default(64_000), strategy: z.enum([ "single-event", "telegram-conversation", "telegram-compaction", "agent-conversation", "atproto-batch", "activity-batch", "coil-public-knowledge", ]).default("single-event"), documents: contextDocumentsSchema.optional(), historyAgentIds: z.array(declarationIdSchema).min(1).max(32).optional(), assistantHistoryMaxTurns: z.number().int().positive().max(100).optional(), compaction: conversationCompactionPolicySchema.optional(), payloadFields: z.array(z.string().min(1).max(100)).min(1).max(100).optional(), atprotoObject: z.boolean().default(false), requireCompletedAgentOutputs: z.boolean().default(false), }).strict(), runner: z.discriminatedUnion("kind", [ deterministicRunnerSchema, piRunnerSchema, lettaAgentRunnerSchema, ]), accounting: inferenceAccountingSchema.optional(), retry: z.object({ initialDelayMs: z.number().int().min(100).max(3_600_000), maxDelayMs: z.number().int().min(100).max(86_400_000), }).strict().superRefine((value, context) => { if (value.maxDelayMs < value.initialDelayMs) { context.addIssue({ code: "custom", path: ["maxDelayMs"], message: "Retry maximum delay must be at least the initial delay" }); } }).optional(), prompt: z.string().min(1), emit: z.array(z.string().min(1)).length(1), policy: z.object({ tools: z.array(z.enum(AGENT_TOOL_NAMES)).max(8).default([]), proposals: z.array(z.enum(PROPOSAL_DECLARATION_NAMES)).max(3).default([]).superRefine((values, context) => { if (new Set(values).size !== values.length) context.addIssue({ code: "custom", message: "Proposal names must be unique" }); }), externalActions: z.literal(false), }).strict(), enabled: z.boolean().default(true), enabledEnv: z.string().regex(/^[A-Z_][A-Z0-9_]*$/, "enabledEnv must be an environment variable name").optional(),}).strict().superRefine((value, context) => { const conceptualizationOutput = value.outputContract.id === CONCEPTUALIZATION_OUTPUT_CONTRACT_ID; const conceptGraphEvent = value.emit[0] === "stream.thought.derived.concept.graph"; const conversationText = (value.runner.kind === "pi" && value.runner.outputMode === "conversation-text") || (value.runner.kind === "letta-agent-sdk" && value.runner.responseMode === "conversation-text"); if (value.role === "standard" && conceptualizationOutput !== conceptGraphEvent) { context.addIssue({ code: "custom", path: conceptualizationOutput ? ["emit"] : ["outputContract"], message: "Standard conceptualization output and stream.thought.derived.concept.graph must be selected together", }); } if (conceptualizationOutput && conversationText) { context.addIssue({ code: "custom", path: ["runner"], message: "Conceptualization output requires strict JSON mode", }); } if (value.runner.kind === "pi" && value.runner.profile === "openai-json-default" && !conceptualizationOutput) { context.addIssue({ code: "custom", path: ["runner", "profile"], message: "The strict OpenAI JSON profile is currently bound only to the conceptualization output contract", }); } if (conversationText && value.outputContract.id !== OBSERVATION_OUTPUT_CONTRACT_ID) { context.addIssue({ code: "custom", path: ["outputContract"], message: "Conversation-text mode requires the observation output contract", }); } if (value.runner.kind === "pi" && !value.runner.profile) { context.addIssue({ code: "custom", path: ["runner", "profile"], message: "Pi agents require a trusted provider profile" }); } if (value.runner.kind === "pi" && !value.runner.model && !value.runner.tier && !value.runner.adapter) { context.addIssue({ code: "custom", path: ["runner", "model"], message: "Pi agents require a model or capability tier" }); } if (value.runner.kind !== "pi" && value.policy.tools.length > 0) { context.addIssue({ code: "custom", path: ["policy", "tools"], message: "Only Pi agents may declare tools" }); } if (value.runner.kind !== "pi" && value.policy.proposals.length > 0) { context.addIssue({ code: "custom", path: ["policy", "proposals"], message: "Only Pi agents may declare proposals" }); } if ((value.runner.kind === "pi" || value.runner.kind === "letta-agent-sdk") && !value.accounting) { context.addIssue({ code: "custom", path: ["accounting"], message: "Model-backed agents require durable inference accounting" }); } if (value.runner.kind === "deterministic" && value.accounting) { context.addIssue({ code: "custom", path: ["accounting"], message: "Deterministic agents cannot reserve provider inference" }); } if (value.accounting) { if (value.accounting.leaseMs < value.runner.timeoutMs + 5_000) { context.addIssue({ code: "custom", path: ["accounting", "leaseMs"], message: "Accounting lease must exceed the runner timeout by at least five seconds" }); } if (value.accounting.reservation.outputTokens < value.runner.maxOutputTokens) { context.addIssue({ code: "custom", path: ["accounting", "reservation", "outputTokens"], message: "Output reservation must cover runner maxOutputTokens" }); } if (value.accounting.reservation.inputTokens < Math.ceil(value.context.maxChars / 2)) { context.addIssue({ code: "custom", path: ["accounting", "reservation", "inputTokens"], message: "Input reservation must conservatively cover at least half of maxChars" }); } } if (value.role === "repair" && value.runner.kind !== "pi") { context.addIssue({ code: "custom", path: ["role"], message: "Repair agents must use the Pi runner" }); } if (value.role === "repair" && value.policy.tools.length > 0) { context.addIssue({ code: "custom", path: ["policy", "tools"], message: "Repair agents cannot declare tools" }); } if (value.role === "repair" && value.policy.proposals.length > 0) { context.addIssue({ code: "custom", path: ["policy", "proposals"], message: "Repair agents cannot declare proposal tools" }); } if (value.role === "repair" && ( value.subscribe.types.length !== 1 || value.subscribe.types[0] !== "stream.thought.agent.repair.requested" || value.emit[0] !== "stream.thought.derived.output.correction.proposed" )) { context.addIssue({ code: "custom", path: ["role"], message: "Repair agents must subscribe to repair requests and emit correction proposals only" }); } if (value.role !== "repair" && value.subscribe.types.includes("stream.thought.agent.repair.requested")) { context.addIssue({ code: "custom", path: ["role"], message: "Repair requests require an explicit repair role" }); } if (["telegram-conversation", "telegram-compaction"].includes(value.context.strategy) && ( value.subscribe.types.length !== 1 || value.subscribe.types[0] !== "stream.thought.source.telegram.message" || value.subscribe.privacy.length !== 1 || value.subscribe.privacy[0] !== "sensitive" )) { context.addIssue({ code: "custom", path: ["context", "strategy"], message: "Telegram conversation and compaction context require one sensitive Telegram message subscription", }); } if (value.context.documents && !["telegram-conversation", "telegram-compaction", "agent-conversation"].includes(value.context.strategy)) { context.addIssue({ code: "custom", path: ["context", "documents"], message: "Subscribed document context is currently supported only for Telegram conversations, compactors, and private agent conversations", }); } if (value.context.documents && value.runner.kind !== "pi") { context.addIssue({ code: "custom", path: ["context", "documents"], message: "Subscribed trusted documents currently require the Pi runner", }); } if (value.context.documents && value.context.documents.maxChars > value.context.maxChars - 1_024) { context.addIssue({ code: "custom", path: ["context", "documents", "maxChars"], message: "Subscribed documents must leave at least 1024 characters for source context", }); } if (value.context.historyAgentIds && !["telegram-conversation", "telegram-compaction"].includes(value.context.strategy)) { context.addIssue({ code: "custom", path: ["context", "historyAgentIds"], message: "Conversation history agent ids are valid only for Telegram conversation roles", }); } if (value.context.assistantHistoryMaxTurns && value.context.strategy !== "telegram-conversation") { context.addIssue({ code: "custom", path: ["context", "assistantHistoryMaxTurns"], message: "Assistant conversation history limits are valid only for Telegram conversation context", }); } if (value.context.compaction?.mode === "consume" && ( value.role !== "standard" || value.context.strategy !== "telegram-conversation" || value.context.compaction.agentId === value.id )) { context.addIssue({ code: "custom", path: ["context", "compaction"], message: "Compaction consumption requires a standard Telegram agent and a separate compactor id", }); } if (value.context.compaction?.mode === "produce" && ( value.role !== "compactor" || value.context.strategy !== "telegram-compaction" || value.context.compaction.targetAgentId === value.id )) { context.addIssue({ code: "custom", path: ["context", "compaction"], message: "Compaction production requires a separate compactor role targeting one Telegram agent", }); } if (value.role === "compactor" && ( value.runner.kind !== "pi" || value.runner.outputMode !== "compaction-text" || value.runner.adapter !== undefined || value.outputContract.id !== CONVERSATION_COMPACTION_OUTPUT_CONTRACT_ID || value.emit[0] !== CONVERSATION_COMPACTION_EVENT_TYPE || value.policy.tools.length > 0 || value.policy.proposals.length > 0 || value.context.compaction?.mode !== "produce" )) { context.addIssue({ code: "custom", path: ["role"], message: "Compactors require tool-free plain-text Pi summarization, Telegram compaction context, and the compaction event/contract", }); } if (value.role !== "compactor" && value.context.strategy === "telegram-compaction") { context.addIssue({ code: "custom", path: ["role"], message: "Telegram compaction context requires the compactor role" }); } if (value.role !== "compactor" && value.emit[0] === CONVERSATION_COMPACTION_EVENT_TYPE) { context.addIssue({ code: "custom", path: ["emit"], message: "Conversation compaction events require the compactor role" }); } const lettaAtprotoObjectContext = value.runner.kind === "letta-agent-sdk" && ["single-event", "atproto-batch"].includes(value.context.strategy) && value.context.maxEvents === 1 && value.subscribe.types.some((type) => ( type === "stream.thought.source.atproto.commit" || type === "stream.thought.derived.event.batch" || type === "*" )); const conceptualizerAtprotoBatchContext = value.runner.kind === "pi" && conceptualizationOutput && conceptGraphEvent && value.context.strategy === "atproto-batch" && value.context.maxEvents === 1 && value.subscribe.types.length === 1 && value.subscribe.types[0] === "stream.thought.derived.event.batch" && value.policy.tools.length === 0 && value.policy.proposals.length === 0 && value.policy.externalActions === false; if (value.context.atprotoObject && ( value.context.maxChars < 2_048 || (!lettaAtprotoObjectContext && !conceptualizerAtprotoBatchContext) )) { context.addIssue({ code: "custom", path: ["context", "atprotoObject"], message: "ATProto object context requires either a bounded Letta Agent SDK source/batch declaration or a tool-free Pi conceptualizer bound to one batch event", }); } if (value.context.requireCompletedAgentOutputs && ( value.context.strategy !== "activity-batch" || value.context.maxEvents !== 1 || value.subscribe.types.length !== 1 || value.subscribe.types[0] !== "stream.thought.derived.event.batch" || value.subscribe.sources.some((source) => !source.startsWith("batch:")) )) { context.addIssue({ code: "custom", path: ["context", "requireCompletedAgentOutputs"], message: "Completed agent-output context requires one derived-observation batch trigger", }); } if (value.runner.kind === "letta-agent-sdk") { if (!["single-event", "atproto-batch", "activity-batch", "coil-public-knowledge"].includes(value.context.strategy) || value.context.maxEvents !== 1) { context.addIssue({ code: "custom", path: ["context"], message: "Letta Agent SDK declarations require one trigger event context; the Letta conversation owns history", }); } if (value.subscribe.sources.length > 8 || value.subscribe.sources.some((source) => source.includes("*"))) { context.addIssue({ code: "custom", path: ["subscribe", "sources"], message: "Letta Agent SDK declarations require one to eight concrete source namespaces", }); } const publicKnowledge = value.context.strategy === "coil-public-knowledge"; const fileTypes = new Set([ "stream.thought.source.file.added", "stream.thought.source.file.changed", "stream.thought.source.file.renamed", "stream.thought.source.file.deleted", ]); if (publicKnowledge && ( value.runner.backend !== "local" || value.runner.conversation !== "per-document" || value.runner.responseMode !== "strict-json" || !value.runner.outputOnly || value.runner.proposalTool !== "public-knowledge-diff" || value.subscribe.replay !== "now" || value.outputContract.id !== PUBLIC_KNOWLEDGE_PROPOSED_DIFF_OUTPUT_CONTRACT_ID || value.outputContract.version !== PUBLIC_KNOWLEDGE_PROPOSED_DIFF_OUTPUT_CONTRACT_VERSION || value.emit[0] !== PUBLIC_KNOWLEDGE_PROPOSED_DIFF_EVENT_TYPE || value.subscribe.sources.length !== 1 || value.subscribe.privacy.length !== 1 || value.subscribe.privacy[0] !== "sensitive" || value.subscribe.types.length < 1 || value.subscribe.types.some((type) => !fileTypes.has(type)) )) { context.addIssue({ code: "custom", path: ["context", "strategy"], message: "Coil Public Knowledge context requires replay now, one sensitive filesystem source, the fixed proposal tool, output-only per-document strict JSON, and the canonical proposed-diff contract/event", }); } if (!publicKnowledge && value.runner.conversation === "per-document") { context.addIssue({ code: "custom", path: ["runner", "conversation"], message: "Per-document conversations require Coil Public Knowledge context" }); } if (!publicKnowledge && value.runner.proposalTool) { context.addIssue({ code: "custom", path: ["runner", "proposalTool"], message: "The Public Knowledge proposal tool requires Coil Public Knowledge context" }); } } const telegramConversationText = value.context.strategy === "telegram-conversation" && value.emit[0] === "stream.thought.derived.message.observation"; const agentConversationText = value.context.strategy === "agent-conversation" && value.emit[0] === AGENT_MESSAGE_RESPONSE_EVENT_TYPE; if (value.runner.kind === "pi" && value.runner.outputMode === "conversation-text" && ( value.role !== "standard" || (!telegramConversationText && !agentConversationText) || value.policy.tools.length > 0 )) { context.addIssue({ code: "custom", path: ["runner", "outputMode"], message: "Conversation-text output requires a standard Pi Telegram or private agent conversation declaration with no read-only model tools and the matching message output", }); } if (value.policy.proposals.length > 0 && ( value.runner.kind !== "pi" || value.runner.outputMode !== "conversation-text" || value.role !== "standard" || value.context.strategy !== "telegram-conversation" || value.emit[0] !== "stream.thought.derived.message.observation" || value.subscribe.privacy.length !== 1 || value.subscribe.privacy[0] !== "sensitive" || !value.context.documents )) { context.addIssue({ code: "custom", path: ["policy", "proposals"], message: "Proposal tools require a sensitive standard Pi Telegram conversation with exact subscribed documents", }); } if (value.context.strategy === "agent-conversation" && ( value.id !== STREAM_AGENT_ENDPOINT_ID || value.role !== "standard" || value.outputContract.id !== OBSERVATION_OUTPUT_CONTRACT_ID || value.runner.kind !== "pi" || value.runner.outputMode !== "conversation-text" || value.subscribe.types.length !== 1 || value.subscribe.types[0] !== AGENT_MESSAGE_SOURCE_EVENT_TYPE || value.subscribe.sources.length !== 1 || value.subscribe.sources[0] !== `agent-message:${CO_AGENT_ID}` || value.subscribe.privacy.length !== 1 || value.subscribe.privacy[0] !== "sensitive" || value.subscribe.replay !== "now" || value.emit[0] !== AGENT_MESSAGE_RESPONSE_EVENT_TYPE || value.policy.tools.length > 0 || value.policy.proposals.length > 0 || !value.context.documents || value.context.compaction !== undefined || value.context.historyAgentIds !== undefined || value.context.assistantHistoryMaxTurns !== undefined )) { context.addIssue({ code: "custom", path: ["context", "strategy"], message: "Agent conversation context is fixed to the sensitive replay-now Co-to-Stream route with operator documents, observation text output, and no tools, proposals, compaction, or Telegram history", }); } if (value.policy.proposals.includes("memory-change") && !value.context.documents?.subscriptions.some((subscription) => ( subscription.source === "filesystem:telegram-agent-context" && subscription.required && subscription.paths.includes("memory.md") ))) { context.addIssue({ code: "custom", path: ["policy", "proposals"], message: "Memory proposals require the exact required Stream memory.md subscription", }); }});
export async function loadAgentDeclarations( directory: string, environment: NodeJS.ProcessEnv = process.env, catalog?: LoadedAdapterCatalog,): Promise<ThoughtAgentDeclaration[]> { const registry = createDefaultRegistry(); const outputContracts = createOutputContractRegistry(); const entries = await fs.readdir(directory, { withFileTypes: true }).catch((error: NodeJS.ErrnoException) => { if (error.code === "ENOENT") return []; throw error; }); const projectRoot = path.resolve(directory, ".."); const adapterCatalog = catalog ?? await loadAdapterCatalogIfPresent(projectRoot, environment); const declarations: ThoughtAgentDeclaration[] = []; for (const entry of entries.sort((left, right) => left.name.localeCompare(right.name))) { if (!entry.isFile() || !/\.(ya?ml|json)$/i.test(entry.name)) continue; const filePath = path.join(directory, entry.name); const raw = await fs.readFile(filePath, "utf8"); const parsed = entry.name.endsWith(".json") ? JSON.parse(raw) : YAML.parse(raw); const file = declarationFileSchema.parse(parsed); const effectiveEnabled = file.enabled || enabledByEnvironment(file.enabledEnv, environment); for (const type of file.emit) { if (!registry.has(type, 1)) throw new Error(`Agent ${file.id} emits an unregistered event type: ${type}@1`); } const promptPath = path.resolve(projectRoot, file.prompt); if (!isWithin(projectRoot, promptPath)) throw new Error(`Agent prompt escapes project root: ${file.prompt}`); const systemPrompt = await fs.readFile(promptPath, "utf8"); const compiledEventTypes = registry.list() .map((definition) => definition.type) .filter((type) => file.subscribe.types.some((pattern) => matchesPattern(pattern, type))); if (compiledEventTypes.length === 0) { throw new Error(`Agent ${file.id} subscription compiles to no registered event types`); } const selectedAdapter = file.runner.kind === "pi" && file.runner.adapter ? adapterCatalog?.selectedByDeclaration.get(`${file.id}@${file.version}`) : undefined; const selectedRelease = selectedAdapter ? adapterCatalog?.releases.find((release) => ( release.id === selectedAdapter.id && release.version === selectedAdapter.version && release.manifestSha256 === selectedAdapter.manifestSha256 )) : undefined; if (file.runner.kind === "pi" && file.runner.adapter && (!adapterCatalog || !selectedAdapter || selectedAdapter.id !== file.runner.adapter.id || selectedAdapter.version !== file.runner.adapter.version)) { throw new Error(`Agent ${file.id}@${file.version} adapter selection does not match the loaded deployment catalog`); } if (selectedAdapter && !selectedRelease) { throw new Error(`Agent ${file.id}@${file.version} selected release is absent from the loaded release catalog`); } if (file.runner.kind === "pi" && file.runner.adapter && (file.runner.model || file.runner.tier)) { throw new Error(`Agent ${file.id}@${file.version} cannot combine adapter selection with model or tier`); } const concreteModel = selectedAdapter?.baseModel ?? resolveRunnerModel(file.runner, environment, effectiveEnabled); const lettaAgentId = file.runner.kind === "letta-agent-sdk" ? resolveLettaAgentId(file.runner.agentIdEnv, environment, effectiveEnabled) : undefined; const outputContract = outputContracts.get(file.outputContract.id, file.outputContract.version).identity; const declaration: ThoughtAgentDeclaration = { id: file.id, version: file.version, name: file.name ?? file.id, description: file.description, mode: file.runner.kind, role: file.role, outputContract: { ...outputContract }, ...(file.runner.kind === "pi" && (selectedAdapter?.providerProfile ?? file.runner.profile) ? { provider: providerKindForProfile((selectedAdapter?.providerProfile ?? file.runner.profile)!), providerProfile: (selectedAdapter?.providerProfile ?? file.runner.profile)!, } : {}), ...(selectedAdapter ? { modelAdapterRelease: selectedRelease, modelAdapter: selectedAdapter, adapterCatalogDigest: adapterCatalog!.digest, adapterCatalogGeneration: adapterCatalog!.generation, } : {}), ...(file.runner.kind === "letta-agent-sdk" ? { provider: file.runner.backend === "cloud" ? "letta-cloud" as const : "letta-local" as const, lettaAgent: { backend: file.runner.backend, agentIdEnv: file.runner.agentIdEnv, ...(lettaAgentId ? { agentId: lettaAgentId } : {}), conversation: file.runner.conversation, responseMode: file.runner.responseMode, ...(file.runner.batchResponseMode ? { batchResponseMode: file.runner.batchResponseMode } : {}), outputOnly: file.runner.outputOnly, ...(file.runner.proposalTool ? { proposalTool: file.runner.proposalTool } : {}), permissionMode: file.runner.permissionMode, ...(file.runner.skillSources ? { skillSources: file.runner.skillSources } : {}), ...(file.runner.memoryDirEnv ? { memoryDirEnv: file.runner.memoryDirEnv } : {}), ...(file.runner.reasoningEffort ? { reasoningEffort: file.runner.reasoningEffort } : {}), dreaming: file.runner.dreaming, sandbox: file.runner.sandbox, }, } : {}), ...(file.runner.kind === "pi" && file.runner.tier ? { modelTier: file.runner.tier } : {}), ...(concreteModel ? { model: concreteModel } : {}), ...(file.runner.kind === "pi" ? { outputMode: file.runner.outputMode, ...(file.runner.thinkingLevel ? { thinkingLevel: file.runner.thinkingLevel } : {}), } : {}), eventTypes: file.subscribe.types, compiledEventTypes, sourcePatterns: file.subscribe.sources, acceptedPrivacy: file.subscribe.privacy, ...(file.privacyFloor !== "public-source" ? { privacyFloor: file.privacyFloor } : {}), initialReplay: file.subscribe.replay, outputEventType: file.emit[0]!, emit: file.emit, promptRef: file.prompt, systemPrompt, enabled: effectiveEnabled, maxEvents: file.context.maxEvents, maxInputChars: file.context.maxChars, contextStrategy: file.context.strategy, ...(file.context.documents ? { contextDocumentMaxChars: file.context.documents.maxChars, contextDocumentSubscriptions: file.context.documents.subscriptions, } : {}), ...(file.context.historyAgentIds ? { conversationHistoryAgentIds: file.context.historyAgentIds } : {}), ...(file.context.assistantHistoryMaxTurns ? { conversationAssistantHistoryMaxTurns: file.context.assistantHistoryMaxTurns, } : {}), ...(file.context.compaction ? { conversationCompaction: file.context.compaction } : {}), ...(file.context.payloadFields ? { payloadFields: file.context.payloadFields } : {}), ...(file.context.atprotoObject ? { atprotoObjectContext: true } : {}), ...(file.context.requireCompletedAgentOutputs ? { requireCompletedAgentOutputs: true } : {}), maxOutputTokens: file.runner.maxOutputTokens, timeoutMs: file.runner.timeoutMs, ...(file.accounting ? { accounting: file.accounting } : {}), ...(file.retry ? { retry: file.retry } : {}), tools: file.policy.tools, proposals: file.policy.proposals, externalActions: file.policy.externalActions, }; declaration.declarationFingerprint = declarationFingerprint(declaration); const compiledDeclaration = deepFreeze(declaration); if (selectedAdapter) privateDeclarationCheckpoints.set(compiledDeclaration, privateCheckpointFor(selectedAdapter)); declarations.push(compiledDeclaration); } const duplicate = findDuplicate(declarations.map((declaration) => declaration.id)); if (duplicate) throw new Error(`Duplicate agent declaration id: ${duplicate}`); const duplicateMainConversation = findDuplicate(declarations .filter((declaration) => declaration.enabled && declaration.lettaAgent?.conversation === "main" && declaration.lettaAgent.agentId) .map((declaration) => `${declaration.lettaAgent!.backend}:${declaration.lettaAgent!.agentId!}`)); if (duplicateMainConversation) { throw new Error("Enabled main-conversation Letta declarations must use distinct agent identities"); } return declarations;}
function resolveRunnerModel(runner: { kind: "deterministic" | "pi" | "letta-agent-sdk"; profile?: "tinker-default" | "openai-json-default" | "openai-compatible-default" | undefined; tier?: "triage-small" | "reasoning-small" | "escalation" | undefined; model?: string | undefined; adapter?: { id: string; version: number } | undefined;}, environment: NodeJS.ProcessEnv, enabled: boolean): string | undefined { if (runner.model) return runner.model; if (!runner.tier) return undefined; const provider = runner.profile ? providerKindForProfile(runner.profile) : "openai-compatible"; const environmentKey = `THOUGHTSTREAM_${provider === "tinker" ? "TINKER" : "MODEL"}_${runner.tier.toUpperCase().replaceAll("-", "_")}_MODEL`; const configured = environment[environmentKey]; if (configured) return configured; if (provider === "tinker" && runner.tier === "triage-small") return "Qwen/Qwen3.5-4B"; if (provider === "tinker" && runner.tier === "reasoning-small") return "thinkingmachines/Inkling-Small"; if (!enabled) return undefined; throw new Error(`No concrete model mapping for ${provider}/${runner.tier}; set ${environmentKey}`);}
function resolveLettaAgentId( environmentKey: string, environment: NodeJS.ProcessEnv, enabled: boolean,): string | undefined { const configured = environment[environmentKey]?.trim(); if (configured) { if (!/^agent-[A-Za-z0-9_-]+$/.test(configured)) { throw new Error(`${environmentKey} does not contain a valid Letta agent id`); } return configured; } if (!enabled) return undefined; throw new Error(`No Letta agent configured; set ${environmentKey}`);}
function enabledByEnvironment(environmentKey: string | undefined, environment: NodeJS.ProcessEnv): boolean { if (!environmentKey) return false; const value = environment[environmentKey]?.trim().toLowerCase(); if (!value || value === "0" || value === "false") return false; if (value === "1" || value === "true") return true; throw new Error(`${environmentKey} must be one of 1, true, 0, or false`);}
export function declarationFingerprint(declaration: ThoughtAgentDeclaration): string { const { declarationFingerprint: _ignored, ...body } = declaration; return sha256(canonicalJson(JSON.parse(JSON.stringify(body)) as JsonObject));}
export function assertTrustedAdapterDeclaration(declaration: ThoughtAgentDeclaration): void { if (declaration.modelAdapter && !privateDeclarationCheckpoints.has(declaration)) { throw new Error(`Agent ${declaration.id}@${declaration.version} adapter declaration was not compiled by the trusted declaration loader`); }}
export function privateCheckpointForDeclaration(declaration: ThoughtAgentDeclaration): string { assertTrustedAdapterDeclaration(declaration); const checkpoint = privateDeclarationCheckpoints.get(declaration); if (!checkpoint || !declaration.modelAdapter || sha256(checkpoint) !== declaration.modelAdapter.binding.checkpointReferenceSha256) { throw new Error(`Agent ${declaration.id}@${declaration.version} has no valid private adapter binding`); } return checkpoint;}
export function agentRole(declaration: ThoughtAgentDeclaration): "standard" | "repair" | "compactor" { return declaration.role ?? "standard";}
export function matchesDeclaration(declaration: ThoughtAgentDeclaration, event: { type: string; privacy: string }): boolean { if (!declaration.enabled || !declaration.acceptedPrivacy.includes(event.privacy as never)) return false; const source = "source" in event && typeof event.source === "string" ? event.source : ""; return declaration.eventTypes.some((pattern) => matchesPattern(pattern, event.type)) && declaration.sourcePatterns.some((pattern) => matchesPattern(pattern, source));}
function isWithin(root: string, candidate: string): boolean { const relative = path.relative(root, candidate); return relative === "" || (!relative.startsWith("..") && !path.isAbsolute(relative));}
function matchesPattern(pattern: string, value: string): boolean { if (pattern === "*") return true; if (pattern.endsWith("*")) return value.startsWith(pattern.slice(0, -1)); return pattern === value;}
function findDuplicate(values: string[]): string | undefined { const seen = new Set<string>(); for (const value of values) { if (seen.has(value)) return value; seen.add(value); } return undefined;}
function deepFreeze<T>(value: T): T { if (value && typeof value === "object") { Object.freeze(value); for (const child of Object.values(value as Record<string, unknown>)) deepFreeze(child); } return value;}
async function loadAdapterCatalogIfPresent(projectRoot: string, environment: NodeJS.ProcessEnv): Promise<LoadedAdapterCatalog | undefined> { try { await fs.access(path.join(projectRoot, "adapters", "deployment.yaml")); } catch { return undefined; } return loadAdapterCatalog(projectRoot, environment);}