Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
77 kB · 1606 lines
TypeScript
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607#!/usr/bin/env nodeimport fs from "node:fs/promises";import path from "node:path";import chokidar from "chokidar";import { canonicalJson, sha256, type JsonObject } from "./core/json.js";import { loadAgentDeclarations } from "./agents/declarations.js";import { inspectRunContext } from "./agents/context-inspection.js";import { ThoughtAgentRuntime } from "./agents/runtime.js";import { TelegramChannelDispatcher } from "./bridges/telegram-dispatcher.js";import { DeterministicBatcher, startBatchService } from "./batches/runtime.js";import { FilesystemConnector } from "./connectors/filesystem.js";import { FastmailConnector } from "./connectors/fastmail.js";import { FastmailJmapClient, FastmailJmapConnector, fastmailJmapErrorCode,} from "./connectors/fastmail-jmap.js";import { subscribeJetstream } from "./connectors/jetstream-live.js";import { JetstreamConnector } from "./connectors/jetstream.js";import { RssConnector } from "./connectors/rss.js";import { TELEGRAM_WEBHOOK_ALLOWED_UPDATES, TelegramBotClient, TelegramBotConnector, type TelegramWebhookInfo,} from "./connectors/telegram-bot.js";import { TelegramSpoolConnector } from "./connectors/telegram-spool.js";import { startTelegramWebhookServer } from "./connectors/telegram-webhook.js";import { XActivityConnector } from "./connectors/x-activity.js";import { X_WEBHOOK_REVISION } from "./connectors/x-contract.js";import { XApiClient, assertBoundedXReplayWindow, planXSubscriptions, type XSubscriptionPlan, type XWebhookRecord,} from "./connectors/x-client.js";import { startXWebhookServer } from "./connectors/x-webhook.js";import { JazzThoughtStore } from "./jazz/store.js";import { IncidentLedger } from "./incidents/ledger.js";import { listOperationalIncidents, OperationalIncidentProjector } from "./incidents/projector.js";import { IncidentTelegramDispatcher } from "./incidents/telegram-alerts.js";import { rebuildRootActivity } from "./projections/activity.js";import { updateXSourceControl, type SourceControlUpstreamState,} from "./projections/source-control.js";import { loadThoughtStreamManifest, type FastmailJmapSourceManifest, type TelegramWebhookSourceManifest, type XWebhookSourceManifest,} from "./runtime/manifest.js";import { projectTrainingExamples, recordJudgment, writeTrainingJsonl, type JudgmentKind,} from "./training/judgments.js";import { projectTelegramFeedbackJudgments } from "./training/telegram-reactions.js";import { startInspectorServer } from "./web/inspector.js";import { decodeReviewCapability } from "./review/web-capability.js";import { decodeCourseChatCapability } from "./courses/web-capability.js";import { appendReviewPrompt, createReviewItem, projectReviewQueue,} from "./review/review.js";import type { PrivacyClass } from "./events/types.js";import type { AgentRun } from "./store/types.js";import { decisionsForProposal, projectAcceptedCorrectionDecisions, recordProposalDecision, requireProposal, type ProposalDecisionDisposition,} from "./agent-proposals/review.js";import { materializeMemoryDecision } from "./agent-proposals/memory-materializer.js";import { STREAM_TELEGRAM_BOT_NAME, STREAM_TELEGRAM_COMMANDS, STREAM_TELEGRAM_HELP_TEXT, STREAM_TELEGRAM_HELP_VERSION,} from "./agents/telegram-help.js";import { CORRECTION_PROPOSAL_EVENT_TYPE, MEMORY_MATERIALIZED_EVENT_TYPE, MEMORY_PROPOSAL_EVENT_TYPE, PROPOSAL_DECISION_EVENT_TYPE,} from "./agent-proposals/contracts.js";import { exportPrivateTrainingDataset } from "./training/private-training.js";import { listFocusProposals, loadFocusDeclaration, proposeFocusDeclaration,} from "./focuses/declarations.js";import { CO_AGENT_ID, AGENT_MESSAGE_TEXT_MAX_CHARS, STREAM_AGENT_ENDPOINT_ID, appendCoAgentMessage, parseAgentMessageResponseEvent,} from "./agents/agent-messages.js";
const projectRoot = process.env.THOUGHTSTREAM_ROOT ?? process.cwd();if (process.argv[2] === "stream") { if (!process.argv[3] || process.argv[3].startsWith("--")) process.argv[2] = "serve"; else process.argv.splice(2, 1);}const command = process.argv[2] ?? "status";const store = await JazzThoughtStore.open({ projectRoot });let commandFailed = false;
try { if (command === "scan") { const root = path.resolve(valueAfter("--root") ?? path.join(projectRoot, "fixtures", "vault")); const source = valueAfter("--source") ?? "filesystem:fixture"; const connector = new FilesystemConnector({ id: source, root }); print({ scan: await connector.scan(store) }); } else if (command === "demo") { const root = path.resolve(valueAfter("--root") ?? path.join(projectRoot, "fixtures", "vault")); const source = valueAfter("--source") ?? "filesystem:fixture"; const connector = new FilesystemConnector({ id: source, root }); const declarations = await loadAgentDeclarations(path.join(projectRoot, "agents")); const runtime = await createAgentRuntime(); const consumers = await runtime.startConsumers(declarations); try { const scan = await connector.scan(store); const finalSequence = Math.max(0, ...scan.events.map((event) => event.sourceSequence)); if (finalSequence > 0) { await waitUntil(async () => (await store.listConsumerProgress()).some((progress) => ( progress.source === source && progress.lastSequence >= finalSequence ))); } await consumers.drain(); print({ scan, runs: await store.listRuns(), activity: await rebuildRootActivity(store) }); } finally { await consumers.stop(); } } else if (command === "watch") { const root = path.resolve(valueAfter("--root") ?? path.join(projectRoot, "fixtures", "vault")); const source = valueAfter("--source") ?? "filesystem:fixture"; const debounceMs = Number(valueAfter("--debounce") ?? "250"); if (!Number.isFinite(debounceMs) || debounceMs < 0) throw new Error("--debounce must be a nonnegative number"); const connector = new FilesystemConnector({ id: source, root }); const producerOnly = process.argv.includes("--producer-only"); const declarations = producerOnly ? undefined : await loadAgentDeclarations(path.join(projectRoot, "agents")); const runtime = producerOnly ? undefined : await createAgentRuntime(); const consumers = runtime && declarations ? await runtime.startConsumers(declarations) : undefined; const watcher = chokidar.watch(root, { ignoreInitial: true, followSymlinks: false, ignored: (candidate) => isIgnoredWatchPath(root, candidate), }); let timer: NodeJS.Timeout | undefined; let running = false; let pending = false; let activeCycle: Promise<void> | undefined; const schedule = () => { pending = true; if (timer) clearTimeout(timer); timer = setTimeout(drain, debounceMs); }; const drain = () => { if (running || !pending) return; pending = false; running = true; activeCycle = (async () => { try { const scan = await connector.scan(store); if (scan.events.length > 0) print({ scan }); } catch (error) { process.stderr.write(`thought stream watch cycle failed: ${error instanceof Error ? error.message : String(error)}\n`); } finally { running = false; activeCycle = undefined; if (pending) schedule(); } })(); }; watcher.on("all", schedule); await new Promise<void>((resolve, reject) => { const onReady = () => { watcher.off("error", onError); resolve(); }; const onError = (error: unknown) => { watcher.off("ready", onReady); reject(error); }; watcher.once("ready", onReady); watcher.once("error", onError); }); running = true; try { print({ scan: await connector.scan(store) }); } finally { running = false; if (pending) schedule(); } process.stdout.write(`thought stream watching read-only source ${source}${producerOnly ? " in producer-only mode" : ""}: ${root}\n`); await waitForSignal(); process.stderr.write("thought stream watcher stopping.\n"); pending = false; if (timer) clearTimeout(timer); const cycle = activeCycle; if (cycle) await cycle; await consumers?.drain(); await consumers?.stop(); watcher.removeAllListeners(); await watcher.unwatch(root); await watcher.close(); process.stderr.write("thought stream watcher stopped.\n"); } else if (command === "batch" || command === "batches") { const loaded = await loadThoughtStreamManifest(projectRoot, valueAfter("--config") ?? "thoughtstream.yaml"); const selectedId = valueAfter("--id"); const selected = loaded.manifest.batches.filter((declaration) => declaration.enabled && (!selectedId || declaration.id === selectedId)); if (selectedId && selected.length === 0) throw new Error(`No enabled batch declaration found: ${selectedId}`); if (selected.length === 0) throw new Error("No batch declarations are enabled"); if (process.argv.includes("--once") || command === "batch") { const batcher = new DeterministicBatcher(store); const results = []; for (const declaration of selected) results.push(await batcher.cycle(declaration)); print({ batches: results }); } else { const handle = await startBatchService(store, selected, { projectRoot, onCycle: (result) => { if (result.emitted > 0) print({ batch: result }); }, }); process.stdout.write(`thought stream batching ${selected.map((declaration) => `${declaration.id}@${declaration.version}`).join(", ")}\n`); try { const maxRuntime = valueAfter("--max-runtime"); await waitForSignalOrTimeout(maxRuntime ? positiveInteger(maxRuntime, "--max-runtime", 86_400) * 1_000 : undefined); } finally { await handle.stop(); } } } else if (command === "rss") { const feedUrl = valueAfter("--url"); if (!feedUrl) throw new Error("Usage: thought stream rss --url <feed-url> [--source rss:name]"); const source = valueAfter("--source") ?? `rss:${new URL(feedUrl).hostname}`; const connector = new RssConnector({ id: source, feedUrl }); const poll = await connector.poll(store); print({ poll }); } else if (command === "jetstream") { const source = valueAfter("--source"); const collections = commaSeparated(valueAfter("--collections")); const dids = commaSeparated(valueAfter("--dids")); if (!source || collections.length === 0) { throw new Error("Usage: thought stream jetstream --source jetstream:name --collections <nsid[,nsid...]> [--dids <did[,did...]>] [--producer-only] [--max-runtime 300] [--max-messages 10000]"); } const maxRuntimeSeconds = positiveInteger(valueAfter("--max-runtime") ?? "300", "--max-runtime", 86_400); const maxMessages = positiveInteger(valueAfter("--max-messages") ?? "10000", "--max-messages", 1_000_000); const rewindUs = nonnegativeInteger(valueAfter("--rewind-us") ?? "2000000", "--rewind-us", 300_000_000); const maxReconnects = nonnegativeInteger(valueAfter("--max-reconnects") ?? "8", "--max-reconnects", 100); const connector = new JetstreamConnector({ id: source, collections, dids }); const producerOnly = process.argv.includes("--producer-only"); const declarations = producerOnly ? undefined : await loadAgentDeclarations(path.join(projectRoot, "agents")); const runtime = producerOnly ? undefined : await createAgentRuntime(); const runsBefore = producerOnly ? 0 : (await store.listRuns()).length; const consumers = runtime && declarations ? await runtime.startConsumers(declarations) : undefined; let consumersStopped = producerOnly; const controller = new AbortController(); const abort = () => controller.abort(); const timer = setTimeout(abort, maxRuntimeSeconds * 1_000); process.once("SIGINT", abort); process.once("SIGTERM", abort); try { const subscription = await subscribeJetstream(store, connector, { ...(process.env.THOUGHTSTREAM_JETSTREAM_URL ? { endpoint: process.env.THOUGHTSTREAM_JETSTREAM_URL } : {}), signal: controller.signal, rewindUs, maxMessages, maxReconnects, }); if (consumers && runtime && declarations) { // A bounded subscription can finish before a wildcard consumer's source-discovery // callback has attached to a source that did not exist at startup. Reconcile the // durable backlog after stopping subscriptions, avoiding concurrent transactions // while still preserving any work completed by the live consumer. await consumers.stop(); consumersStopped = true; await runtime.consumeBacklog(declarations); } print({ subscription, producerOnly, consumerRuns: producerOnly ? 0 : (await store.listRuns()).length - runsBefore, }); } finally { if (consumers && !consumersStopped) await consumers.stop(); clearTimeout(timer); process.removeListener("SIGINT", abort); process.removeListener("SIGTERM", abort); } } else if (command === "x-webhook") { const sourceConfig = await xSourceConfig(); const consumerSecret = process.env[sourceConfig.consumerSecretEnv]; if (!consumerSecret) throw new Error(`X consumer secret environment variable is unavailable: ${sourceConfig.consumerSecretEnv}`); const connector = new XActivityConnector({ id: sourceConfig.id, lane: sourceConfig.lane, expectedSubscriptions: sourceConfig.expectedSubscriptions, }); const receiver = await startXWebhookServer({ connector, store, consumerSecret, path: sourceConfig.webhookPath, host: sourceConfig.listenHost, port: sourceConfig.listenPort, maxBodyBytes: sourceConfig.maxBodyBytes, requestTimeoutMs: sourceConfig.requestTimeoutMs, pendingRequestLimit: sourceConfig.pendingRequestLimit, onCrcDiagnostic: (diagnostic) => print({ xWebhookCrc: { source: sourceConfig.id, ...diagnostic, }, }), }); try { await updateXSourceControl(store, sourceConfig, { runtime: { state: "ready", readyAt: new Date().toISOString(), revision: X_WEBHOOK_REVISION, }, }); } catch (error) { await receiver.close(); throw error; } print({ xWebhook: { source: sourceConfig.id, lane: sourceConfig.lane, listenHost: receiver.host, listenPort: receiver.port, path: receiver.path, webhookUrlSha256: sha256(sourceConfig.webhookUrl), expectedSubscriptionCount: sourceConfig.expectedSubscriptions.length, }, }); try { const maxRuntime = valueAfter("--max-runtime"); await waitForSignalOrTimeout(maxRuntime ? positiveInteger(maxRuntime, "--max-runtime", 86_400) * 1_000 : undefined); } finally { await receiver.close(); } process.stderr.write("thought stream X webhook stopped.\n"); } else if (command === "x-webhook-status") { const sourceConfig = await xSourceConfig(); const client = xManagementClient(sourceConfig); const [webhooks, subscriptions] = await Promise.all([ client.listWebhooks(), client.listSubscriptions(), ]); const webhook = exactXWebhook(sourceConfig, webhooks, false); const plan = webhook ? planXSubscriptions(sourceConfig.id, webhook.id, sourceConfig.expectedSubscriptions, subscriptions) : undefined; await updateXSourceControl(store, sourceConfig, { ...(process.argv.includes("--record-runtime-ready") ? { runtime: { state: "ready" as const, readyAt: new Date().toISOString(), revision: X_WEBHOOK_REVISION, }, } : {}), upstream: xUpstreamControl(sourceConfig, webhook, plan), }); print({ xWebhookStatus: { source: sourceConfig.id, lane: sourceConfig.lane, webhookUrlSha256: sha256(sourceConfig.webhookUrl), webhook: webhook ? { id: webhook.id, valid: webhook.valid } : null, subscriptionCount: subscriptions.length, desiredSubscriptionCount: sourceConfig.expectedSubscriptions.length, plan: plan ? publicXPlan(plan) : null, }, }); } else if (command === "x-webhook-register") { const sourceConfig = await xSourceConfig(); requireExactFlag("--confirm-url-hash", sha256(sourceConfig.webhookUrl)); const client = xManagementClient(sourceConfig); const before = exactXWebhook(sourceConfig, await client.listWebhooks(), false); let status: "created" | "validated" | "unchanged"; let webhook = before; if (!webhook) { webhook = await client.createWebhook(sourceConfig.webhookUrl); status = "created"; } else if (!webhook.valid) { if (!await client.validateWebhook(webhook.id)) throw new Error("X webhook validation returned false"); status = "validated"; } else { status = "unchanged"; } const verified = exactXWebhook(sourceConfig, await client.listWebhooks(), true); if (!verified || verified.id !== webhook.id) throw new Error("X webhook registration verification failed"); const subscriptions = await client.listSubscriptions(); const plan = planXSubscriptions(sourceConfig.id, verified.id, sourceConfig.expectedSubscriptions, subscriptions); await updateXSourceControl(store, sourceConfig, { upstream: xUpstreamControl(sourceConfig, verified, plan), }); print({ xWebhookRegistration: { source: sourceConfig.id, status, webhookId: verified.id, valid: verified.valid } }); } else if (command === "x-subscriptions-plan") { const sourceConfig = await xSourceConfig(); const client = xManagementClient(sourceConfig); const webhook = exactXWebhook(sourceConfig, await client.listWebhooks(), true); if (!webhook) throw new Error("No registered X webhook matches the source URL"); const subscriptions = await client.listSubscriptions(); const plan = planXSubscriptions( sourceConfig.id, webhook.id, sourceConfig.expectedSubscriptions, subscriptions, ); await updateXSourceControl(store, sourceConfig, { upstream: xUpstreamControl(sourceConfig, webhook, plan), }); print({ xSubscriptionPlan: publicXPlan(plan) }); } else if (command === "x-subscriptions-apply") { const sourceConfig = await xSourceConfig(); const client = xManagementClient(sourceConfig); const webhook = exactXWebhook(sourceConfig, await client.listWebhooks(), true); if (!webhook) throw new Error("No registered valid X webhook matches the source URL"); const plan = planXSubscriptions( sourceConfig.id, webhook.id, sourceConfig.expectedSubscriptions, await client.listSubscriptions(), ); requireExactFlag("--confirm-plan-hash", plan.planHash); if (plan.conflicts.length > 0) throw new Error(`X subscription plan has conflicts: ${plan.conflicts.join(", ")}`); for (const update of plan.update) await client.updateSubscription(update.subscriptionId, update); for (const deletion of plan.delete) { if (!await client.deleteSubscription(deletion.subscriptionId)) { throw new Error(`X subscription deletion returned false: ${deletion.subscriptionId}`); } } for (const creation of plan.create) await client.createSubscription(creation); const liveSubscriptions = await client.listSubscriptions(); const verification = planXSubscriptions( sourceConfig.id, webhook.id, sourceConfig.expectedSubscriptions, liveSubscriptions, ); if (verification.create.length > 0 || verification.update.length > 0 || verification.delete.length > 0 || verification.conflicts.length > 0) { throw new Error("X subscription apply did not converge to the desired state"); } await updateXSourceControl(store, sourceConfig, { upstream: xUpstreamControl(sourceConfig, webhook, verification), }); print({ xSubscriptionApply: { source: sourceConfig.id, planHash: plan.planHash, created: plan.create.length, updated: plan.update.length, deleted: plan.delete.length, verifiedUnchanged: verification.unchanged.length, }, }); } else if (command === "x-webhook-replay") { const sourceConfig = await xSourceConfig(); const fromDate = valueAfter("--from"); const toDate = valueAfter("--to"); if (!fromDate || !toDate) { throw new Error("Usage: thought stream x-webhook-replay --source <id> --from YYYYMMDDHHmm --to YYYYMMDDHHmm --confirm-webhook-id <id>"); } assertBoundedXReplayWindow(fromDate, toDate); const client = xManagementClient(sourceConfig); const webhook = exactXWebhook(sourceConfig, await client.listWebhooks(), true); if (!webhook) throw new Error("No registered valid X webhook matches the source URL"); requireExactFlag("--confirm-webhook-id", webhook.id); const replay = await client.createReplay(webhook.id, fromDate, toDate); print({ xWebhookReplay: { source: sourceConfig.id, webhookId: webhook.id, fromDate, toDate, ...replay } }); } else if (command === "x-webhook-delete") { const sourceConfig = await xSourceConfig(); const client = xManagementClient(sourceConfig); const webhook = exactXWebhook(sourceConfig, await client.listWebhooks(), false); if (!webhook) throw new Error("No registered X webhook matches the source URL"); requireExactFlag("--confirm-webhook-id", webhook.id); const linkedSubscriptions = (await client.listSubscriptions()).filter((subscription) => subscription.webhookId === webhook.id); if (linkedSubscriptions.length > 0) { throw new Error(`Refusing to delete X webhook while ${linkedSubscriptions.length} linked subscriptions remain`); } if (!await client.deleteWebhook(webhook.id)) throw new Error("X webhook deletion returned false"); if (exactXWebhook(sourceConfig, await client.listWebhooks(), false)) { throw new Error("X webhook deletion could not be verified"); } await updateXSourceControl(store, sourceConfig, { upstream: xUpstreamControl(sourceConfig, undefined, undefined), }); print({ xWebhookRegistration: { source: sourceConfig.id, status: "deleted", webhookId: webhook.id } }); } else if (command === "x-user-lookup") { const sourceConfig = await xSourceConfig(); const username = valueAfter("--username"); if (!username) throw new Error("Usage: thought stream x-user-lookup --source <id> --username <handle>"); print({ xUser: await xManagementClient(sourceConfig).lookupUser(username) }); } else if (command === "telegram-spool") { const file = valueAfter("--file"); const source = valueAfter("--source"); if (!file || !source) { throw new Error("Usage: thought stream telegram-spool --file <ndjson-file> --source telegram:name [--max-records 1000] [--max-read-bytes 8388608]"); } const maxRecords = positiveInteger(valueAfter("--max-records") ?? "1000", "--max-records", 100_000); const maxReadBytes = positiveInteger(valueAfter("--max-read-bytes") ?? "8388608", "--max-read-bytes", 64 * 1024 * 1024); const connector = new TelegramSpoolConnector({ id: source, file, maxRecords, maxReadBytes }); const ingest = await connector.ingest(store); print({ ingest }); } else if (command === "telegram-webhook") { const sourceConfig = await telegramSourceConfig(); const channels = sourceConfig.channels.filter((channel) => channel.enabled); const chatIds = channels.map((channel) => channel.id); const token = process.env[sourceConfig.tokenEnv]; if (!token) throw new Error(`Telegram bot credential environment variable is unavailable: ${sourceConfig.tokenEnv}`); const secretToken = process.env[sourceConfig.webhookSecretEnv]; if (!secretToken) throw new Error(`Telegram webhook secret environment variable is unavailable: ${sourceConfig.webhookSecretEnv}`); const source = sourceConfig.id; const client = new TelegramBotClient({ token, ...(process.env.THOUGHTSTREAM_TELEGRAM_API_BASE_URL ? { baseUrl: process.env.THOUGHTSTREAM_TELEGRAM_API_BASE_URL } : {}), requestTimeoutMs: sourceConfig.requestTimeoutMs, }); const connector = new TelegramBotConnector({ id: source, bot: await client.identity(), allowedChatIds: chatIds, reactionFeedback: channels .filter((channel) => channel.reactionFeedback?.enabled) .map((channel) => ({ chatId: channel.id, allowedUserIds: channel.reactionFeedback!.allowedUserIds, })), imageExtraction: { client, artifactRoot: store.getArtifactRoot(), }, }); const feedbackCatchup = await projectTelegramFeedbackJudgments(store, source); if (feedbackCatchup.projected > 0 || feedbackCatchup.retracted > 0) { print({ telegramFeedbackCatchup: { source, examined: feedbackCatchup.examined, projected: feedbackCatchup.projected, retracted: feedbackCatchup.retracted, skipped: feedbackCatchup.skipped, }, }); } const receiver = await startTelegramWebhookServer({ connector, store, secretToken, path: sourceConfig.webhookPath, host: sourceConfig.listenHost, port: sourceConfig.listenPort, maxBodyBytes: sourceConfig.maxBodyBytes, }); process.stdout.write(`thought stream Telegram webhook listening on http://${receiver.host}:${receiver.port}${receiver.path}\n`); try { const maxRuntime = valueAfter("--max-runtime"); await waitForSignalOrTimeout(maxRuntime ? positiveInteger(maxRuntime, "--max-runtime", 86_400) * 1_000 : undefined); } finally { await receiver.close(); } process.stderr.write("thought stream Telegram webhook stopped.\n"); } else if (command === "telegram-webhook-register") { const sourceConfig = await telegramSourceConfig(); const client = telegramClient(sourceConfig); const secretToken = process.env[sourceConfig.webhookSecretEnv]; if (!secretToken) throw new Error(`Telegram webhook secret environment variable is unavailable: ${sourceConfig.webhookSecretEnv}`); await client.setWebhook({ url: sourceConfig.webhookUrl, secretToken, dropPendingUpdates: process.argv.includes("--drop-pending-updates"), }); const info = await client.getWebhookInfo(); assertWebhookRegistration(sourceConfig, info); print({ telegramWebhookRegistration: { source: sourceConfig.id, status: "registered", webhookUrlHash: sha256(sourceConfig.webhookUrl), pendingUpdateCount: info.pending_update_count, maxConnections: info.max_connections, allowedUpdates: info.allowed_updates ?? [], }, }); } else if (command === "telegram-webhook-delete") { const sourceConfig = await telegramSourceConfig(); const client = telegramClient(sourceConfig); await client.deleteWebhook({ dropPendingUpdates: process.argv.includes("--drop-pending-updates") }); const info = await client.getWebhookInfo(); if (info.url) throw new Error("Telegram webhook deletion could not be verified"); print({ telegramWebhookRegistration: { source: sourceConfig.id, status: "deleted" } }); } else if (command === "telegram-menu-register") { const sourceConfig = await telegramSourceConfig(); const client = telegramClient(sourceConfig); await client.setName(STREAM_TELEGRAM_BOT_NAME); if (await client.getName() !== STREAM_TELEGRAM_BOT_NAME) { throw new Error("Telegram bot display name verification failed"); } await client.setCommands([...STREAM_TELEGRAM_COMMANDS]); const commands = await client.getCommands(); if (canonicalJson(commands as unknown as JsonObject[]) !== canonicalJson([...STREAM_TELEGRAM_COMMANDS] as unknown as JsonObject[])) { throw new Error("Telegram bot command menu verification failed"); } const sendHelpNotice = !process.argv.includes("--skip-help-notice"); const deliveries = []; for (const channel of sourceConfig.channels.filter((candidate) => candidate.enabled)) { await client.setCommandsMenuButton(channel.id); if (await client.getMenuButton(channel.id) !== "commands") { throw new Error("Telegram chat menu button verification failed"); } if (!sendHelpNotice) continue; const notification = channel.notifications; const dispatcher = new TelegramChannelDispatcher({ id: `telegram-dispatcher:${sourceConfig.id}:${channel.id}`, client, chatId: channel.id, allowedSources: notification?.allowedSources ?? ["system:telegram-menu"], ...(notification?.allowedActors ? { allowedActors: notification.allowedActors } : {}), ...(notification?.directReplyAgentIds ? { directReplyAgentIds: notification.directReplyAgentIds } : {}), ...(notification?.directReplySources ? { directReplySources: notification.directReplySources } : {}), ...(notification?.notificationProposalAgentIds ? { notificationProposalAgentIds: notification.notificationProposalAgentIds } : {}), ...(notification?.notificationProposalSources ? { notificationProposalSources: notification.notificationProposalSources } : {}), runStatuses: notification?.runStatuses ?? ["completed"], maxMessagesPerWindow: notification?.maxMessagesPerWindow ?? 3, windowMs: notification?.windowMs ?? 60_000, maxLikesPerDigest: notification?.maxLikesPerDigest ?? 10, likeDigestDelayMs: notification?.likeDigestDelayMs ?? 60_000, }); deliveries.push({ channelId: channel.id, ...(await dispatcher.sendOperationalNotice(store, STREAM_TELEGRAM_HELP_TEXT, { scope: STREAM_TELEGRAM_HELP_VERSION, })), }); } print({ telegramMenu: { source: sourceConfig.id, version: STREAM_TELEGRAM_HELP_VERSION, botName: STREAM_TELEGRAM_BOT_NAME, commands: commands.map((entry) => entry.command), menuButton: "commands", helpNotice: sendHelpNotice ? "sent" : "skipped", deliveries, }, }); } else if (command === "telegram-dispatcher") { const sourceConfig = await telegramSourceConfig(); const client = telegramClient(sourceConfig); const startedAt = new Date().toISOString(); const channels = sourceConfig.channels.filter((channel) => channel.enabled); const channelRuntimes = channels.map((channel) => { const notification = channel.notifications; const dispatcher = channel.bootMessage?.enabled || notification?.enabled ? new TelegramChannelDispatcher({ id: `telegram-dispatcher:${sourceConfig.id}:${channel.id}`, client, chatId: channel.id, allowedSources: notification?.allowedSources ?? ["system:boot"], ...(notification?.allowedActors ? { allowedActors: notification.allowedActors } : {}), ...(notification?.directReplyAgentIds ? { directReplyAgentIds: notification.directReplyAgentIds } : {}), ...(notification?.directReplySources ? { directReplySources: notification.directReplySources } : {}), ...(notification?.notificationProposalAgentIds ? { notificationProposalAgentIds: notification.notificationProposalAgentIds } : {}), ...(notification?.notificationProposalSources ? { notificationProposalSources: notification.notificationProposalSources } : {}), runStatuses: notification?.runStatuses ?? ["completed"], maxMessagesPerWindow: notification?.maxMessagesPerWindow ?? 3, windowMs: notification?.windowMs ?? 60_000, maxLikesPerDigest: notification?.maxLikesPerDigest ?? 10, likeDigestDelayMs: notification?.likeDigestDelayMs ?? 60_000, }) : undefined; return { channel, notification, dispatcher }; }); const activeSince = new Map<string, string>(); for (const runtime of channelRuntimes) { if (!runtime.dispatcher || !runtime.notification?.enabled) continue; activeSince.set(runtime.channel.id, await runtime.dispatcher.activate(store)); } const maxRuntimeSeconds = positiveInteger(valueAfter("--max-runtime") ?? "300", "--max-runtime", 86_400); const deadline = Date.now() + maxRuntimeSeconds * 1_000; let stopping = false; const stop = () => { stopping = true; }; process.once("SIGINT", stop); process.once("SIGTERM", stop); let cycles = 0; let delivered = 0; let failed = 0; const bootedChannels = new Set(channelRuntimes .filter((runtime) => !runtime.channel.bootMessage?.enabled) .map((runtime) => runtime.channel.id)); try { do { cycles += 1; const bootMessages = []; for (const runtime of channelRuntimes) { if (!runtime.dispatcher || !runtime.channel.bootMessage?.enabled || bootedChannels.has(runtime.channel.id)) continue; const result = await runtime.dispatcher.sendOperationalNotice(store, runtime.channel.bootMessage.text, { scope: startedAt, }); bootMessages.push({ channelId: runtime.channel.id, ...result }); if (result.status === "delivered") delivered += 1; if (result.status === "failed") failed += 1; if (result.status !== "rate-limited") bootedChannels.add(runtime.channel.id); } const notifications = []; for (const runtime of channelRuntimes) { if (!runtime.dispatcher || !runtime.notification?.enabled || !bootedChannels.has(runtime.channel.id)) continue; const [, result] = await Promise.all([ runtime.dispatcher.refreshTyping(store), runtime.dispatcher.sendPending(store, { ...(process.argv.includes("--notify-existing") ? {} : { since: activeSince.get(runtime.channel.id)! }), includeNormal: runtime.notification.includeNormal, }), ]); notifications.push({ channelId: runtime.channel.id, ...result }); delivered += result.delivered; failed += result.failed; } if (bootMessages.length > 0 || notifications.some((result) => result.delivered || result.failed)) { print({ ...(bootMessages.length > 0 ? { bootMessages } : {}), ...(notifications.length > 0 ? { notifications } : {}) }); } if (!stopping && !process.argv.includes("--once") && Date.now() < deadline) { await delay(sourceConfig.dispatchIntervalMs); } } while (!stopping && !process.argv.includes("--once") && Date.now() < deadline); print({ telegramDispatcher: { source: sourceConfig.id, cycles, notifications: { delivered, failed }, stopped: stopping ? "signal" : process.argv.includes("--once") ? "once" : "runtime-limit", }, }); } finally { process.removeListener("SIGINT", stop); process.removeListener("SIGTERM", stop); } } else if (command === "fastmail-jmap") { const sourceConfig = await fastmailSourceConfig(); const token = process.env[sourceConfig.tokenEnv]; if (!token) throw new Error(`Fastmail credential environment variable is unavailable: ${sourceConfig.tokenEnv}`); if (sourceConfig.credentialCustody === "unprovisioned") { throw new Error(`Fastmail credential custody is unprovisioned: ${sourceConfig.id}`); } const connector = new FastmailJmapConnector({ id: sourceConfig.id, client: new FastmailJmapClient({ token, timeoutMs: sourceConfig.timeoutMs, maxResponseBytes: sourceConfig.maxResponseBytes, }), maxChanges: sourceConfig.maxChanges, maxPages: sourceConfig.maxPages, resnapshotLimit: sourceConfig.resnapshotLimit, credentialCustody: sourceConfig.credentialCustody, }); const once = process.argv.includes("--once"); const controller = new AbortController(); const stop = () => controller.abort(); const maxRuntime = valueAfter("--max-runtime"); const timer = maxRuntime ? setTimeout(stop, positiveInteger(maxRuntime, "--max-runtime", 86_400) * 1_000) : undefined; let cycles = 0; let failures = 0; let consecutiveFailures = 0; process.once("SIGINT", stop); process.once("SIGTERM", stop); try { if (!sourceConfig.pollOnStart && !once) await abortableDelay(sourceConfig.intervalMs, controller.signal); while (!controller.signal.aborted) { try { const poll = await connector.poll(store, controller.signal); cycles += 1; consecutiveFailures = 0; if (poll.initialized || poll.resnapshot || poll.inserted > 0 || once) { print({ fastmailJmap: { source: poll.source, initialized: poll.initialized, resnapshot: poll.resnapshot, inserted: poll.inserted, unchanged: poll.unchanged, created: poll.created, updated: poll.updated, destroyed: poll.destroyed, }, }); } } catch (error) { if (controller.signal.aborted) break; failures += 1; consecutiveFailures += 1; process.stderr.write(`thought stream Fastmail poll failed: ${fastmailJmapErrorCode(error)}\n`); if (once) throw error; } if (once) break; await abortableDelay(fastmailRetryDelay(sourceConfig.intervalMs, consecutiveFailures), controller.signal); } } finally { if (timer) clearTimeout(timer); process.removeListener("SIGINT", stop); process.removeListener("SIGTERM", stop); } if (!once) print({ fastmailJmap: { source: sourceConfig.id, cycles, failures, stopped: controller.signal.aborted } }); } else if (command === "fastmail-capture") { const file = valueAfter("--file"); const source = valueAfter("--source"); const accountId = valueAfter("--account-id"); if (!file || !source || !accountId) { throw new Error("Usage: thought stream fastmail-capture --file <jmap-response.json> --source fastmail:name --account-id <jmap-account-id>"); } const raw = JSON.parse(await fs.readFile(path.resolve(file), "utf8")) as unknown; const connector = new FastmailConnector({ id: source, accountId }); const ingest = await connector.ingestCapture(store, raw); print({ ingest }); } else if (command === "consume") { const declarations = await loadAgentDeclarations(path.join(projectRoot, "agents")); const runtime = await createAgentRuntime(); if (process.argv.includes("--once")) { print({ runs: await runtime.consumeBacklog(declarations) }); } else { const consumers = await runtime.startConsumers(declarations); process.stdout.write("thought stream consumers following Jazz subscriptions.\n"); try { await waitForSignal(); } finally { await consumers.stop(); } process.stderr.write("thought stream consumers stopped.\n"); } } else if (command === "incidents") { const loaded = await loadThoughtStreamManifest(projectRoot, valueAfter("--config") ?? "thoughtstream.yaml"); const incidentConfig = loaded.manifest.incidents; if (!incidentConfig.enabled) throw new Error("Operational incidents are disabled in the thought stream manifest"); const ledger = await IncidentLedger.open(projectRoot, incidentConfig.ledgerPath); const projector = new OperationalIncidentProjector(); const alerts = incidentConfig.telegramAlerts; const dispatchers: Array<{ dispatcher: IncidentTelegramDispatcher; since: string }> = []; if (alerts.enabled) { const source = loaded.manifest.sources.find((candidate) => candidate.id === alerts.sourceId); if (!source || source.kind !== "telegram-webhook" || !source.enabled) { throw new Error("Operational incident alerts require an enabled telegram-webhook source"); } const client = telegramClient(source); for (const chatId of alerts.channelIds) { const dispatcher = new IncidentTelegramDispatcher({ id: `telegram-incident-dispatcher:${source.id}:${chatId}`, client, chatId, categories: alerts.categories, cooldownMs: alerts.cooldownMs, maxMessagesPerWindow: alerts.maxMessagesPerWindow, windowMs: alerts.windowMs, }); dispatchers.push({ dispatcher, since: await dispatcher.activate(store) }); } } const maxRuntimeSeconds = positiveInteger(valueAfter("--max-runtime") ?? "86400", "--max-runtime", 86_400); const deadline = Date.now() + maxRuntimeSeconds * 1_000; let stopping = false; const stop = () => { stopping = true; }; process.once("SIGINT", stop); process.once("SIGTERM", stop); let cycles = 0; let projectedTotal = 0; let ledgerTotal = 0; let deliveredTotal = 0; let failedTotal = 0; try { do { cycles += 1; const projection = await projector.project(store); projectedTotal += projection.inserted; const incidents = await listOperationalIncidents(store); const ledgerResult = await ledger.append(incidents.map((value) => value.incident)); ledgerTotal += ledgerResult.appended; const alertResults = []; for (const runtime of dispatchers) { const result = await runtime.dispatcher.sendPending(store, { since: runtime.since }); alertResults.push({ dispatcher: runtime.dispatcher.id, ...result }); deliveredTotal += result.delivered; failedTotal += result.failed; } if (projection.inserted > 0 || ledgerResult.appended > 0 || alertResults.some((result) => result.delivered || result.failed)) { print({ operationalIncidents: { projection, ledger: ledgerResult, alerts: alertResults } }); } if (!stopping && !process.argv.includes("--once") && Date.now() < deadline) { await delay(incidentConfig.intervalMs); } } while (!stopping && !process.argv.includes("--once") && Date.now() < deadline); print({ incidentService: { cycles, projected: projectedTotal, ledgerAppended: ledgerTotal, alerts: { delivered: deliveredTotal, failed: failedTotal }, stopped: stopping ? "signal" : process.argv.includes("--once") ? "once" : "runtime-limit", }, }); } finally { process.removeListener("SIGINT", stop); process.removeListener("SIGTERM", stop); } } else if (command === "events") { const limit = Number(valueAfter("--limit") ?? "100"); print(await store.listEvents({ limit })); } else if (command === "event") { const id = process.argv[3]; if (!id) throw new Error("Usage: thought stream event <event-id>"); const event = await store.getEvent(id); if (!event) throw new Error(`Event not found: ${id}`); print(event); } else if (command === "runs") { print(await store.listRuns()); } else if (command === "run") { const id = process.argv[3]; if (!id) throw new Error("Usage: thought stream run <run-id>"); const run = await store.getRun(id); if (!run) throw new Error(`Run not found: ${id}`); print({ run, trace: await store.listTrace(id) }); } else if (command === "context") { const id = process.argv[3]; const includeContent = process.argv.includes("--show-content"); const acknowledged = process.argv.includes("--acknowledge-sensitive-private"); if (!id) { throw new Error("Usage: thought stream context <run-id> [--show-content --acknowledge-sensitive-private]"); } if (includeContent && !acknowledged) { throw new Error("--show-content requires --acknowledge-sensitive-private"); } print({ context: await inspectRunContext(store, id, { includeContent }) }); } else if (command === "agent-send") { const from = valueAfter("--from"); const to = valueAfter("--to"); const threadId = valueAfter("--thread"); const messageId = valueAfter("--message-id"); const file = valueAfter("--file"); const waitSecondsValue = valueAfter("--wait-seconds"); const showResponse = process.argv.includes("--show-response"); const acknowledged = process.argv.includes("--acknowledge-sensitive-private"); if (from !== CO_AGENT_ID || to !== STREAM_AGENT_ENDPOINT_ID || !threadId || !messageId || !file) { throw new Error("Usage: thought stream agent-send --from co --to stream-agent-conversation --thread <id> --message-id <id> --file <text-file> [--wait-seconds <1-600>] [--show-response --acknowledge-sensitive-private]"); } if (showResponse && !acknowledged) { throw new Error("--show-response requires --acknowledge-sensitive-private"); } if (showResponse && waitSecondsValue === undefined) { throw new Error("--show-response requires --wait-seconds"); } const messagePath = path.resolve(file); const messageStat = await fs.lstat(messagePath); if (!messageStat.isFile() || messageStat.isSymbolicLink()) { throw new Error("--file must be a regular non-symlink file"); } if (messageStat.size > AGENT_MESSAGE_TEXT_MAX_CHARS * 4) { throw new Error("--file is too large for one agent message"); } const text = await fs.readFile(messagePath, "utf8"); const appended = await appendCoAgentMessage(store, { messageId, threadId, text }); const receipt: JsonObject = { eventId: appended.event.id, inserted: appended.inserted, sourceSequence: appended.event.sourceSequence, messageId, threadId, senderAgentId: CO_AGENT_ID, recipientAgentId: STREAM_AGENT_ENDPOINT_ID, }; if (waitSecondsValue === undefined) { print({ agentMessage: receipt }); } else { const waitSeconds = positiveInteger(waitSecondsValue, "--wait-seconds", 600); const deadline = Date.now() + waitSeconds * 1_000; let terminal: AgentRun | undefined; while (Date.now() <= deadline) { const runs = (await store.getRunsForTriggerEvents([appended.event.id])) .filter((run) => run.agentId === STREAM_AGENT_ENDPOINT_ID) .sort((left, right) => right.attempt - left.attempt || right.createdAt.localeCompare(left.createdAt)); const latest = runs[0]; terminal = latest && ["completed", "failed", "blocked", "abandoned", "skipped"].includes(latest.status) && !agentRunRetryPending(latest) ? latest : undefined; if (terminal) break; await delay(250); } if (!terminal) { print({ agentMessage: receipt, run: { status: "pending", waitedSeconds: waitSeconds } }); process.exitCode = 2; } else { const runReceipt: JsonObject = { id: terminal.id, status: terminal.status, attempt: terminal.attempt, agentId: terminal.agentId, agentVersion: terminal.agentVersion, outputEventIds: terminal.outputEventIds, }; const output = terminal.status === "completed" && terminal.outputEventIds.length === 1 ? await store.getEvent(terminal.outputEventIds[0]!) : undefined; const response = output ? parseAgentMessageResponseEvent(output) : undefined; if (terminal.status === "completed" && (!output || !response)) { throw new Error("Completed agent-message run has no valid typed response event"); } print({ agentMessage: receipt, run: runReceipt, ...(output && response ? { response: { eventId: output.id, messageId: response.messageId, inReplyToMessageId: response.inReplyToMessageId, threadId: response.threadId, senderAgentId: response.senderAgentId, recipientAgentId: response.recipientAgentId, ...(showResponse ? { text: response.summary } : { chars: response.summary.length }), }, } : {}), }); if (terminal.status !== "completed") process.exitCode = 1; } } } else if (command === "review-prompt") { const file = valueAfter("--file"); const externalId = valueAfter("--external-id"); const privacy = valueAfter("--privacy") ?? "private"; if (!file || !externalId || !["public-source", "private", "sensitive"].includes(privacy)) { throw new Error("Usage: thought stream review-prompt --file <payload.json> --external-id <stable-id> [--privacy <public-source|private|sensitive>] [--source review:prompt]"); } const prompt = await appendReviewPrompt(store, { source: valueAfter("--source") ?? "review:prompt", externalId, privacy: privacy as PrivacyClass, payload: jsonObject(JSON.parse(await fs.readFile(path.resolve(file), "utf8")), "--file"), }); print({ reviewPrompt: prompt }); } else if (command === "review-item") { const promptEventId = valueAfter("--prompt-event"); const candidateRunIds = commaSeparated(valueAfter("--candidate-runs")); if (!promptEventId || candidateRunIds.length !== 2) { throw new Error("Usage: thought stream review-item --prompt-event <event-id> --candidate-runs <run-a>,<run-b>"); } const item = await createReviewItem(store, { promptEventId, candidateRunIds: [candidateRunIds[0]!, candidateRunIds[1]!], }); print({ reviewItem: item }); } else if (command === "review-queue") { print(await projectReviewQueue(store)); } else if (command === "focus-propose") { const file = valueAfter("--file"); if (!file) throw new Error("Usage: thought stream focus-propose --file <focus.yaml>"); const declaration = await loadFocusDeclaration(file); const appended = await proposeFocusDeclaration(store, declaration); print({ focusProposal: { eventId: appended.event.id, focusId: declaration.id, focusVersion: declaration.version, declarationFingerprint: String(appended.event.payload.declarationFingerprint), inserted: appended.inserted, authority: "proposal-only", }, }); } else if (command === "focus-list") { print({ focusProposals: await listFocusProposals(store) }); } else if (command === "proposal-list") { const proposals = (await store.listEvents({ types: [MEMORY_PROPOSAL_EVENT_TYPE, CORRECTION_PROPOSAL_EVENT_TYPE], })).sort((left, right) => left.observedAt.localeCompare(right.observedAt) || left.id.localeCompare(right.id)); const rows = []; for (const proposal of proposals) { const decisions = await decisionsForProposal(store, proposal.id); const decision = decisions.at(-1); const materialized = decision ? (await store.listEvents({ types: [MEMORY_MATERIALIZED_EVENT_TYPE] })).find((event) => event.payload.decisionEventId === decision.id) : undefined; rows.push({ proposalEventId: proposal.id, proposalType: proposal.type === MEMORY_PROPOSAL_EVENT_TYPE ? "memory-change" : "self-correction", observedAt: proposal.observedAt, decisionEventId: decision?.id ?? null, disposition: decision?.payload.disposition ?? null, materializedEventId: materialized?.id ?? null, }); } print({ proposals: rows }); } else if (command === "proposal-decision") { const proposalEventId = process.argv[3]; const disposition = valueAfter("--disposition") as ProposalDecisionDisposition | undefined; if (!proposalEventId || !disposition || !["accept", "edit", "reject"].includes(disposition)) { throw new Error("Usage: thought stream proposal-decision <proposal-event-id> --disposition <accept|edit|reject> [--replacement-file <private-text-file>] [--context-root <configured-stream-context-root>] [--submission-id <stable-id>]"); } const replacementFile = valueAfter("--replacement-file"); if ((disposition === "edit") !== Boolean(replacementFile)) { throw new Error("Edit decisions require exactly one --replacement-file; accept/reject decisions forbid it"); } const proposal = await requireProposal(store, proposalEventId); const contextRoot = valueAfter("--context-root"); if (proposal.type === MEMORY_PROPOSAL_EVENT_TYPE && disposition !== "reject" && !contextRoot) { throw new Error("Accepted or edited memory decisions require --context-root for exact materialization"); } const replacementText = replacementFile ? await fs.readFile(path.resolve(replacementFile), "utf8") : undefined; const result = await recordProposalDecision(store, { proposalEventId, disposition, ...(replacementText !== undefined ? { replacementText } : {}), ...(valueAfter("--submission-id") ? { submissionId: valueAfter("--submission-id")! } : {}), actor: "operator:local", }); const materialization = proposal.type === MEMORY_PROPOSAL_EVENT_TYPE && disposition !== "reject" ? await materializeMemoryDecision(store, result.decision.id, { contextRoot: path.resolve(contextRoot!) }) : undefined; print({ proposalDecision: { proposalEventId, decisionEventId: result.decision.id, disposition, projectedJudgmentEventId: result.projectedJudgment?.id ?? null, materializationStatus: materialization?.status ?? null, materializationEventId: materialization?.event.id ?? null, }, }); } else if (command === "proposal-project") { const contextRoot = valueAfter("--context-root"); const correction = await projectAcceptedCorrectionDecisions(store); const materialized: Array<{ decisionEventId: string; status: string; eventId: string }> = []; for (const decision of await store.listEvents({ types: [PROPOSAL_DECISION_EVENT_TYPE] })) { if (decision.payload.proposalType !== "memory-change" || decision.payload.disposition === "reject") continue; if (!contextRoot) throw new Error("Memory proposal projection requires --context-root <configured-stream-context-root>"); const result = await materializeMemoryDecision(store, decision.id, { contextRoot: path.resolve(contextRoot) }); materialized.push({ decisionEventId: decision.id, status: result.status, eventId: result.event.id }); } print({ proposalProjection: { correction, materialized } }); } else if (command === "judgment") { const runId = process.argv[3]; const kind = valueAfter("--kind") as JudgmentKind | undefined; if (!runId || !kind || !["accept", "reject", "correct", "prefer"].includes(kind)) { throw new Error("Usage: thought stream judgment <run-id> --kind <accept|reject|correct|prefer> [--criterion output-quality] [--criterion-version 1] [--compared-run <run-id>] [--replacement <json-file>] [--external-export-eligible] [--authorize-sensitive-external-export] [--notes <text>]"); } if (process.argv.includes("--export-eligible")) { throw new Error("--export-eligible is ambiguous and retired; use --external-export-eligible with the sensitive authorization flag when required"); } const externalExportEligible = process.argv.includes("--external-export-eligible"); const sensitiveExternalExportAuthorized = process.argv.includes("--authorize-sensitive-external-export"); if (sensitiveExternalExportAuthorized && !externalExportEligible) { throw new Error("--authorize-sensitive-external-export requires --external-export-eligible"); } const replacementPath = valueAfter("--replacement"); const replacementOutput = replacementPath ? jsonObject(JSON.parse(await fs.readFile(path.resolve(replacementPath), "utf8")), "--replacement") : undefined; const judgment = await recordJudgment(store, { runId, kind, criterion: valueAfter("--criterion") ?? "output-quality", criterionVersion: positiveInteger(valueAfter("--criterion-version") ?? "1", "--criterion-version", 1_000_000), qualityEligible: true, externalExportEligible, sensitiveExternalExportAuthorized, ...(valueAfter("--compared-run") ? { comparedRunId: valueAfter("--compared-run")! } : {}), ...(replacementOutput ? { replacementOutput } : {}), ...(valueAfter("--notes") ? { notes: valueAfter("--notes")! } : {}), }); print({ judgment }); } else if (command === "private-training-export") { const output = valueAfter("--output"); if (!output || !process.argv.includes("--acknowledge-sensitive-private-training")) { throw new Error("Usage: thought stream private-training-export --output <private-file.jsonl> --acknowledge-sensitive-private-training [--include-restricted-model-adapters]"); } const manifest = await exportPrivateTrainingDataset(store, path.resolve(output), { acknowledgeSensitivePrivate: true, ...(process.argv.includes("--include-restricted-model-adapters") ? { includeRestrictedModelAdapters: true } : {}), }); print({ privateTrainingExport: { output: path.resolve(output), manifest: `${path.resolve(output)}.manifest.json`, ...manifest } }); } else if (command === "training-export") { const includeSensitivePrivate = process.argv.includes("--include-sensitive-private"); const authorizeSensitivePrivateExport = process.argv.includes("--authorize-sensitive-private-export"); if (includeSensitivePrivate !== authorizeSensitivePrivateExport) { throw new Error("Sensitive/private export requires both --include-sensitive-private and --authorize-sensitive-private-export"); } const output = valueAfter("--output"); if (includeSensitivePrivate && !output) { throw new Error("Sensitive/private export requires --output and is never written to stdout"); } const examples = await projectTrainingExamples(store, { includeSensitivePrivate }); if (output) { const manifest = await writeTrainingJsonl(output, examples, { authorizeSensitivePrivateExport }); print({ trainingExport: { output: path.resolve(output), manifest: `${path.resolve(output)}.manifest.json`, ...manifest } }); } else { for (const example of examples) process.stdout.write(`${JSON.stringify(example)}\n`); } } else if (command === "exe-artifact-transport") { const host = valueAfter("--host"); const workspaceRoot = valueAfter("--workspace-root"); const workspaceId = valueAfter("--workspace-id"); const leaseId = valueAfter("--lease-id"); const stagingRoot = valueAfter("--staging-root"); if (host !== "cameron.exe.xyz" || workspaceRoot !== "/srv/thoughtstream/workspaces/stream-v1" || workspaceId !== "stream-v1" || !leaseId || !stagingRoot) { throw new Error("Usage: thought exe-artifact-transport --host cameron.exe.xyz --workspace-root /srv/thoughtstream/workspaces/stream-v1 --workspace-id stream-v1 --lease-id <id> --staging-root <absolute-private-path> [--poll-seconds <5-300>]"); } const { runExeArtifactTransportOnce } = await import("./artifacts/exe-transport.js"); const config = { host, workspaceRoot, workspaceId, leaseId, stagingRoot } as const; const poll = valueAfter("--poll-seconds"); if (!poll) print(await runExeArtifactTransportOnce(store, config)); else { const seconds = positiveInteger(poll, "--poll-seconds", 300); if (seconds < 5) throw new Error("--poll-seconds must be at least 5"); while (true) { const result = await runExeArtifactTransportOnce(store, config); if (result.receipts.length || result.failures.length) print(result); await new Promise((resolve) => setTimeout(resolve, seconds * 1_000)); } } } else if (command === "artifact-request" || command === "request-artifact-storage") { const file = valueAfter("--file"); const root = valueAfter("--workspace-root"); const workspaceLeaseId = valueAfter("--workspace-lease-id"); const requestId = valueAfter("--request-id"); const artifactId = valueAfter("--artifact-id"); const kind = valueAfter("--kind"); const title = valueAfter("--title"); const summary = valueAfter("--summary"); const mediaType = valueAfter("--media-type") ?? "text/markdown"; const visibility = valueAfter("--visibility") ?? "private"; const provenanceSource = valueAfter("--provenance-source") ?? "agent-memory"; const provenanceLabel = valueAfter("--provenance-label") ?? "agent-memory skill source"; const supersedes = valueAfter("--supersedes"); if (!root || !workspaceLeaseId || !requestId || !file || !artifactId || !kind || !title || !summary) { throw new Error("Usage: thought artifact-request --workspace-root <absolute-local-lease-root> --workspace-lease-id <id> --file <relative-path> --request-id <id> --artifact-id <id> --kind <skill|spec|report|design|code|document|image> --title <title> --summary <summary> [--media-type <type>] [--visibility <private|sensitive>]"); } const { requestArtifactStorage } = await import("./artifacts/append.js"); const { ARTIFACT_KINDS, ARTIFACT_MEDIA_TYPES, ARTIFACT_VISIBILITY } = await import("./artifacts/types.js"); if (!(ARTIFACT_KINDS as readonly string[]).includes(kind)) { throw new Error(`--kind must be one of: ${ARTIFACT_KINDS.join(", ")}`); } if (!(ARTIFACT_MEDIA_TYPES as readonly string[]).includes(mediaType)) { throw new Error(`--media-type must be one of: ${ARTIFACT_MEDIA_TYPES.join(", ")}`); } if (!(ARTIFACT_VISIBILITY as readonly string[]).includes(visibility)) { throw new Error(`--visibility must be one of: ${ARTIFACT_VISIBILITY.join(", ")}`); } const result = await requestArtifactStorage({ store, workspaceRoot: root, workspaceLeaseId, workspaceRootLabel: valueAfter("--workspace-root-label") ?? "explicit workspace lease root", requestId, relativeFilePath: file, artifactId, artifactVersion: positiveInteger(valueAfter("--artifact-version") ?? "1", "--artifact-version", 1_000_000), kind: kind as import("./artifacts/types.js").ArtifactKind, title, summary, mediaType: mediaType as import("./artifacts/types.js").ArtifactMediaType, provenance: { source: provenanceSource, label: provenanceLabel }, visibility: visibility as import("./artifacts/types.js").ArtifactVisibility, ...(valueAfter("--requesting-agent") ? { requestingAgentId: valueAfter("--requesting-agent")! } : {}), ...(valueAfter("--requesting-run") ? { requestingRunId: valueAfter("--requesting-run")! } : {}), ...(valueAfter("--source-event") ? { sourceEventId: valueAfter("--source-event")! } : {}), ...(supersedes ? { supersedesArtifactEventId: supersedes } : {}), }); print({ artifactRequest: { eventId: result.requestEvent?.id, requestId }, artifact: { eventId: result.event.id, inserted: result.inserted, blobInserted: result.blobInserted, bodySha256: result.bodySha256, byteCount: result.byteCount, artifactId: result.event.payload.artifactId, artifactVersion: result.event.payload.artifactVersion, blobPath: (result.event.payload.blob as { relativePath: string }).relativePath } }); } else if (command === "artifacts") { const { buildArtifactCatalog } = await import("./artifacts/catalog.js"); print(await buildArtifactCatalog(store)); } else if (command === "status") { const events = await store.listEvents(); const runs = await store.listRuns(); const projection = await rebuildRootActivity(store); const review = await projectReviewQueue(store); const { buildArtifactCatalog } = await import("./artifacts/catalog.js"); const artifacts = await buildArtifactCatalog(store); print({ database: path.join(projectRoot, ".thoughtstream", "state", "jazz.sqlite"), events: events.length, runs: runs.length, completedRuns: runs.filter((run) => run.status === "completed").length, failedRuns: runs.filter((run) => run.status === "failed").length, review: review.counts, artifacts: { total: artifacts.total, active: artifacts.active, superseded: artifacts.superseded }, projection, }); } else if (command === "serve") { const host = valueAfter("--host") ?? "127.0.0.1"; const port = Number(valueAfter("--port") ?? "4317"); const reviewCapability = decodeReviewCapability(process.env.THOUGHTSTREAM_REVIEW_CAPABILITY_B64); const courseChatCapability = decodeCourseChatCapability(process.env.THOUGHTSTREAM_COURSE_CHAT_CAPABILITY_B64); const agentContextRoot = valueAfter("--agent-context-root"); if (agentContextRoot && !path.isAbsolute(agentContextRoot)) throw new Error("--agent-context-root must be absolute"); const server = await startInspectorServer(store, { host, port, ...(reviewCapability ? { reviewCapability } : {}), ...(courseChatCapability ? { courseChatCapability } : {}), ...(agentContextRoot ? { agentContextRoot: path.resolve(agentContextRoot) } : {}), }); process.stdout.write(`thought stream inspector: http://${host}:${port}\n`); await waitForShutdown(server); } else { throw new Error(`Unknown command: ${command}`); }} catch (error) { commandFailed = true; throw error;} finally { await store.close(); if (!commandFailed) { await Promise.all([flushStream(process.stdout), flushStream(process.stderr)]); process.exit(0); }}
function valueAfter(flag: string): string | undefined { const index = process.argv.indexOf(flag); return index >= 0 ? process.argv[index + 1] : undefined;}
function commaSeparated(value: string | undefined): string[] { return value ? value.split(",").map((item) => item.trim()).filter(Boolean) : [];}
function agentRunRetryPending(run: AgentRun): boolean { if (run.status !== "failed" && run.status !== "blocked") return false; const diagnostic = run.result?.failureDiagnostic; if (!diagnostic || typeof diagnostic !== "object" || Array.isArray(diagnostic)) return false; return diagnostic.progressDisposition === "retry-delayed" || diagnostic.progressDisposition === "deferred";}
function jsonObject(value: unknown, flag: string): Record<string, import("./core/json.js").JsonValue> { if (!value || typeof value !== "object" || Array.isArray(value)) throw new Error(`${flag} must contain a JSON object`); return value as Record<string, import("./core/json.js").JsonValue>;}
async function telegramSourceConfig(): Promise<TelegramWebhookSourceManifest> { const loaded = await loadThoughtStreamManifest(projectRoot, valueAfter("--config") ?? "thoughtstream.yaml"); const requestedSource = valueAfter("--source"); const source = loaded.manifest.sources.find((candidate) => ( candidate.kind === "telegram-webhook" && (requestedSource ? candidate.id === requestedSource : candidate.enabled) )); if (!source || source.kind !== "telegram-webhook") { throw new Error("No matching telegram-webhook source in the thought stream manifest"); } if (!source.enabled) throw new Error(`Telegram webhook source is disabled: ${source.id}`); return source;}
async function fastmailSourceConfig(): Promise<FastmailJmapSourceManifest> { const loaded = await loadThoughtStreamManifest(projectRoot, valueAfter("--config") ?? "thoughtstream.yaml"); const requestedSource = valueAfter("--source"); const source = loaded.manifest.sources.find((candidate) => ( candidate.kind === "fastmail-jmap" && (requestedSource ? candidate.id === requestedSource : candidate.enabled) )); if (!source || source.kind !== "fastmail-jmap") { throw new Error("No matching fastmail-jmap source in the thought stream manifest"); } if (!source.enabled) throw new Error(`Fastmail JMAP source is disabled: ${source.id}`); return source;}
async function xSourceConfig(): Promise<XWebhookSourceManifest> { const loaded = await loadThoughtStreamManifest(projectRoot, valueAfter("--config") ?? "thoughtstream.yaml"); const requestedSource = valueAfter("--source"); const source = loaded.manifest.sources.find((candidate) => ( candidate.kind === "x-webhook" && (requestedSource ? candidate.id === requestedSource : candidate.enabled) )); if (!source || source.kind !== "x-webhook") { throw new Error("No matching x-webhook source in the thought stream manifest"); } if (!source.enabled) throw new Error(`X webhook source is disabled: ${source.id}`); return source;}
function xManagementClient(source: XWebhookSourceManifest): XApiClient { const bearerToken = process.env[source.managementBearerTokenEnv]; if (!bearerToken) { throw new Error(`X management bearer token environment variable is unavailable: ${source.managementBearerTokenEnv}`); } return new XApiClient({ bearerToken, ...(process.env.THOUGHTSTREAM_X_API_BASE_URL ? { baseUrl: process.env.THOUGHTSTREAM_X_API_BASE_URL } : {}), timeoutMs: Math.max(1_000, source.requestTimeoutMs), });}
function exactXWebhook( source: XWebhookSourceManifest, webhooks: XWebhookRecord[], requireValid: boolean,): XWebhookRecord | undefined { const matches = webhooks.filter((webhook) => webhook.url === source.webhookUrl); if (matches.length > 1) throw new Error("Multiple X webhooks match the configured source URL"); const webhook = matches[0]; if (webhook && requireValid && !webhook.valid) throw new Error("The configured X webhook exists but is invalid"); return webhook;}
function publicXPlan(plan: XSubscriptionPlan) { return { source: plan.sourceId, webhookId: plan.webhookId, planHash: plan.planHash, create: plan.create.map((item) => ({ eventType: item.eventType, userId: item.userId, tag: item.tag })), update: plan.update, delete: plan.delete, unchangedCount: plan.unchanged.length, conflicts: plan.conflicts, unmanagedCount: plan.unmanagedCount, };}
function xUpstreamControl( source: XWebhookSourceManifest, webhook: XWebhookRecord | undefined, plan: XSubscriptionPlan | undefined,): SourceControlUpstreamState { const sourceOwnedLiveCount = plan ? plan.unchanged.length + plan.update.length + plan.delete.length : 0; const subscriptionsConverged = Boolean( webhook?.valid && plan && plan.create.length === 0 && plan.update.length === 0 && plan.delete.length === 0 && plan.conflicts.length === 0, ); return { registered: Boolean(webhook), valid: Boolean(webhook?.valid), ...(webhook ? { webhookId: webhook.id } : {}), desiredSubscriptionCount: source.expectedSubscriptions.length, liveSubscriptionCount: sourceOwnedLiveCount, subscriptionsConverged, checkedAt: new Date().toISOString(), };}
function requireExactFlag(flag: string, expected: string): void { const supplied = valueAfter(flag); if (supplied !== expected) throw new Error(`${flag} must exactly match ${expected}`);}
function telegramClient(source: TelegramWebhookSourceManifest): TelegramBotClient { const token = process.env[source.tokenEnv]; if (!token) throw new Error(`Telegram bot credential environment variable is unavailable: ${source.tokenEnv}`); return new TelegramBotClient({ token, ...(process.env.THOUGHTSTREAM_TELEGRAM_API_BASE_URL ? { baseUrl: process.env.THOUGHTSTREAM_TELEGRAM_API_BASE_URL } : {}), requestTimeoutMs: source.requestTimeoutMs, });}
function assertWebhookRegistration(source: TelegramWebhookSourceManifest, info: TelegramWebhookInfo): void { if (info.url !== source.webhookUrl) throw new Error("Telegram webhook URL verification failed"); if (info.max_connections !== 1) throw new Error("Telegram webhook registration must use exactly one upstream connection"); const actual = new Set(info.allowed_updates ?? []); if (TELEGRAM_WEBHOOK_ALLOWED_UPDATES.some((update) => !actual.has(update))) { throw new Error("Telegram webhook allowed-update verification failed"); }}
async function createAgentRuntime(): Promise<ThoughtAgentRuntime> { const manifestArgument = valueAfter("--config"); const defaultManifest = path.join(projectRoot, "thoughtstream.yaml"); const hasDefaultManifest = await fs.stat(defaultManifest).then((stat) => stat.isFile()).catch(() => false); const configuredScheduler = manifestArgument || hasDefaultManifest ? (await loadThoughtStreamManifest(projectRoot, manifestArgument ?? "thoughtstream.yaml")).manifest.runtime.scheduler : { maxConcurrentOperations: 4, reconcileIntervalMs: 1_000 }; const override = valueAfter("--max-concurrent-operations"); const maxConcurrentOperations = override ? positiveInteger(override, "--max-concurrent-operations", 64) : configuredScheduler.maxConcurrentOperations; const publicKnowledgePolicy = process.env.THOUGHTSTREAM_PUBLIC_KNOWLEDGE_POLICY_PATH?.trim(); const publicKnowledgeCatalog = process.env.THOUGHTSTREAM_PUBLIC_KNOWLEDGE_CATALOG_ROOT?.trim(); if (Boolean(publicKnowledgePolicy) !== Boolean(publicKnowledgeCatalog)) { throw new Error("Public Knowledge runtime requires both policy and catalog-root environment paths"); } return new ThoughtAgentRuntime(store, undefined, { maxConcurrentOperations, reconcileIntervalMs: configuredScheduler.reconcileIntervalMs, ...(publicKnowledgePolicy && publicKnowledgeCatalog ? { publicKnowledgeContext: { policyPath: publicKnowledgePolicy, catalogRoot: publicKnowledgeCatalog, }, } : {}), });}
async function delay(milliseconds: number): Promise<void> { await new Promise((resolve) => setTimeout(resolve, milliseconds));}
async function abortableDelay(milliseconds: number, signal: AbortSignal): Promise<void> { if (signal.aborted) return; await new Promise<void>((resolve) => { const timer = setTimeout(done, milliseconds); function done() { clearTimeout(timer); signal.removeEventListener("abort", done); resolve(); } signal.addEventListener("abort", done, { once: true }); });}
function fastmailRetryDelay(intervalMs: number, consecutiveFailures: number): number { if (consecutiveFailures <= 0) return intervalMs; return Math.min(intervalMs * (2 ** Math.min(consecutiveFailures - 1, 10)), 15 * 60_000);}
async function waitUntil(predicate: () => Promise<boolean>, timeoutMs = 5_000): Promise<void> { const deadline = Date.now() + timeoutMs; while (Date.now() < deadline) { if (await predicate()) return; await new Promise((resolve) => setTimeout(resolve, 10)); } throw new Error("Timed out waiting for consumer subscription progress");}
function positiveInteger(value: string, flag: string, maximum: number): number { const parsed = Number(value); if (!Number.isSafeInteger(parsed) || parsed <= 0 || parsed > maximum) { throw new Error(`${flag} must be a positive integer <= ${maximum}`); } return parsed;}
function nonnegativeInteger(value: string, flag: string, maximum: number): number { const parsed = Number(value); if (!Number.isSafeInteger(parsed) || parsed < 0 || parsed > maximum) { throw new Error(`${flag} must be a nonnegative integer <= ${maximum}`); } return parsed;}
function print(value: unknown): void { process.stdout.write(`${JSON.stringify(value, null, 2)}\n`);}
async function flushStream(stream: NodeJS.WriteStream): Promise<void> { await new Promise<void>((resolve, reject) => stream.write("", (error) => error ? reject(error) : resolve()));}
async function waitForShutdown(server: import("node:http").Server): Promise<void> { await waitForSignal(); server.closeAllConnections(); await Promise.race([ new Promise<void>((resolve) => server.close(() => resolve())), new Promise<void>((resolve) => setTimeout(resolve, 2_000)), ]);}
async function waitForSignal(): Promise<void> { await new Promise<void>((resolve) => { process.once("SIGINT", resolve); process.once("SIGTERM", resolve); });}
async function waitForSignalOrTimeout(timeoutMs: number | undefined): Promise<void> { await new Promise<void>((resolve) => { let timer: NodeJS.Timeout | undefined; const done = () => { if (timer) clearTimeout(timer); process.removeListener("SIGINT", done); process.removeListener("SIGTERM", done); resolve(); }; process.once("SIGINT", done); process.once("SIGTERM", done); if (timeoutMs !== undefined) timer = setTimeout(done, timeoutMs); });}
function isIgnoredWatchPath(root: string, candidate: string): boolean { const relative = path.relative(root, candidate).split(path.sep); return relative.some((part) => part === ".git" || part === ".thoughtstream" || part === "node_modules") || /(^|\/)\.env(?:\.|$)/.test(relative.join("/"));}