#!/usr/bin/env node import 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 | 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((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 [--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 [--dids ] [--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 --from YYYYMMDDHHmm --to YYYYMMDDHHmm --confirm-webhook-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 --username "); 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 --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(); 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 --source fastmail:name --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 "); 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 "); 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 [--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 --message-id --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 --external-id [--privacy ] [--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 --candidate-runs ,"); } 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 "); 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 --disposition [--replacement-file ] [--context-root ] [--submission-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 "); 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 --kind [--criterion output-quality] [--criterion-version 1] [--compared-run ] [--replacement ] [--external-export-eligible] [--authorize-sensitive-external-export] [--notes ]"); } 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 --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 --staging-root [--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 --workspace-lease-id --file --request-id --artifact-id --kind --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("/")); }