Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
71 kB · 1655 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656import { canonicalJson, sha256, type JsonObject } from "../core/json.js";import { stableKey } from "../core/ids.js";import type { EventCandidate, ThoughtEvent } from "../events/types.js";import { inferenceAccountingEnabled } from "../jazz/schema.js";import { CORRECTION_PROPOSAL_EVENT_TYPE, MEMORY_PROPOSAL_EVENT_TYPE } from "../agent-proposals/contracts.js";import { FOCUS_DECLARATION_PROPOSED_EVENT_TYPE, focusDeclarationFingerprint,} from "../focuses/types.js";import type { JazzThoughtStore } from "../jazz/store.js";import { appendSchedulerExhaustedIncident, type SchedulerExhaustionEvidence } from "../incidents/projector.js";import { rebuildEffectiveOutputForRun } from "../projections/effective-output.js";import type { AgentRun, ConsumerEventQuery, ConsumerProgress, InferenceBudgetPolicy, InferenceUsage,} from "../store/types.js";import { buildAtprotoBatchContextPacket, buildActivityBatchContextPacket, buildDurableAtprotoObjectContextPacket, buildContextPacket, buildSubscribedAgentConversationContextPacket, buildRepairContextPacket, buildSubscribedTelegramConversationContextPacket, buildTelegramConversationCompactionContextPacket, buildTelegramConversationContextPacket, type AtprotoObjectContextOptions,} from "./context.js";import { AGENT_MESSAGE_RESPONSE_EVENT_TYPE, AgentMessageNotAdmitted, assertAgentMessageRoute,} from "./agent-messages.js";import { assertTrustedAdapterDeclaration, declarationFingerprint } from "./declarations.js";import { executionAdapterRevisionFor } from "./execution-adapters.js";import { declarationPrivacy } from "../security/privacy.js";import { DeterministicTriageRunner } from "./deterministic.js";import { LETTA_AGENT_SDK_ADAPTER_REVISION, LettaAgentSdkRunner } from "./letta-agent-sdk.js";import { canonicalStructuredOutput, CONCEPTUALIZATION_OUTPUT_CONTRACT_ID, CONVERSATION_COMPACTION_OUTPUT_CONTRACT_ID, conversationCompactionOutputSchema, createOutputContractRegistry, OBSERVATION_OUTPUT_CONTRACT_ID, PUBLIC_KNOWLEDGE_PROPOSED_DIFF_OUTPUT_CONTRACT_ID, PUBLIC_KNOWLEDGE_RECOMMENDATION_OUTPUT_CONTRACT_ID, REVIEW_RESPONSE_OUTPUT_CONTRACT_ID, reviewResponseSummary, outputContractForDeclaration, outputContractIdentityJson, parseOutputContractIdentity, OutputContractValidationError, type OutputContractRegistry,} from "./output-contracts.js";import { REVIEW_RESPONSE_EVENT_TYPE } from "../review/types.js";import { PUBLIC_KNOWLEDGE_PROPOSED_DIFF_EVENT_TYPE } from "../public-knowledge/types.js";import { buildPublicKnowledgeContextPacket, PublicKnowledgeEligibilitySkip, type PublicKnowledgeContextOptions,} from "../public-knowledge/context.js";import { PiAgentRunner } from "./pi.js";import { CONVERSATION_COMPACTION_EVENT_TYPE, ConversationCompactionNotNeeded, conversationCompactionBoundarySha256, conversationCompactionContextSnapshotSchema, conversationCompactionPlanSchema,} from "./conversation-compaction.js";import { RepairRequestCoordinator } from "./repairs.js";import { ConsumerScheduler } from "./scheduler.js";import { AgentRunFailure, type AgentOutput, type AgentRunner, type ThoughtAgentDeclaration } from "./types.js";import { proposalCapabilitiesSchema, resolveCorrectionTarget, validateCapturedProposalAgainstCapabilities,} from "./proposals.js";import { telegramFocusDescription } from "./telegram-help.js";
export interface ConsumerHandle { drain(): Promise<void>; stop(): Promise<void>;}
export interface ProcessEventResult { declaration: ThoughtAgentDeclaration; runId: string; output?: AgentOutput; derivedEvent?: ThoughtEvent; error?: string; retryable?: boolean;}
export interface ThoughtAgentRuntimeOptions { maxConcurrentOperations?: number; reconcileIntervalMs?: number; atprotoObjectContext?: AtprotoObjectContextOptions | undefined; publicKnowledgeContext?: PublicKnowledgeContextOptions | undefined;}
export class ThoughtAgentRuntime { private readonly runners = new Map<string, AgentRunner>(); private readonly outputContracts: OutputContractRegistry; private readonly repairCoordinator: RepairRequestCoordinator; private readonly declarationsByVersion = new Map<string, ThoughtAgentDeclaration>(); private readonly maxConcurrentOperations: number; private readonly reconcileIntervalMs: number; private readonly atprotoObjectContext: AtprotoObjectContextOptions; private readonly publicKnowledgeContext: PublicKnowledgeContextOptions | undefined;
constructor( private readonly store: JazzThoughtStore, runners?: AgentRunner[], options: ThoughtAgentRuntimeOptions = {}, ) { this.maxConcurrentOperations = options.maxConcurrentOperations ?? 4; this.reconcileIntervalMs = options.reconcileIntervalMs ?? 1_000; this.atprotoObjectContext = options.atprotoObjectContext ?? {}; this.publicKnowledgeContext = options.publicKnowledgeContext; if (!Number.isSafeInteger(this.maxConcurrentOperations) || this.maxConcurrentOperations < 1) { throw new Error("Maximum concurrent consumer operations must be a positive integer"); } if (!Number.isSafeInteger(this.reconcileIntervalMs) || this.reconcileIntervalMs < 1) { throw new Error("Consumer reconciliation interval must be a positive integer"); } this.outputContracts = createOutputContractRegistry(); this.repairCoordinator = new RepairRequestCoordinator(store, undefined, this.outputContracts); const selected = runners ?? [ new DeterministicTriageRunner(), new PiAgentRunner({ artifactRoot: store.getArtifactRoot(), outputContracts: this.outputContracts }), new LettaAgentSdkRunner({ artifactRoot: store.getArtifactRoot(), outputContracts: this.outputContracts, conversationStore: store, }), ]; for (const runner of selected) this.runners.set(runner.mode, runner); }
async registerDeclarations(declarations: ThoughtAgentDeclaration[]): Promise<void> { assertUniqueLettaAgentOwners(declarations); assertConversationCompactionPairing(declarations); for (const declaration of declarations) { assertTrustedAdapterDeclaration(declaration); assertOutputContractEventBinding(declaration); assertRuntimeAccountingPolicy(declaration); this.declarationsByVersion.set(`${declaration.id}@${declaration.version}`, declaration); const declarationJson = asJsonObject(declaration); await this.store.upsertAgent({ id: declaration.id, version: declaration.version, enabled: declaration.enabled, spec: declarationJson, specHash: sha256(canonicalJson(declarationJson)), updatedAt: new Date().toISOString(), }); } await this.repairCoordinator.reconcile(declarations); }
async consumeBacklog(declarations: ThoughtAgentDeclaration[]): Promise<ProcessEventResult[]> { await this.registerDeclarations(declarations); const runs: ProcessEventResult[] = []; for (const declaration of declarations.filter((candidate) => candidate.enabled)) { for (const source of await this.resolveSources(declaration)) { for (;;) { const progress = await this.consumerProgressForStart(declaration, source); const events = await this.store.queryConsumerEvents(this.consumerQuery(declaration, source, progress?.lastSequence ?? 0)); let retryableFailure = false; for (const event of events) { const result = await this.executeEvent(event, declaration); runs.push(result); if (result.retryable) { retryableFailure = true; break; } } if (retryableFailure || events.length < 1_000) break; } } } return runs; }
async startConsumers(declarations: ThoughtAgentDeclaration[]): Promise<ConsumerHandle> { await this.registerDeclarations(declarations); const unsubscribers: Array<() => void> = []; const subscriptions = new Set<string>(); const reconcilers = new Map<string, () => void>(); const scheduler = new ConsumerScheduler({ concurrency: this.maxConcurrentOperations, attempts: 3, retryDelayMs: (attempt) => 25 * attempt, errorMessage: "thought stream consumer cycles failed", onError: async (_error, operationKey, context) => { const incident = await appendSchedulerExhaustedIncident(this.store, { operationKey, ...schedulerIncidentContext(context), }); process.stderr.write(`thought stream consumer cycle exhausted retries; recorded incident ${incident.id}.\n`); }, }); let stopped = false; const enqueue = (key: string, operation: () => Promise<void>, context?: Omit<SchedulerExhaustionEvidence, "operationKey">): void => { scheduler.enqueue(key, operation, context); }; const enabled = declarations.filter((candidate) => candidate.enabled); const install = async (declaration: ThoughtAgentDeclaration, source: string): Promise<void> => { const subscriptionKey = `${declaration.id}@${declaration.version}:${source}`; const operationKey = consumerOperationKey(declaration, source); if (stopped || subscriptions.has(subscriptionKey)) return; subscriptions.add(subscriptionKey); try { const progress = await this.consumerProgressForStart(declaration, source); const query = this.consumerQuery(declaration, source, progress?.lastSequence ?? 0); const consumeAvailable = async (): Promise<boolean> => { const currentProgress = await this.store.getConsumerProgress(progressId(declaration, source)); const events = await this.store.queryConsumerEvents(this.consumerQuery( declaration, source, currentProgress?.lastSequence ?? 0, )); let retryableFailure = false; for (const event of events) { const current = await this.store.getConsumerProgress(progressId(declaration, source)); if ((current?.lastSequence ?? 0) >= event.sourceSequence) continue; const result = await this.executeEvent(event, declaration); if (result.retryable) { retryableFailure = true; break; } } return !retryableFailure && events.length === 1_000; }; let queued = false; let requested = false; const requestConsume = (): void => { if (stopped) return; requested = true; if (queued) return; queued = true; enqueue(operationKey, async () => { let moreAvailable = false; try { requested = false; moreAvailable = await consumeAvailable(); } finally { queued = false; // A full batch becomes a new scheduler operation so deep backlogs // yield the global slot. Interval ticks that arrive during a run // collapse into the same single follow-up request. if (!stopped && (requested || moreAvailable)) requestConsume(); } }, { source, agentId: declaration.id, agentVersion: declaration.version, }); }; reconcilers.set(subscriptionKey, requestConsume); unsubscribers.push(this.store.subscribeConsumerEvents(query, () => { requestConsume(); })); requestConsume(); } catch (error) { subscriptions.delete(subscriptionKey); reconcilers.delete(subscriptionKey); throw error; } }; for (const declaration of enabled) { for (const source of await this.resolveSources(declaration)) await install(declaration, source); } if (enabled.some((declaration) => declaration.sourcePatterns.some((pattern) => pattern.includes("*")))) { unsubscribers.push(this.store.subscribeSources((sources) => { if (stopped) return; enqueue("source-discovery", async () => { for (const declaration of enabled) { for (const source of sources) { if (declaration.sourcePatterns.some((pattern) => matchesPatternValue(pattern, source.id))) { await install(declaration, source.id); } } } }); })); } // Local Jazz persistence is query-visible across processes but does not // guarantee that a subscription callback in one process fires for another // process's write. Subscriptions remain the low-latency wakeup; this bounded // reconciliation loop makes durable progress authoritative. const reconciliationTimer = setInterval(() => { if (stopped) return; for (const reconcile of reconcilers.values()) reconcile(); }, this.reconcileIntervalMs); reconciliationTimer.unref(); return { drain: async () => scheduler.drain(), stop: async () => { stopped = true; clearInterval(reconciliationTimer); for (const unsubscribe of unsubscribers) unsubscribe(); await scheduler.drain(); reconcilers.clear(); }, }; }
private async executeEvent(event: ThoughtEvent, declaration: ThoughtAgentDeclaration): Promise<ProcessEventResult> { const progressKey = progressId(declaration, event.source); const progress = await this.store.getConsumerProgress(progressKey); if ((progress?.lastSequence ?? 0) >= event.sourceSequence) { const prior = await this.store.latestRunForExecution(executionKey(event, declaration)); if (!prior) throw new Error(`Progress exists without execution evidence for ${event.id}`); return { declaration, runId: prior.id, ...(prior.errorText ? { error: prior.errorText } : {}) }; }
const key = executionKey(event, declaration); const prior = await this.store.latestRunForExecution(key); const pendingRetryAt = prior ? retryAtForPendingRun(prior, declaration) : undefined; if (pendingRetryAt && Date.now() < Date.parse(pendingRetryAt)) { return { declaration, runId: prior!.id, error: prior!.errorText ?? "Inference retry is delayed", retryable: true, }; } if (prior?.status === "blocked" && !isDeferredExhaustion(declaration)) { await this.settleInterruptedBlockedRun(prior, declaration, event); return { declaration, runId: prior.id, error: prior.errorText ?? "Inference budget exhausted" }; } if (prior?.status === "running") { if ((declaration.role ?? "standard") === "repair") { await this.abandonInterruptedRepairRun(prior, declaration, event); return { declaration, runId: prior.id, error: "Interrupted repair run was terminally abandoned without another proposal attempt", }; } await this.abandonInterruptedRun(prior, declaration, event); } const attempt = (prior?.attempt ?? 0) + 1; const runId = stableKey("run", key, String(attempt)); let context; try { context = await this.contextFor(declaration, event); } catch (error) { if (error instanceof PublicKnowledgeEligibilitySkip || error instanceof ConversationCompactionNotNeeded || error instanceof AgentMessageNotAdmitted) { return this.settlePreflightSkip(declaration, event, key, runId, attempt, error); } throw error; } const runner = this.runners.get(declaration.mode); if (!runner) return { declaration, runId, error: `No runner for mode ${declaration.mode}` }; const startedAt = new Date().toISOString(); let run: AgentRun = { id: runId, executionKey: key, triggerEventId: event.id, agentId: declaration.id, agentVersion: declaration.version, status: "running", inputEventIds: [event.id], outputEventIds: [], attempt, provider: declaration.provider ?? declaration.mode, model: declaration.model ?? declaration.mode, privacy: outputPrivacy(declaration, event), adapterRevision: executionAdapterRevisionFor(declaration), executionAdapterRevision: executionAdapterRevisionFor(declaration), ...(declaration.modelAdapter ? { modelAdapter: declaration.modelAdapter } : {}), ...(declaration.adapterCatalogDigest ? { adapterCatalogDigest: declaration.adapterCatalogDigest } : {}), ...(declaration.adapterCatalogGeneration ? { adapterCatalogGeneration: declaration.adapterCatalogGeneration } : {}), promptHash: sha256(declaration.systemPrompt), contextManifest: { ...context.manifest, privacy: executionPrivacy(declaration, event), ...adapterEvidence(declaration), }, createdAt: startedAt, startedAt, updatedAt: startedAt, }; if (inferenceAccountingEnabled && declaration.accounting) { const reservationId = stableKey("inference-reservation", runId); run = { ...run, accountingReservationId: reservationId }; // Persist the execution owner before reserving. A crash after reservation must leave // a recoverable run, never an orphaned lease with no attempt evidence. await this.store.upsertRun(run); const reservation = await this.store.reserveInference({ reservationId, scopeType: "agent", scopeKey: declaration.id, runId, agentId: declaration.id, agentVersion: declaration.version, provider: run.provider, model: run.model, policy: declaration.accounting, estimate: { calls: 1, ...declaration.accounting.reservation }, reservedAt: startedAt, }); if (!reservation.approved) return this.settleBudgetBlockedRun(run, declaration, event, reservation.record.limitingWindowKeys ?? []); if (!reservation.acquired) { return { declaration, runId, error: "Inference reservation is already owned by another worker", retryable: true, }; } } else { await this.store.upsertRun(run); } await this.appendLifecycleEvent("started", declaration, event, run, { status: "running", executionKey: key, sourceSequence: event.sourceSequence, });
let traceSequence = 0; let output: AgentOutput; let preparedAt: string | undefined; let preparedOutputCandidate: EventCandidate | undefined; let preparedProposalCandidates: EventCandidate[] | undefined; let usage: InferenceUsage | undefined; try { const candidateOutput = await runner.run({ runId, attempt, declaration, event, context }, async (trace) => { traceSequence += 1; const at = new Date().toISOString(); await this.store.appendTrace({ id: stableKey("trace", runId, String(traceSequence)), runId, sequence: traceSequence, type: trace.kind, payload: { data: asJsonValue(trace.data) }, createdAt: trace.at ?? at, }); await this.appendLifecycleEvent("trace", declaration, event, run, { status: "running", sequence: traceSequence, traceType: trace.kind, data: asJsonValue(trace.data), }); }); try { const { model: reportedModel, enrichments, usage: providerUsage, proposals, ...semanticOutput } = candidateOutput; usage = providerUsage; const structured = this.outputContracts.validate( outputContractForDeclaration(declaration), semanticOutput, ); output = { ...structured, ...(declaration.modelAdapter ? { model: { provider: declaration.provider ?? declaration.modelAdapter.providerProfile, id: declaration.modelAdapter.baseModel, revision: publicAdapterRevision(declaration), }, } : reportedModel ? { model: reportedModel } : {}), ...(enrichments ? { enrichments } : {}), ...(providerUsage ? { usage: providerUsage } : {}), ...(proposals && proposals.length > 0 ? { proposals } : {}), }; } catch (error) { if (!(error instanceof OutputContractValidationError)) throw error; throw new AgentRunFailure("Agent final output rejected by canonical contract", { diagnostic: { code: "invalid-final-output", stage: "final-output-validation", reason: "output-contract-invalid", outputContract: outputContractIdentityJson(error.identity), validationIssues: error.issues.map((issue) => ({ code: issue.code, path: issue.path })), }, }); } preparedAt = new Date().toISOString(); preparedOutputCandidate = this.outputCandidate(declaration, event, run, output, preparedAt); preparedProposalCandidates = await this.proposalCandidates( declaration, event, run, output, deterministicEventId(preparedOutputCandidate), preparedAt, ); } catch (error) { const completedAt = new Date().toISOString(); if (error instanceof AgentRunFailure && error.usage) usage = error.usage; await this.settleAccountingReservation(run, usage, completedAt); const message = error instanceof AgentRunFailure ? error.message : "Agent runner failed"; const rawFailureDiagnostic = error instanceof AgentRunFailure ? error.diagnostic : { code: "unclassified-run-failure", stage: "runner" }; const isRepairRun = (declaration.role ?? "standard") === "repair"; const shouldAdvance = isRepairRun || !(error instanceof AgentRunFailure && !error.advanceProgress); const retryAt = !shouldAdvance && declaration.retry ? retryAtForFailedAttempt(completedAt, attempt, declaration) : undefined; const retryableDiagnostic = { ...(rawFailureDiagnostic ?? {}), ...(retryAt ? { retryAt, progressDisposition: "retry-delayed" } : {}), }; const failureDiagnostic = declaration.modelAdapter ? { ...retryableDiagnostic, provider: declaration.provider ?? declaration.modelAdapter.providerProfile, model: declaration.modelAdapter.baseModel, checkpointRevision: publicAdapterRevision(declaration), } : retryableDiagnostic; const diagnosticProvider = jsonStringField(failureDiagnostic, "provider"); const diagnosticModel = jsonStringField(failureDiagnostic, "model"); const diagnosticRevision = jsonStringField(failureDiagnostic, "checkpointRevision"); run = { ...run, status: "failed", errorText: message, ...(diagnosticProvider ? { provider: diagnosticProvider } : {}), ...(diagnosticModel ? { model: diagnosticModel } : {}), ...(diagnosticRevision ? { checkpointRevision: diagnosticRevision } : {}), ...(failureDiagnostic ? { result: { failureDiagnostic } } : {}), executionAdapterRevision: executionAdapterRevisionFor(declaration), completedAt, updatedAt: completedAt, }; const terminal = this.lifecycleCandidate("failed", declaration, event, run, completedAt, { status: "failed", error: message, ...(failureDiagnostic ? { failureDiagnostic } : {}), }); if (!shouldAdvance) { await this.store.settleConsumerAbandonment(run, terminal); } else { await this.store.settleConsumerFailure({ run, inputEvent: event, terminal, progress: progressFor(declaration, event, completedAt), }); if (!isRepairRun) await this.repairCoordinator.appendForRun(run, declaration); await this.rebuildEffectiveOutput(run.id); } return { declaration, runId, error: message, ...(!shouldAdvance ? { retryable: true } : {}), }; }
if (!preparedAt || !preparedOutputCandidate || !preparedProposalCandidates) { throw new Error("Successful agent run is missing prepared settlement candidates"); } const completedAt = preparedAt; await this.settleAccountingReservation(run, usage, completedAt); const outputCandidate = preparedOutputCandidate; const outputEventId = deterministicEventId(outputCandidate); const proposalCandidates = preparedProposalCandidates; const proposalEventIds = proposalCandidates.map((candidate) => deterministicEventId(candidate)); run = { ...run, status: "completed", outputEventIds: [outputEventId], ...(output.model ? { provider: output.model.provider, model: output.model.id, ...(output.model.revision ? { checkpointRevision: output.model.revision } : {}), } : {}), result: persistedRunResult(output), executionAdapterRevision: executionAdapterRevisionFor(declaration), completedAt, updatedAt: completedAt, }; const completed = this.lifecycleCandidate("completed", declaration, event, run, completedAt, { status: "completed", outputEventId, ...(proposalEventIds.length > 0 ? { proposalEventIds } : {}), }); const settled = await this.store.settleConsumerSuccess({ run, inputEvent: event, output: outputCandidate, ...(proposalCandidates.length > 0 ? { sideEffects: proposalCandidates } : {}), completed, progress: progressFor(declaration, event, completedAt), }); await this.rebuildEffectiveOutput(run.id); return { declaration, runId, output, derivedEvent: settled.outputEvent, }; }
private async settleAccountingReservation( run: AgentRun, usage: InferenceUsage | undefined, at: string, ): Promise<void> { if (!run.accountingReservationId) return; await this.store.settleInferenceReservation(run.accountingReservationId, usage, at); }
private async settleBudgetBlockedRun( run: AgentRun, declaration: ThoughtAgentDeclaration, event: ThoughtEvent, limitingWindowKeys: string[], ): Promise<ProcessEventResult> { const at = new Date().toISOString(); const defer = isDeferredExhaustion(declaration); const retryAt = defer && declaration.accounting ? retryAtForBudget(declaration.accounting, limitingWindowKeys, at) : undefined; const diagnostic: JsonObject = { code: "inference-budget-exhausted", stage: "accounting-reservation", reservationId: run.accountingReservationId ?? "", limitingWindowKeys, ...(retryAt ? { retryAt, progressDisposition: "deferred" } : {}), }; const blocked: AgentRun = { ...run, status: "blocked", errorText: "Inference budget exhausted before provider dispatch", result: { failureDiagnostic: diagnostic }, completedAt: at, updatedAt: at, }; const terminal = this.blockedLifecycleCandidate(declaration, event, blocked, at); if (defer) { await this.store.settleConsumerAbandonment(blocked, terminal); } else { await this.store.settleConsumerFailure({ run: blocked, inputEvent: event, terminal, progress: progressFor(declaration, event, at), }); } return { declaration, runId: blocked.id, error: blocked.errorText ?? "Inference budget exhausted", ...(defer ? { retryable: true } : {}), }; }
private async settleInterruptedBlockedRun( run: AgentRun, declaration: ThoughtAgentDeclaration, event: ThoughtEvent, ): Promise<void> { const at = run.completedAt ?? run.updatedAt; await this.store.settleConsumerFailure({ run, inputEvent: event, terminal: this.blockedLifecycleCandidate(declaration, event, run, at), progress: progressFor(declaration, event, at), }); }
private blockedLifecycleCandidate( declaration: ThoughtAgentDeclaration, event: ThoughtEvent, run: AgentRun, at: string, ): EventCandidate { const diagnostic = run.result?.failureDiagnostic; return this.lifecycleCandidate("blocked", declaration, event, run, at, { status: "blocked", reason: "inference-budget-exhausted", ...(diagnostic && typeof diagnostic === "object" && !Array.isArray(diagnostic) ? { failureDiagnostic: diagnostic as JsonObject } : {}), }); }
private async abandonInterruptedRepairRun( prior: AgentRun, declaration: ThoughtAgentDeclaration, event: ThoughtEvent, ): Promise<void> { const at = new Date().toISOString(); const abandoned: AgentRun = { ...prior, status: "abandoned", errorText: "Interrupted repair run is terminal; no second proposal attempt is permitted", completedAt: at, updatedAt: at, }; await this.settleInterruptedAccounting(prior, at); await this.store.settleConsumerFailure({ run: abandoned, inputEvent: event, terminal: this.lifecycleCandidate("abandoned", declaration, event, abandoned, at, { status: "abandoned", reason: "repair-restart-no-retry", }), progress: progressFor(declaration, event, at), }); }
private async abandonInterruptedRun( prior: AgentRun, declaration: ThoughtAgentDeclaration, event: ThoughtEvent, ): Promise<void> { const at = new Date().toISOString(); const abandoned: AgentRun = { ...prior, status: "abandoned", errorText: "Interrupted before terminal consumer settlement", completedAt: at, updatedAt: at, }; await this.settleInterruptedAccounting(prior, at); await this.store.settleConsumerAbandonment(abandoned, this.lifecycleCandidate("abandoned", declaration, event, abandoned, at, { status: "abandoned", reason: "restart-recovery", })); }
private async settleInterruptedAccounting(run: AgentRun, at: string): Promise<void> { if (!run.accountingReservationId) return; const reservation = await this.store.getInferenceAccountingRecord(run.accountingReservationId); if (reservation?.status === "reserved") { await this.store.settleInferenceReservation(run.accountingReservationId, undefined, at); } }
private consumerQuery(declaration: ThoughtAgentDeclaration, source: string, afterSequence: number): ConsumerEventQuery { return { consumerId: declaration.id, consumerVersion: declaration.version, source, eventTypes: declaration.compiledEventTypes, acceptedPrivacy: declaration.acceptedPrivacy, afterSequence, limit: 1_000, }; }
private async resolveSources(declaration: ThoughtAgentDeclaration): Promise<string[]> { const exact = declaration.sourcePatterns.filter((pattern) => !pattern.includes("*")); const discovered = (await this.store.listSources()) .map((source) => source.id) .filter((source) => declaration.sourcePatterns.some((pattern) => matchesPatternValue(pattern, source))); return [...new Set([...exact, ...discovered])].sort(); }
private async rebuildEffectiveOutput(runId: string): Promise<void> { try { await rebuildEffectiveOutputForRun(this.store, runId); } catch (error) { const message = error instanceof Error ? error.message : "unknown projection error"; console.error(`effective-output projection rebuild failed for ${runId}: ${message}`); } }
private async contextFor(declaration: ThoughtAgentDeclaration, event: ThoughtEvent) { if ((declaration.role ?? "standard") !== "repair") { if (declaration.atprotoObjectContext && event.type === "stream.thought.derived.event.batch") { return buildAtprotoBatchContextPacket( this.store, declaration, event, this.atprotoObjectContext, ); } if (declaration.contextStrategy === "activity-batch" && event.type === "stream.thought.derived.event.batch") { return buildActivityBatchContextPacket(this.store, declaration, event); } if (declaration.atprotoObjectContext && event.type === "stream.thought.source.atproto.commit") { return buildDurableAtprotoObjectContextPacket( this.store, declaration, event, this.atprotoObjectContext, ); } if (declaration.contextStrategy === "coil-public-knowledge") { if (!this.publicKnowledgeContext) { throw new Error("Coil Public Knowledge context is not configured for this consumer process"); } return buildPublicKnowledgeContextPacket( this.store, declaration, event, this.publicKnowledgeContext, ); } if (declaration.contextStrategy === "telegram-conversation") { return declaration.contextDocumentSubscriptions?.length ? buildSubscribedTelegramConversationContextPacket(declaration, event, this.store) : buildTelegramConversationContextPacket(declaration, event, this.store); } if (declaration.contextStrategy === "telegram-compaction") { return buildTelegramConversationCompactionContextPacket(declaration, event, this.store); } if (declaration.contextStrategy === "agent-conversation") { return buildSubscribedAgentConversationContextPacket(declaration, event, this.store); } return buildContextPacket(declaration, event); } if (event.type !== "stream.thought.agent.repair.requested") { throw new Error(`Repair agent ${declaration.id} received a non-repair event`); } const originalTriggerEventId = jsonStringField(event.payload, "originalTriggerEventId"); const originalAgentId = jsonStringField(event.payload, "originalAgentId"); const originalAgentVersion = jsonPositiveIntegerField(event.payload, "originalAgentVersion"); if (!originalTriggerEventId || !originalAgentId || !originalAgentVersion) { throw new Error("Repair request is missing original trigger or declaration references"); } const originalDeclaration = this.declarationsByVersion.get(`${originalAgentId}@${originalAgentVersion}`); if (!originalDeclaration) throw new Error(`Repair declaration evidence is unavailable: ${originalAgentId}@${originalAgentVersion}`); const originalSource = await this.store.getEvent(originalTriggerEventId); if (!originalSource) throw new Error(`Repair source evidence is unavailable: ${originalTriggerEventId}`); return buildRepairContextPacket(declaration, event, originalSource, originalDeclaration); }
private async settlePreflightSkip( declaration: ThoughtAgentDeclaration, event: ThoughtEvent, execution: string, runId: string, attempt: number, skip: { code: string; evidence: JsonObject; message: string }, ): Promise<ProcessEventResult> { const at = new Date().toISOString(); const run: AgentRun = { id: runId, executionKey: execution, triggerEventId: event.id, agentId: declaration.id, agentVersion: declaration.version, status: "skipped", inputEventIds: [event.id], outputEventIds: [], attempt, provider: declaration.provider ?? declaration.mode, model: declaration.model ?? declaration.mode, privacy: outputPrivacy(declaration, event), adapterRevision: executionAdapterRevisionFor(declaration), executionAdapterRevision: executionAdapterRevisionFor(declaration), promptHash: sha256(declaration.systemPrompt), contextManifest: { inputEventIds: [event.id], includedEventIds: [], omittedEventIds: [event.id], skip: { code: skip.code, evidence: skip.evidence }, privacy: executionPrivacy(declaration, event), ...adapterEvidence(declaration), }, errorText: skip.message, createdAt: at, startedAt: at, completedAt: at, updatedAt: at, }; await this.store.upsertRun(run); await this.store.settleConsumerFailure({ run, inputEvent: event, terminal: this.lifecycleCandidate("skipped", declaration, event, run, at, { status: "skipped", reasonCode: skip.code, evidence: skip.evidence, }), progress: progressFor(declaration, event, at), }); return { declaration, runId, error: skip.message }; }
private async consumerProgressForStart( declaration: ThoughtAgentDeclaration, source: string, ): Promise<ConsumerProgress | undefined> { const id = progressId(declaration, source); const existing = await this.store.getConsumerProgress(id); if (existing || declaration.initialReplay !== "now") return existing; const sourceState = (await this.store.listSources()).find((candidate) => candidate.id === source); const lastSequence = sourceState?.lastSequence ?? 0; const lastEvent = lastSequence > 0 ? await this.store.latestSourceEvent(source) : undefined; if (lastSequence > 0 && (!lastEvent || lastEvent.sourceSequence !== lastSequence)) { throw new Error(`Source head evidence mismatch for ${source}`); } const updatedAt = new Date().toISOString(); return this.store.initializeConsumerProgress({ id, consumerId: declaration.id, consumerVersion: declaration.version, source, lastSequence, lastEventId: lastEvent?.id ?? "", updatedAt, }); }
private async proposalCandidates( declaration: ThoughtAgentDeclaration, event: ThoughtEvent, run: AgentRun, output: AgentOutput, outputEventId: string, at: string, ): Promise<EventCandidate[]> { const proposals = output.proposals ?? []; if (proposals.length === 0) return []; if ((declaration.role ?? "standard") !== "standard" || declaration.outputMode !== "conversation-text") { throw new Error("Only standard conversation outputs may settle agent proposals"); } const capabilities = proposalCapabilitiesSchema.parse(run.contextManifest.proposalCapabilities); const snapshot = asJsonObject(run.contextManifest.contextSnapshot); const contextSnapshotId = jsonStringField(snapshot, "id"); if (!contextSnapshotId) throw new Error("Agent proposal is missing its immutable context snapshot"); const admittedEvidence = new Set(capabilities.evidenceEventIds); const proposer = { runId: run.id, outputEventId, triggerEventId: event.id, agentId: declaration.id, agentVersion: declaration.version, declarationFingerprint: declaration.declarationFingerprint ?? declarationFingerprint(declaration), provider: output.model?.provider ?? run.provider, model: output.model?.id ?? run.model, contextSnapshotId, }; const candidates: EventCandidate[] = []; for (let index = 0; index < proposals.length; index += 1) { const proposal = validateCapturedProposalAgainstCapabilities(proposals[index]!, capabilities); if (proposal.kind !== "focus-declaration") { for (const evidenceId of proposal.arguments.evidence_event_ids) { if (!admittedEvidence.has(evidenceId)) throw new Error("Agent proposal cited evidence outside its context snapshot"); } } const base = { schemaVersion: 1, source: `agent:${declaration.id}`, sourceKind: "agent" as const, externalId: `${run.id}:proposal:${index + 1}`, idempotencyKey: `${run.id}:proposal:${index + 1}:${proposal.kind}`, occurredAt: at, actor: declaration.id, privacy: "sensitive" as const, traceId: run.id, }; if (proposal.kind === "memory-change") { if (!capabilities.memoryTarget) throw new Error("Memory proposal has no snapshot-bound target"); candidates.push({ ...base, type: MEMORY_PROPOSAL_EVENT_TYPE, rootEventId: event.rootEventId, parentEventId: outputEventId, correlationId: run.id, payload: { proposalState: "agent-proposed", proposer, target: capabilities.memoryTarget, operation: proposal.arguments.operation, proposedText: proposal.arguments.proposed_text, proposedTextChars: proposal.arguments.proposed_text.length, proposedTextSha256: sha256(proposal.arguments.proposed_text), reason: proposal.arguments.reason, evidenceEventIds: proposal.arguments.evidence_event_ids, publicationEligible: false, }, }); continue; } if (proposal.kind === "focus-declaration") { if (telegramFocusDescription(event) === undefined) { throw new Error("Focus proposals require an exact description-bearing Telegram /focus command"); } const focus = proposal.arguments.declaration; const focusIdentity = `${focus.id}@${focus.version}`; candidates.push({ ...base, type: FOCUS_DECLARATION_PROPOSED_EVENT_TYPE, externalId: focusIdentity, idempotencyKey: base.idempotencyKey, rootEventId: event.rootEventId, parentEventId: outputEventId, correlationId: run.id, payload: { declaration: focus as unknown as JsonObject, declarationFingerprint: focusDeclarationFingerprint(focus), proposer: { kind: "agent", runId: proposer.runId, outputEventId: proposer.outputEventId, triggerEventId: proposer.triggerEventId, agentId: proposer.agentId, agentVersion: proposer.agentVersion, agentDeclarationFingerprint: proposer.declarationFingerprint, provider: proposer.provider, model: proposer.model, contextSnapshotId: proposer.contextSnapshotId, }, resolution: { eventTypes: "unchecked", sourcePatterns: "unchecked", requestedPermissions: "unchecked", parentLineage: "unchecked", }, authority: { activatesFocus: false, expandsDataAccess: false, grantsExternalActions: false, startsTraining: false, promotesAdapter: false, }, }, }); continue; } const resolvedTargetOutput = resolveCorrectionTarget(proposal.arguments.target_output, capabilities); const target = capabilities.correctionTargets.find((candidate) => ( candidate.outputEventId === resolvedTargetOutput )); if (!target) throw new Error("Correction proposal target is outside its context snapshot"); const [targetRun, targetOutput, delivery] = await Promise.all([ this.store.getRun(target.runId), this.store.getEvent(target.outputEventId), this.store.getEvent(target.deliveryReceiptEventId), ]); const chatId = jsonStringField(event.payload, "chatId"); if ( !targetRun || targetRun.status !== "completed" || targetRun.outputEventIds.length !== 1 || targetRun.outputEventIds[0] !== target.outputEventId || !targetOutput || targetOutput.type !== "stream.thought.derived.message.observation" || targetOutput.source !== `agent:${targetRun.agentId}` || targetOutput.sourceKind !== "agent" || targetOutput.payload.runId !== target.runId || targetOutput.rootEventId !== target.sourceRootEventId || !delivery || delivery.type !== "stream.thought.action.telegram.send.delivered" || delivery.source !== `telegram-dispatcher:${event.source}:${chatId}` || delivery.parentEventId !== target.outputEventId || delivery.payload.chatId !== chatId || delivery.rootEventId !== target.sourceRootEventId || !Array.isArray(delivery.payload.runIds) || delivery.payload.runIds.length !== 1 || delivery.payload.runIds[0] !== target.runId ) { throw new Error("Correction proposal target evidence is incomplete or inconsistent"); } const outputContract = outputContractIdentityJson(parseOutputContractIdentity(targetRun.contextManifest.outputContract)); if ( canonicalJson(outputContract) !== canonicalJson(target.outputContract as unknown as JsonObject) || canonicalJson(outputContract) !== canonicalJson(asJsonObject(targetOutput.payload.outputContract)) ) { throw new Error("Correction proposal target output contract is inconsistent"); } const originalOutput = canonicalStructuredOutput( this.outputContracts, target.outputContract, asJsonObject(targetOutput.payload.structuredOutput), ); const replacementOutput = canonicalStructuredOutput( this.outputContracts, target.outputContract, { ...originalOutput, summary: proposal.arguments.replacement }, ); candidates.push({ ...base, type: CORRECTION_PROPOSAL_EVENT_TYPE, rootEventId: target.sourceRootEventId, parentEventId: target.outputEventId, correlationId: target.runId, payload: { proposalState: "agent-proposed", proposer, target, replacementOutput, replacementText: proposal.arguments.replacement, replacementTextChars: proposal.arguments.replacement.length, replacementTextSha256: sha256(proposal.arguments.replacement), reason: proposal.arguments.reason, evidenceEventIds: proposal.arguments.evidence_event_ids, qualityEligible: false, externalExportEligible: false, publicationEligible: false, }, }); } return candidates; }
private outputCandidate( declaration: ThoughtAgentDeclaration, event: ThoughtEvent, run: AgentRun, result: AgentOutput, at: string, ): EventCandidate { if ((declaration.role ?? "standard") === "repair") { return this.repairProposalCandidate(declaration, event, run, result, at); } if ((declaration.role ?? "standard") === "compactor") { if (!isConversationCompactionOutput(result) || declaration.outputEventType !== CONVERSATION_COMPACTION_EVENT_TYPE) { throw new Error("Compactor output does not satisfy the conversation boundary shape"); } const compactionPlan = conversationCompactionPlanSchema.parse(run.contextManifest.compactionPlan); const structuredOutput = canonicalStructuredOutput( this.outputContracts, outputContractForDeclaration(declaration), semanticOutputForPersistence(result), ); const semanticCompaction = conversationCompactionOutputSchema.parse(structuredOutput); const contextSnapshot = conversationCompactionContextSnapshotSchema.parse( run.contextManifest.contextSnapshot, ); const boundarySha256 = conversationCompactionBoundarySha256(compactionPlan, semanticCompaction); return { type: declaration.outputEventType, schemaVersion: 1, source: `agent:${declaration.id}`, sourceKind: "agent", externalId: run.id, idempotencyKey: `${run.id}:output`, occurredAt: at, actor: declaration.id, rootEventId: event.rootEventId, parentEventId: event.id, correlationId: event.correlationId, privacy: "sensitive", payload: { runId: run.id, executionKey: run.executionKey, inputEventId: event.id, inputSourceSequence: event.sourceSequence, compactorAgentVersion: declaration.version, compactorDeclarationFingerprint: declaration.declarationFingerprint ?? declarationFingerprint(declaration), promptHash: run.promptHash, contextSnapshot: contextSnapshot as unknown as JsonObject, summary: result.summary, confidence: result.confidence, outputContract: outputContractIdentityJson(outputContractForDeclaration(declaration)), structuredOutput, compactionPlan: compactionPlan as unknown as JsonObject, boundarySha256, ...(result.model ? { model: asJsonObject(result.model) } : {}), }, traceId: run.id, }; } if (declaration.outputEventType === AGENT_MESSAGE_RESPONSE_EVENT_TYPE) { const observation = isObservationOutput(result) ? result : undefined; if (!observation || declaration.contextStrategy !== "agent-conversation") { throw new Error("Agent message response requires an observation from agent-conversation context"); } const input = assertAgentMessageRoute(event, declaration.id); const structuredOutput = canonicalStructuredOutput( this.outputContracts, outputContractForDeclaration(declaration), semanticOutputForPersistence(result), ); return { type: declaration.outputEventType, schemaVersion: 1, source: `agent:${declaration.id}`, sourceKind: "agent", externalId: run.id, idempotencyKey: `${run.id}:output`, occurredAt: at, actor: declaration.id, rootEventId: event.rootEventId, parentEventId: event.id, correlationId: input.threadId, privacy: "sensitive", payload: { messageId: run.id, inReplyToMessageId: input.messageId, threadId: input.threadId, senderAgentId: declaration.id, recipientAgentId: input.senderAgentId, runId: run.id, executionKey: run.executionKey, inputEventId: event.id, inputSourceSequence: event.sourceSequence, agentVersion: declaration.version, declarationFingerprint: declaration.declarationFingerprint ?? declarationFingerprint(declaration), summary: observation.summary, tags: observation.tags, importance: observation.importance, confidence: observation.confidence, outputContract: outputContractIdentityJson(outputContractForDeclaration(declaration)), structuredOutput, ...(result.model ? { model: asJsonObject(result.model) } : {}), }, traceId: run.id, }; } const observation = isObservationOutput(result) ? result : undefined; const reviewResponse = isReviewResponseOutput(result) ? result : undefined; return { type: declaration.outputEventType, schemaVersion: 1, source: `agent:${declaration.id}`, sourceKind: "agent", externalId: run.id, idempotencyKey: `${run.id}:output`, occurredAt: at, actor: declaration.id, rootEventId: event.rootEventId, parentEventId: event.id, correlationId: event.correlationId, privacy: outputPrivacy(declaration, event), payload: { runId: run.id, executionKey: run.executionKey, inputEventId: event.id, inputSourceSequence: event.sourceSequence, summary: reviewResponse ? reviewResponseSummary(reviewResponse.response) : result.summary, ...(observation ? { tags: observation.tags, importance: observation.importance } : {}), ...(observation?.recommendation ? { recommendation: observation.recommendation as JsonObject } : {}), ...(reviewResponse ? {} : { confidence: result.confidence }), outputContract: outputContractIdentityJson(outputContractForDeclaration(declaration)), structuredOutput: canonicalStructuredOutput( this.outputContracts, outputContractForDeclaration(declaration), semanticOutputForPersistence(result), ), ...(result.model ? { model: asJsonObject(result.model) } : {}), ...adapterEvidence(declaration), }, traceId: run.id, }; }
private repairProposalCandidate( declaration: ThoughtAgentDeclaration, event: ThoughtEvent, run: AgentRun, result: AgentOutput, at: string, ): EventCandidate { if (event.type !== "stream.thought.agent.repair.requested") throw new Error("Repair proposals require a repair-request input"); const originalRunId = jsonStringField(event.payload, "originalRunId"); const originalTriggerEventId = jsonStringField(event.payload, "originalTriggerEventId"); const sourceRootEventId = jsonStringField(event.payload, "sourceRootEventId"); if (!originalRunId || !originalTriggerEventId || !sourceRootEventId || sourceRootEventId !== event.rootEventId) { throw new Error("Repair proposal lineage is incomplete or inconsistent"); } const requestContract = event.payload.outputContract; const declarationContract = outputContractIdentityJson(outputContractForDeclaration(declaration)); if (canonicalJson(asJsonObject(requestContract)) !== canonicalJson(declarationContract)) { throw new Error("Repair proposal output contract does not match the request"); } return { type: declaration.outputEventType, schemaVersion: 1, source: `agent:${declaration.id}`, sourceKind: "agent", externalId: run.id, idempotencyKey: `${run.id}:proposal`, occurredAt: at, actor: declaration.id, rootEventId: sourceRootEventId, parentEventId: event.id, correlationId: originalRunId, privacy: outputPrivacy(declaration, event), payload: { originalRunId, repairRequestEventId: event.id, repairRunId: run.id, originalTriggerEventId, sourceRootEventId, outputContract: declarationContract, originalModel: asJsonObject(event.payload.model), repairModel: runModelEvidence(run, declaration), structuredOutput: canonicalStructuredOutput( this.outputContracts, outputContractForDeclaration(declaration), semanticOutputForPersistence(result), ), ...adapterEvidence(declaration), }, traceId: run.id, }; }
private async appendLifecycleEvent( phase: "started" | "trace", declaration: ThoughtAgentDeclaration, event: ThoughtEvent, run: AgentRun, payload: JsonObject, ): Promise<void> { await this.store.appendEvent(this.lifecycleCandidate(phase, declaration, event, run, new Date().toISOString(), payload)); }
private lifecycleCandidate( phase: "started" | "trace" | "completed" | "failed" | "blocked" | "skipped" | "abandoned", declaration: ThoughtAgentDeclaration, event: ThoughtEvent, run: AgentRun, at: string, payload: JsonObject, ): EventCandidate { const sequence = phase === "trace" && typeof payload.sequence === "number" ? `:${payload.sequence}` : ""; return { type: `stream.thought.agent.run.${phase}`, schemaVersion: 1, source: `agent:${declaration.id}`, sourceKind: "agent", externalId: run.id, idempotencyKey: `${run.id}:${phase}${sequence}`, occurredAt: at, actor: declaration.id, rootEventId: event.rootEventId, parentEventId: event.id, correlationId: event.correlationId, privacy: executionPrivacy(declaration, event), payload: { runId: run.id, executionKey: run.executionKey, agentId: declaration.id, agentVersion: declaration.version, agentRole: declaration.role ?? "standard", declarationFingerprint: declaration.declarationFingerprint ?? declarationExecutionFingerprint(declaration), outputContract: outputContractIdentityJson(outputContractForDeclaration(declaration)), inputEventIds: [event.id], attempt: run.attempt, ...adapterEvidence(declaration), ...payload, }, traceId: run.id, }; }}
function assertRuntimeAccountingPolicy(declaration: ThoughtAgentDeclaration): void { if (!inferenceAccountingEnabled) return; if ((declaration.mode === "pi" || declaration.mode === "letta-agent-sdk") && !declaration.accounting) { throw new Error(`Model-backed agent ${declaration.id}@${declaration.version} has no inference accounting policy`); } if (declaration.mode === "deterministic" && declaration.accounting) { throw new Error(`Deterministic agent ${declaration.id}@${declaration.version} cannot reserve provider inference`); }}
function isDeferredExhaustion(declaration: ThoughtAgentDeclaration): boolean { return declaration.accounting?.onExhaustion === "defer";}
function retryAtForPendingRun(run: AgentRun, declaration: ThoughtAgentDeclaration): string | undefined { if (run.status === "blocked" && !isDeferredExhaustion(declaration)) return undefined; if (run.status === "failed" && !declaration.retry) return undefined; if (run.status !== "failed" && run.status !== "blocked") return undefined; const diagnostic = run.result?.failureDiagnostic; if (!diagnostic || typeof diagnostic !== "object" || Array.isArray(diagnostic)) return undefined; const retryAt = (diagnostic as JsonObject).retryAt; if (typeof retryAt !== "string" || !Number.isFinite(Date.parse(retryAt))) return undefined; return retryAt;}
function retryAtForFailedAttempt(at: string, attempt: number, declaration: ThoughtAgentDeclaration): string { if (!declaration.retry) throw new Error("Retry timing requires a declaration retry policy"); const exponent = Math.max(0, Math.min(30, attempt - 1)); const delayMs = Math.min(declaration.retry.maxDelayMs, declaration.retry.initialDelayMs * (2 ** exponent)); return new Date(Date.parse(at) + delayMs).toISOString();}
function retryAtForBudget(policy: InferenceBudgetPolicy, limitingWindowKeys: string[], at: string): string { const deniedAt = Date.parse(at); const candidates = limitingWindowKeys.map((key) => { if (key.startsWith("rolling:")) { const durationMs = Number(key.slice("rolling:".length)); return Number.isSafeInteger(durationMs) && durationMs > 0 ? deniedAt + durationMs : deniedAt + policy.leaseMs; } if (key === "hour") { const boundary = new Date(deniedAt); boundary.setUTCMinutes(0, 0, 0); return boundary.getTime() + 3_600_000; } if (key === "day") { const boundary = new Date(deniedAt); boundary.setUTCHours(0, 0, 0, 0); return boundary.getTime() + 86_400_000; } return deniedAt + policy.leaseMs; }); return new Date(Math.max(deniedAt + 1_000, ...candidates)).toISOString();}
function assertOutputContractEventBinding(declaration: ThoughtAgentDeclaration): void { const contractId = outputContractForDeclaration(declaration).id; const conceptualizationOutput = contractId === CONCEPTUALIZATION_OUTPUT_CONTRACT_ID; const conceptGraphEvent = declaration.outputEventType === "stream.thought.derived.concept.graph"; if ((declaration.role ?? "standard") === "standard" && conceptualizationOutput !== conceptGraphEvent) { throw new Error(`Agent ${declaration.id}@${declaration.version} must bind conceptualization output and concept graph events together`); } const reviewResponseOutput = contractId === REVIEW_RESPONSE_OUTPUT_CONTRACT_ID; const reviewResponseEvent = declaration.outputEventType === REVIEW_RESPONSE_EVENT_TYPE; if ((declaration.role ?? "standard") === "standard" && reviewResponseOutput !== reviewResponseEvent) { throw new Error(`Agent ${declaration.id}@${declaration.version} must bind review response output and review response events together`); } const conversationCompactionOutput = contractId === CONVERSATION_COMPACTION_OUTPUT_CONTRACT_ID; const conversationCompactionEvent = declaration.outputEventType === CONVERSATION_COMPACTION_EVENT_TYPE; if ((declaration.role ?? "standard") === "compactor" && (!conversationCompactionOutput || !conversationCompactionEvent)) { throw new Error(`Compactor ${declaration.id}@${declaration.version} must bind the conversation compaction contract and event`); } if ((declaration.role ?? "standard") !== "compactor" && (conversationCompactionOutput || conversationCompactionEvent)) { throw new Error(`Agent ${declaration.id}@${declaration.version} cannot use compaction output without the compactor role`); } const publicKnowledgeOutput = contractId === PUBLIC_KNOWLEDGE_RECOMMENDATION_OUTPUT_CONTRACT_ID; const publicKnowledgeEvent = declaration.outputEventType === "stream.thought.derived.public-knowledge.recommendation"; if ((declaration.role ?? "standard") === "standard" && publicKnowledgeOutput !== publicKnowledgeEvent) { throw new Error(`Agent ${declaration.id}@${declaration.version} must bind Public Knowledge recommendation output and events together`); } const publicKnowledgeProposedDiffOutput = contractId === PUBLIC_KNOWLEDGE_PROPOSED_DIFF_OUTPUT_CONTRACT_ID; const publicKnowledgeProposedDiffEvent = declaration.outputEventType === PUBLIC_KNOWLEDGE_PROPOSED_DIFF_EVENT_TYPE; if ((declaration.role ?? "standard") === "standard" && publicKnowledgeProposedDiffOutput !== publicKnowledgeProposedDiffEvent) { throw new Error(`Agent ${declaration.id}@${declaration.version} must bind Public Knowledge proposed-diff output and events together`); } const conversationText = (declaration.mode === "pi" && declaration.outputMode === "conversation-text") || (declaration.mode === "letta-agent-sdk" && declaration.lettaAgent?.responseMode === "conversation-text"); if (conceptualizationOutput && conversationText) { throw new Error(`Agent ${declaration.id}@${declaration.version} requires strict JSON for conceptualization output`); } if (publicKnowledgeOutput && conversationText) { throw new Error(`Agent ${declaration.id}@${declaration.version} requires strict JSON for Public Knowledge recommendations`); } if (publicKnowledgeProposedDiffOutput && conversationText) { throw new Error(`Agent ${declaration.id}@${declaration.version} requires strict JSON for Public Knowledge proposed diffs`); } if (conversationText && contractId !== OBSERVATION_OUTPUT_CONTRACT_ID) { throw new Error(`Agent ${declaration.id}@${declaration.version} requires the observation contract for conversation-text output`); }}
function assertUniqueLettaAgentOwners(declarations: ThoughtAgentDeclaration[]): void { const owners = new Map<string, string>(); for (const declaration of declarations) { if (!declaration.enabled || declaration.mode !== "letta-agent-sdk") continue; const agentId = declaration.lettaAgent?.agentId; if (!agentId) throw new Error(`Enabled Letta Agent SDK declaration ${declaration.id} has no resolved Cloud agent id`); const owner = owners.get(agentId); if (owner) { throw new Error(`Letta Cloud agent ${agentId} is shared by enabled declarations ${owner} and ${declaration.id}`); } owners.set(agentId, declaration.id); }}
function assertConversationCompactionPairing(declarations: ThoughtAgentDeclaration[]): void { const byId = new Map(declarations.map((declaration) => [declaration.id, declaration])); for (const declaration of declarations) { const policy = declaration.conversationCompaction; if (!policy) continue; if (policy.mode === "consume") { const compactor = byId.get(policy.agentId); if (!compactor || (compactor.role ?? "standard") !== "compactor" || compactor.conversationCompaction?.mode !== "produce" || compactor.conversationCompaction.targetAgentId !== declaration.id) { throw new Error(`Conversation agent ${declaration.id} does not have one reciprocal compactor ${policy.agentId}`); } if (declaration.enabled !== compactor.enabled) { throw new Error(`Conversation agent ${declaration.id} and compactor ${compactor.id} must activate together`); } const parentHistory = [...new Set(declaration.conversationHistoryAgentIds ?? [declaration.id])].sort(); const compactorHistory = [...new Set(compactor.conversationHistoryAgentIds ?? [declaration.id])].sort(); if (canonicalJson(parentHistory) !== canonicalJson(compactorHistory) || canonicalJson(declaration.contextDocumentSubscriptions as unknown as JsonObject[] ?? []) !== canonicalJson(compactor.contextDocumentSubscriptions as unknown as JsonObject[] ?? []) || declaration.contextDocumentMaxChars !== compactor.contextDocumentMaxChars || declaration.provider !== compactor.provider || declaration.providerProfile !== compactor.providerProfile || declaration.model !== compactor.model || canonicalJson(declaration.sourcePatterns) !== canonicalJson(compactor.sourcePatterns) || canonicalJson(declaration.acceptedPrivacy) !== canonicalJson(compactor.acceptedPrivacy)) { throw new Error(`Compactor ${compactor.id} is not an identity/context clone of ${declaration.id}`); } continue; } const target = byId.get(policy.targetAgentId); if (!target || target.conversationCompaction?.mode !== "consume" || target.conversationCompaction.agentId !== declaration.id) { throw new Error(`Compactor ${declaration.id} does not have one reciprocal target ${policy.targetAgentId}`); } }}
function schedulerIncidentContext(value: unknown): Omit<SchedulerExhaustionEvidence, "operationKey"> { if (!value || typeof value !== "object" || Array.isArray(value)) return {}; const context = value as Record<string, unknown>; return { ...(typeof context.source === "string" ? { source: context.source } : {}), ...(typeof context.agentId === "string" ? { agentId: context.agentId } : {}), ...(typeof context.agentVersion === "number" && Number.isSafeInteger(context.agentVersion) && context.agentVersion > 0 ? { agentVersion: context.agentVersion } : {}), };}
function consumerOperationKey(declaration: ThoughtAgentDeclaration, source: string): string { if (declaration.mode === "letta-agent-sdk") { const agentId = declaration.lettaAgent?.agentId; if (!agentId) throw new Error(`Enabled Letta Agent SDK declaration ${declaration.id} has no resolved Cloud agent id`); return `letta-agent:${agentId}`; } return `${declaration.id}@${declaration.version}:${source}`;}
function executionPrivacy( declaration: ThoughtAgentDeclaration, event: ThoughtEvent,): ThoughtEvent["privacy"] { return declaration.mode === "letta-agent-sdk" ? declarationPrivacy(declaration, ...declaration.acceptedPrivacy, event.privacy) : declarationPrivacy(declaration, event.privacy);}
function outputPrivacy( declaration: ThoughtAgentDeclaration, event: ThoughtEvent,): ThoughtEvent["privacy"] { const privacy = executionPrivacy(declaration, event); if (outputContractForDeclaration(declaration).id !== CONCEPTUALIZATION_OUTPUT_CONTRACT_ID) return privacy; return privacy === "public-source" ? "private" : privacy;}
function semanticOutputForPersistence(output: AgentOutput): JsonObject { const { model: _model, enrichments: _enrichments, usage: _usage, proposals: _proposals, ...semanticOutput } = output; return semanticOutput as JsonObject;}
function persistedRunResult(output: AgentOutput): JsonObject { const { proposals: _proposals, ...result } = output; return asJsonObject(result);}
function isObservationOutput( output: AgentOutput,): output is Extract<AgentOutput, { tags: string[] }> { return "tags" in output && "importance" in output;}
function isReviewResponseOutput( output: AgentOutput,): output is Extract<AgentOutput, { response: string }> { return "response" in output;}
function isConversationCompactionOutput( output: AgentOutput,): output is Extract<AgentOutput, { boundary: string }> { return "boundary" in output && Array.isArray(output.openLoops) && Array.isArray(output.decisions) && Array.isArray(output.exactReferences) && Array.isArray(output.unresolved) && Array.isArray(output.lookupHints);}
function adapterEvidence(declaration: ThoughtAgentDeclaration): JsonObject { const catalog = declaration.adapterCatalogDigest ? { adapterCatalogDigest: declaration.adapterCatalogDigest, ...(declaration.adapterCatalogGeneration ? { adapterCatalogGeneration: declaration.adapterCatalogGeneration } : {}), } : {}; if (!declaration.modelAdapter) return { executionAdapterRevision: executionAdapterRevisionFor(declaration), ...catalog }; return { executionAdapterRevision: executionAdapterRevisionFor(declaration), modelAdapter: declaration.modelAdapter as unknown as JsonObject, ...catalog, };}
function runModelEvidence(run: AgentRun, declaration: ThoughtAgentDeclaration): JsonObject { return { provider: run.provider, id: run.model, ...(run.checkpointRevision ? { checkpointRevision: run.checkpointRevision } : {}), ...adapterEvidence(declaration), };}
function publicAdapterRevision(declaration: ThoughtAgentDeclaration): string { if (!declaration.modelAdapter) throw new Error(`Agent ${declaration.id}@${declaration.version} has no model adapter`); return `sha256:${declaration.modelAdapter.binding.checkpointReferenceSha256}`;}
function executionKey(event: ThoughtEvent, declaration: ThoughtAgentDeclaration): string { return stableKey("execution", declaration.id, String(declaration.version), event.id);}
function progressId(declaration: ThoughtAgentDeclaration, source: string): string { return stableKey("consumer-progress", declaration.id, String(declaration.version), source);}
function progressFor(declaration: ThoughtAgentDeclaration, event: ThoughtEvent, at: string): ConsumerProgress { return { id: progressId(declaration, event.source), consumerId: declaration.id, consumerVersion: declaration.version, source: event.source, lastSequence: event.sourceSequence, lastEventId: event.id, updatedAt: at, };}
function deterministicEventId(candidate: EventCandidate): string { return `evt_${sha256(`${candidate.source}\u0000${candidate.idempotencyKey}`).slice(0, 48)}`;}
export function declarationExecutionFingerprint(declaration: ThoughtAgentDeclaration): string { return declarationFingerprint(declaration);}
function matchesPatternValue(pattern: string, value: string): boolean { if (pattern === "*") return true; if (pattern.endsWith("*")) return value.startsWith(pattern.slice(0, -1)); return pattern === value;}
function jsonPositiveIntegerField(value: JsonObject | undefined, key: string): number | undefined { const field = value?.[key]; return Number.isSafeInteger(field) && Number(field) > 0 ? Number(field) : undefined;}
function jsonStringField(value: JsonObject | undefined, key: string): string | undefined { const field = value?.[key]; return typeof field === "string" && field.length > 0 ? field : undefined;}
function asJsonObject(value: unknown): JsonObject { const normalized = JSON.parse(JSON.stringify(value)) as unknown; if (!normalized || typeof normalized !== "object" || Array.isArray(normalized)) throw new Error("Expected a JSON object"); return normalized as JsonObject;}
function asJsonValue(value: unknown): JsonObject[string] { return JSON.parse(JSON.stringify(value)) as JsonObject[string];}