import { 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; stop(): Promise; } 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(); private readonly outputContracts: OutputContractRegistry; private readonly repairCoordinator: RepairRequestCoordinator; private readonly declarationsByVersion = new Map(); 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 { 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 { 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 { await this.registerDeclarations(declarations); const unsubscribers: Array<() => void> = []; const subscriptions = new Set(); const reconcilers = new Map 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, context?: Omit): void => { scheduler.enqueue(key, operation, context); }; const enabled = declarations.filter((candidate) => candidate.enabled); const install = async (declaration: ThoughtAgentDeclaration, source: string): Promise => { 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 => { 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 { 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 { if (!run.accountingReservationId) return; await this.store.settleInferenceReservation(run.accountingReservationId, usage, at); } private async settleBudgetBlockedRun( run: AgentRun, declaration: ThoughtAgentDeclaration, event: ThoughtEvent, limitingWindowKeys: string[], ): Promise { 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 { 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 { 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 { 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 { 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 { 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 { 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 { 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 { 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 { 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 { 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(); 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 { if (!value || typeof value !== "object" || Array.isArray(value)) return {}; const context = value as Record; 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 { return "tags" in output && "importance" in output; } function isReviewResponseOutput( output: AgentOutput, ): output is Extract { return "response" in output; } function isConversationCompactionOutput( output: AgentOutput, ): output is Extract { 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]; }