Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
73 kB · 1725 lines
TypeScript
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726import { createServer, type IncomingMessage, type Server, type ServerResponse } from "node:http";import { execFile } from "node:child_process";import fs from "node:fs/promises";import path from "node:path";import { promisify } from "node:util";import { afterEach, describe, expect, test, vi } from "vitest";import { TelegramChannelDispatcher } from "../src/bridges/telegram-dispatcher.js";import { ThoughtAgentRuntime } from "../src/agents/runtime.js";import { loadAgentDeclarations } from "../src/agents/declarations.js";import type { AgentRunner } from "../src/agents/types.js";import { parseTelegramBotUpdate, TelegramBotClient, TelegramBotConnector, type TelegramBotUser,} from "../src/connectors/telegram-bot.js";import { startTelegramWebhookServer } from "../src/connectors/telegram-webhook.js";import { JetstreamConnector } from "../src/connectors/jetstream.js";import type { JazzThoughtStore } from "../src/jazz/store.js";import { describeRunResult } from "../src/projections/activity.js";import { activeJudgments } from "../src/projections/effective-output.js";import { stableKey } from "../src/core/ids.js";import { projectTrainingExamples } from "../src/training/judgments.js";import type { AgentRun } from "../src/store/types.js";import { temporaryProject, testDeclarationEnvironment, testStore } from "./helpers.js";
const stores: JazzThoughtStore[] = [];const roots: string[] = [];const servers: Server[] = [];const execFileAsync = promisify(execFile);
afterEach(async () => { vi.restoreAllMocks(); await Promise.all(stores.splice(0).map((store) => store.close())); await Promise.all(roots.splice(0).map((root) => fs.rm(root, { recursive: true, force: true }))); await Promise.all(servers.splice(0).map((server) => new Promise<void>((resolve) => server.close(() => resolve()))));});
describe("TelegramBotConnector", () => { test("projects feedback once only for accepted feedback-shaped updates", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); let calls = 0; const connector = new TelegramBotConnector({ id: "telegram:thoughtstream", bot: telegramBotIdentity(), allowedChatIds: ["123456789"], reactionFeedback: [{ chatId: "123456789", allowedUserIds: ["123456789"] }], feedbackProjector: async () => { calls += 1; return { examined: 0, projected: 0, retracted: 0, skipped: 0, judgmentEventIds: [] }; }, });
await connector.ingest(store, parseTelegramBotUpdate(telegramMessageUpdate(80, 8, "ordinary"))); expect(calls).toBe(0);
const unauthorized = telegramMessageUpdate(81, 9, "/correct anything") as Record<string, any>; unauthorized.message.chat.id = 999; unauthorized.message.from.id = 999; await connector.ingest(store, parseTelegramBotUpdate(unauthorized)); expect(calls).toBe(0);
await connector.ingest(store, parseTelegramBotUpdate({ update_id: 82, message_reaction: { chat: { id: 123456789, type: "private" }, message_id: 8, user: { id: 123456789, is_bot: false, first_name: "Cameron" }, date: 1784042402, old_reaction: [], new_reaction: [{ type: "emoji", emoji: "👍" }], }, })); expect(calls).toBe(1); });
test("admits at most one image when a malformed update contains both photo and image-document fields", async () => { const fixture = await telegramFixture(); const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const artifactRoot = path.join(project, "artifacts"); await fs.mkdir(artifactRoot); const client = new TelegramBotClient({ token: "fixture-token", baseUrl: fixture.baseUrl }); const connector = new TelegramBotConnector({ id: "telegram:thoughtstream", bot: telegramBotIdentity(), allowedChatIds: ["123456789"], imageExtraction: { client, artifactRoot }, }); const update = telegramMessageUpdate(90, 9, "") as Record<string, any>; delete update.message.text; update.message.photo = [{ file_id: "photo-file", file_unique_id: "photo-unique", width: 100, height: 100, file_size: 8 }]; update.message.document = { file_id: "document-file", file_unique_id: "document-unique", file_name: "second.png", mime_type: "image/png", file_size: 8, }; const result = await connector.ingest(store, parseTelegramBotUpdate(update)); expect(result.events[0]?.payload.attachments).toEqual([ expect.objectContaining({ kind: "image", status: "stored", id: "photo-unique" }), ]); expect(fixture.imageFileDownloads()).toBe(1); });
test("preserves empty non-image attachments without admitting them to the conversation event type", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const connector = new TelegramBotConnector({ id: "telegram:thoughtstream", bot: telegramBotIdentity(), allowedChatIds: ["123456789"], }); const update = telegramMessageUpdate(91, 10, "") as Record<string, any>; delete update.message.text; update.message.document = { file_id: "pdf-file", file_unique_id: "pdf-unique", file_name: "evidence.pdf", mime_type: "application/pdf", file_size: 1_024, }; const result = await connector.ingest(store, parseTelegramBotUpdate(update)); expect(result.events[0]).toMatchObject({ type: "stream.thought.source.telegram.nonconversation", payload: { text: "", attachments: [expect.objectContaining({ kind: "file", id: "pdf-unique" })] }, }); });
test("durably ingests allowlisted webhook deliveries without using the high-water mark as admission", async () => { const fixture = await telegramFixture(); const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const client = new TelegramBotClient({ token: "fixture-token", baseUrl: fixture.baseUrl }); const connector = new TelegramBotConnector({ id: "telegram:thoughtstream", bot: telegramBotIdentity(), allowedChatIds: ["123456789"], });
await expect(client.identity()).resolves.toMatchObject({ id: 8765422491 }); await store.upsertSourceCursor({ id: "cursor:telegram:thoughtstream", source: "telegram:thoughtstream", cursor: { revision: "telegram-bot-api-v2", botId: "8765422491", updateOffset: 100 }, updatedAt: "2026-07-15T00:00:00.000Z", }); const first = await connector.ingest(store, parseTelegramBotUpdate(fixture.updates[0])); expect(first).toMatchObject({ accepted: 1, ignored: 0, inserted: 1 }); expect(first.events[0]).toMatchObject({ type: "stream.thought.source.telegram.message", privacy: "sensitive", payload: { accountId: "8765422491", accountUsername: "CameronStreamBot", chatId: "123456789", messageId: "10", text: "a thoughtstream blip", }, }); expect(first.cursor.cursor).toMatchObject({ botId: "8765422491", highestUpdateId: 100, revision: "telegram-bot-api-webhook-v1", }); const ignored = await connector.ingest(store, parseTelegramBotUpdate(fixture.updates[1])); expect(ignored).toMatchObject({ accepted: 0, ignored: 1, inserted: 0 }); expect(ignored.cursor.cursor).toMatchObject({ highestUpdateId: 101 }); expect((await store.getSourceCursor("cursor:telegram:thoughtstream"))?.cursor).toMatchObject({ highestUpdateId: 101 }); const lowerReplay = await connector.ingest(store, parseTelegramBotUpdate(fixture.updates[0])); expect(lowerReplay).toMatchObject({ accepted: 1, inserted: 0, unchanged: 1 }); expect(lowerReplay.cursor.cursor).toMatchObject({ highestUpdateId: 101 }); expect(JSON.stringify(await store.listEvents())).not.toContain("fixture-token"); });
test("turns only allowlisted thumbs reactions on delivered run messages into append-only judgments", async () => { const fixture = await telegramFixture(); const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const client = new TelegramBotClient({ token: "fixture-token", baseUrl: fixture.baseUrl }); const connector = new TelegramBotConnector({ id: "telegram:thoughtstream-bot", bot: telegramBotIdentity(), allowedChatIds: ["123456789"], reactionFeedback: [{ chatId: "123456789", allowedUserIds: ["123456789"] }], }); const ingress = await connector.ingest(store, parseTelegramBotUpdate(fixture.updates[0])); const trigger = ingress.events[0]!; const declarations = await loadAgentDeclarations(path.join(process.cwd(), "agents"), testDeclarationEnvironment); const loadedDeclaration = declarations.find((candidate) => candidate.id === "telegram-conversation"); const declaration = loadedDeclaration ? structuredClone(loadedDeclaration) : undefined; if (!declaration) throw new Error("Missing Telegram conversation declaration"); declaration.enabled = true; declaration.initialReplay = "beginning"; declaration.sourcePatterns = [trigger.source]; delete declaration.contextDocumentMaxChars; delete declaration.contextDocumentSubscriptions; delete declaration.conversationCompaction; const runner: AgentRunner = { mode: "pi", run: async () => ({ summary: "A reaction target", tags: ["telegram"], importance: "normal", confidence: 1, }), }; const runtime = new ThoughtAgentRuntime(store, [runner]); const [processed] = await runtime.consumeBacklog([declaration]); const run = await store.getRun(processed!.runId); expect(run).toBeDefined(); expect(describeRunResult(run!)).toBe("A reaction target"); const dispatcher = new TelegramChannelDispatcher({ id: "telegram-dispatcher:telegram:thoughtstream-bot:123456789", client, chatId: "123456789", allowedSources: ["telegram:thoughtstream-bot"], allowedActors: ["123456789"], runStatuses: ["completed"], }); await dispatcher.sendPending(store, { includeNormal: true }); const targetMessageId = fixture.sentMessageIds[0]!; const [receipt] = await store.listEvents({ types: ["stream.thought.action.telegram.send.delivered"] });
fixture.updates.push(reactionUpdate(102, targetMessageId, [], [{ type: "emoji", emoji: "👍" }])); const appendEvent = store.appendEvent.bind(store); const projectionFailure = vi.spyOn(store, "appendEvent").mockImplementation(async (candidate) => { if (candidate.type === "stream.thought.judgment.training-example") { throw new Error("fixture projection interruption"); } return appendEvent(candidate); }); await expect(connector.ingest(store, parseTelegramBotUpdate(fixture.updates.at(-1)))).rejects.toThrow("fixture projection interruption"); const interruptedCursor = await store.getSourceCursor("cursor:telegram:thoughtstream-bot"); expect(interruptedCursor).toMatchObject({ cursor: { revision: "telegram-bot-api-webhook-v1", highestUpdateId: 102 }, lastError: "fixture projection interruption", }); projectionFailure.mockRestore(); const recoveryIngest = await connector.ingest(store, parseTelegramBotUpdate(fixture.updates.at(-1))); expect(recoveryIngest).toMatchObject({ accepted: 1, inserted: 0, unchanged: 1 }); const reactions = await store.listEvents({ source: "telegram:thoughtstream-bot", types: ["stream.thought.source.telegram.reaction"], }); expect(reactions[0]).toMatchObject({ rootEventId: trigger.rootEventId, parentEventId: receipt!.id, payload: { messageId: targetMessageId, senderId: "123456789", resolutionStatus: "resolved", feedbackAction: "set", feedbackLabel: "positive", deliveryReceiptEventId: receipt!.id, runId: run!.id, sourceRootEventId: trigger.rootEventId, }, }); const [positive] = await store.listEvents({ source: "judgment:telegram-reaction", types: ["stream.thought.judgment.training-example"], }); expect(positive).toMatchObject({ schemaVersion: 2, rootEventId: trigger.rootEventId, parentEventId: reactions[0]!.id, payload: { runId: run!.id, kind: "accept", criterion: "telegram-reaction", criterionVersion: 1, qualityEligible: true, externalExportEligible: false, feedbackSourceEventId: reactions[0]!.id, deliveryReceiptEventId: receipt!.id, }, }); expect(await projectTrainingExamples(store)).toEqual([]); expect(await projectTrainingExamples(store, { includeSensitivePrivate: true })).toEqual([]); expect(await runtime.consumeBacklog([declaration])).toHaveLength(0);
fixture.updates.push(reactionUpdate( 103, targetMessageId, [{ type: "emoji", emoji: "👍" }], [{ type: "emoji", emoji: "👎" }], )); await connector.ingest(store, parseTelegramBotUpdate(fixture.updates.at(-1))); const judgmentsAfterChange = await store.listEvents({ source: "judgment:telegram-reaction", types: ["stream.thought.judgment.training-example"], }); expect(judgmentsAfterChange).toHaveLength(2); const negative = judgmentsAfterChange.find((event) => event.payload.kind === "reject")!; expect(negative.payload.supersedesJudgmentEventId).toBe(positive!.id); expect(await projectTrainingExamples(store)).toEqual([]); expect(await projectTrainingExamples(store, { includeSensitivePrivate: true })).toEqual([]);
fixture.updates.push(reactionUpdate( 104, targetMessageId, [{ type: "emoji", emoji: "👎" }], [], )); await connector.ingest(store, parseTelegramBotUpdate(fixture.updates.at(-1))); const [retraction] = await store.listEvents({ source: "judgment:telegram-reaction", types: ["stream.thought.judgment.training-example.retracted"], }); expect(retraction).toMatchObject({ rootEventId: trigger.rootEventId, payload: { retractedJudgmentEventId: negative.id, deliveryReceiptEventId: receipt!.id, runId: run!.id, }, }); expect(await projectTrainingExamples(store)).toEqual([]);
fixture.updates.push( reactionUpdate(105, targetMessageId, [], [ { type: "emoji", emoji: "👍" }, { type: "emoji", emoji: "❤️" }, ]), reactionUpdate(106, "999999", [], [{ type: "emoji", emoji: "👍" }]), reactionUpdate(107, targetMessageId, [], [{ type: "emoji", emoji: "👍" }], "999"), ); const filtered = []; for (const update of fixture.updates.slice(-3)) { filtered.push(await connector.ingest(store, parseTelegramBotUpdate(update))); } expect(filtered.map(({ accepted, ignored, inserted }) => ({ accepted, ignored, inserted }))).toEqual([ { accepted: 1, ignored: 0, inserted: 1 }, { accepted: 1, ignored: 0, inserted: 1 }, { accepted: 0, ignored: 1, inserted: 0 }, ]); const finalReactions = await store.listEvents({ source: "telegram:thoughtstream-bot", types: ["stream.thought.source.telegram.reaction"], }); expect(finalReactions).toHaveLength(5); expect(finalReactions.find((event) => event.payload.updateId === "105")?.payload).toMatchObject({ resolutionStatus: "resolved", feedbackAction: "none", newReactions: [ { type: "emoji", emoji: "👍" }, { type: "emoji", emoji: "❤️" }, ], }); expect(finalReactions.find((event) => event.payload.updateId === "106")?.payload).toMatchObject({ resolutionStatus: "unknown-delivery", feedbackAction: "none", }); expect(await store.listEvents({ source: "judgment:telegram-reaction" })).toHaveLength(3); expect(await runtime.consumeBacklog([declaration])).toHaveLength(0); });
test("turns an allowlisted reply-bound /correct command into an exact active correction without inference or egress", async () => { const fixture = await telegramFixture(); const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const connector = new TelegramBotConnector({ id: "telegram:thoughtstream-bot", bot: telegramBotIdentity(), allowedChatIds: ["123456789"], reactionFeedback: [{ chatId: "123456789", allowedUserIds: ["123456789"] }], }); const trigger = (await connector.ingest(store, parseTelegramBotUpdate(fixture.updates[0]))).events[0]!; const declarations = await loadAgentDeclarations(path.join(process.cwd(), "agents"), testDeclarationEnvironment); const loadedDeclaration = declarations.find((candidate) => candidate.id === "telegram-conversation"); const declaration = loadedDeclaration ? structuredClone(loadedDeclaration) : undefined; if (!declaration) throw new Error("Missing Telegram conversation declaration"); declaration.enabled = true; declaration.initialReplay = "beginning"; declaration.sourcePatterns = [trigger.source]; delete declaration.contextDocumentMaxChars; delete declaration.contextDocumentSubscriptions; delete declaration.conversationCompaction; const runner: AgentRunner = { mode: "pi", run: async () => ({ summary: "An answer that should be replaced", tags: ["conversation"], importance: "normal", confidence: 0.5, }), }; const runtime = new ThoughtAgentRuntime(store, [runner]); const [processed] = await runtime.consumeBacklog([declaration]); const run = await store.getRun(processed!.runId); if (!run) throw new Error("Missing completed Telegram run"); const [originalOutputId] = run.outputEventIds; const dispatcher = new TelegramChannelDispatcher({ id: "telegram-dispatcher:telegram:thoughtstream-bot:123456789", client: new TelegramBotClient({ token: "fixture-token", baseUrl: fixture.baseUrl }), chatId: "123456789", allowedSources: ["telegram:thoughtstream-bot"], allowedActors: ["123456789"], runStatuses: ["completed"], }); await dispatcher.sendPending(store, { includeNormal: true }); const targetMessageId = fixture.sentMessageIds[0]!; const [receipt] = await store.listEvents({ types: ["stream.thought.action.telegram.send.delivered"] }); if (!receipt) throw new Error("Missing Telegram delivery receipt");
const negativeUpdate = reactionUpdate(202, targetMessageId, [], [{ type: "emoji", emoji: "👎" }]); await connector.ingest(store, parseTelegramBotUpdate(negativeUpdate)); const [reject] = await store.listEvents({ source: "judgment:telegram-reaction", types: ["stream.thought.judgment.training-example"], }); if (!reject) throw new Error("Missing reaction rejection");
const correction = correctionUpdate(203, 40, "/correct Of course.", targetMessageId); const corrected = await connector.ingest(store, parseTelegramBotUpdate(correction)); expect(corrected).toMatchObject({ accepted: 1, ignored: 0, inserted: 1 }); const [correctionSource] = await store.listEvents({ source: "telegram:thoughtstream-bot", types: ["stream.thought.source.telegram.correction"], }); expect(correctionSource).toMatchObject({ rootEventId: trigger.rootEventId, parentEventId: receipt.id, privacy: "sensitive", payload: { messageId: "40", replyToMessageId: targetMessageId, replacementText: "Of course.", replacementChars: 10, resolutionStatus: "resolved", deliveryReceiptEventId: receipt.id, runId: run.id, outputEventId: originalOutputId, sourceRootEventId: trigger.rootEventId, }, }); const [correctionJudgment] = await store.listEvents({ source: "judgment:telegram-correction", types: ["stream.thought.judgment.training-example"], }); expect(correctionJudgment).toMatchObject({ rootEventId: trigger.rootEventId, parentEventId: correctionSource!.id, privacy: "sensitive", payload: { runId: run.id, outputEventId: originalOutputId, kind: "correct", criterion: "telegram-reaction", criterionVersion: 1, qualityEligible: true, externalExportEligible: false, feedbackSourceEventId: correctionSource!.id, deliveryReceiptEventId: receipt.id, supersedesJudgmentEventId: reject.id, replacementOutput: { summary: "Of course.", tags: ["conversation"], importance: "normal", confidence: 0.5, }, }, }); const active = await activeJudgments(store); expect(active.inactiveIds.has(reject.id)).toBe(true); expect(active.active.some((event) => event.id === correctionJudgment!.id)).toBe(true); const effective = await store.getProjection(stableKey("effective-output", run.id)); expect(effective).toMatchObject({ projectionVersion: 2, lastEventId: correctionJudgment!.id, payload: { status: "corrected", outputEventId: originalOutputId, judgmentEventId: correctionJudgment!.id, feedbackSourceEventId: correctionSource!.id, authority: "correct", structuredOutput: { summary: "Of course.", tags: ["conversation"], importance: "normal", confidence: 0.5, }, }, }); expect(await projectTrainingExamples(store)).toEqual([]); expect(await projectTrainingExamples(store, { includeSensitivePrivate: true })).toEqual([]); expect(fixture.sentMessages).toHaveLength(1); expect(await runtime.consumeBacklog([declaration])).toHaveLength(0); expect(await store.listRuns()).toHaveLength(1);
const replay = await connector.ingest(store, parseTelegramBotUpdate(correction)); expect(replay).toMatchObject({ accepted: 1, inserted: 0, unchanged: 1 }); expect(await store.listEvents({ source: "judgment:telegram-correction" })).toHaveLength(1);
const noReply = await connector.ingest(store, parseTelegramBotUpdate( correctionUpdate(204, 41, "/correct No target", undefined), )); expect(noReply).toMatchObject({ accepted: 1, inserted: 1 }); const unresolved = await connector.ingest(store, parseTelegramBotUpdate( correctionUpdate(205, 42, "/correct Unknown target", "999999"), )); expect(unresolved).toMatchObject({ accepted: 1, inserted: 1 }); const correctionSources = await store.listEvents({ source: "telegram:thoughtstream-bot", types: ["stream.thought.source.telegram.correction"], }); expect(correctionSources.find((event) => event.payload.messageId === "41")?.payload.resolutionStatus) .toBe("missing-reply-target"); expect(correctionSources.find((event) => event.payload.messageId === "42")?.payload.resolutionStatus) .toBe("unknown-delivery");
const wrongTargetMessageId = "998"; await store.appendEvent({ type: "stream.thought.action.telegram.send.delivered", schemaVersion: 1, source: "telegram-dispatcher:telegram:thoughtstream-bot:123456789", sourceKind: "system", externalId: "corrupt-parent-delivery", idempotencyKey: "corrupt-parent-delivery", occurredAt: "2026-07-15T01:00:00.000Z", actor: "telegram-dispatcher:telegram:thoughtstream-bot:123456789", rootEventId: trigger.rootEventId, parentEventId: trigger.id, correlationId: "corrupt-parent-delivery", privacy: "sensitive", payload: { status: "delivered", chatId: "123456789", messageId: wrongTargetMessageId, runIds: [run.id] }, }); await connector.ingest(store, parseTelegramBotUpdate( correctionUpdate(2051, 421, "/correct Corrupt target", wrongTargetMessageId), )); const corruptTarget = (await store.listEvents({ source: "telegram:thoughtstream-bot", types: ["stream.thought.source.telegram.correction"], })).find((event) => event.payload.messageId === "421"); expect(corruptTarget?.payload.resolutionStatus).toBe("invalid-lineage"); expect(await store.listEvents({ source: "judgment:telegram-correction" })).toHaveLength(1);
const ineligible = await connector.ingest(store, parseTelegramBotUpdate( correctionUpdate(206, 43, "/correct Unauthorized", targetMessageId, "999"), )); expect(ineligible).toMatchObject({ accepted: 0, ignored: 1, inserted: 0 }); expect(await store.listEvents({ source: "judgment:telegram-correction" })).toHaveLength(1); expect(await store.listEvents({ types: ["stream.thought.source.telegram.message"] })).toHaveLength(1); expect(await store.listRuns()).toHaveLength(1); expect(fixture.sentMessages).toHaveLength(1); }, 30_000);
test("keeps authenticated webhook ingress send-dark and lets the separate dispatcher honor boot configuration", async () => { const fixture = await telegramFixture(); const project = await temporaryProject(); roots.push(project); const manifestPath = path.join(project, "thoughtstream.yaml"); await fs.writeFile(manifestPath, telegramManifest(true)); const ingressStore = await testStore(project); const connector = new TelegramBotConnector({ id: "telegram:fixture", bot: telegramBotIdentity(), allowedChatIds: ["123456789"], reactionFeedback: [{ chatId: "123456789", allowedUserIds: ["123456789"] }], }); const receiver = await startTelegramWebhookServer({ connector, store: ingressStore, secretToken: "fixture-webhook-secret", path: "/webhooks/telegram", host: "127.0.0.1", port: 0, }); const response = await fetch(`http://${receiver.host}:${receiver.port}${receiver.path}`, { method: "POST", headers: { "content-type": "application/json", "x-telegram-bot-api-secret-token": "fixture-webhook-secret", }, body: JSON.stringify(fixture.updates[0]), }); expect(response.status).toBe(204); await receiver.close(); await ingressStore.close(); expect(fixture.sentMessages).toHaveLength(0); await runTelegramCli("telegram-dispatcher", project, fixture.baseUrl); expect(fixture.sentMessages).toEqual([ { chatId: "123456789", text: "thought stream fixture is live." }, ]); await runTelegramCli("telegram-dispatcher", project, fixture.baseUrl); expect(fixture.sentMessages).toHaveLength(2);
await fs.writeFile(manifestPath, telegramManifest(false)); await runTelegramCli("telegram-dispatcher", project, fixture.baseUrl); expect(fixture.sentMessages).toHaveLength(2); }, 30_000);
test("requires the exact webhook secret and acknowledges only durable valid updates", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const connector = new TelegramBotConnector({ id: "telegram:webhook-security", bot: telegramBotIdentity(), allowedChatIds: ["123456789"], }); const receiver = await startTelegramWebhookServer({ connector, store, secretToken: "fixture-secret", path: "/webhooks/telegram", host: "127.0.0.1", port: 0, maxBodyBytes: 1_024, }); const endpoint = `http://${receiver.host}:${receiver.port}${receiver.path}`; const valid = telegramMessageUpdate(200, 20, "webhook fixture");
expect((await fetch(endpoint, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify(valid) })).status).toBe(401); expect((await fetch(endpoint, { method: "POST", headers: { "content-type": "application/json", "x-telegram-bot-api-secret-token": "wrong" }, body: JSON.stringify(valid) })).status).toBe(401); expect((await fetch(`${endpoint}?query=forbidden`, { method: "POST", headers: webhookHeaders(), body: JSON.stringify(valid) })).status).toBe(404); expect((await fetch(endpoint, { method: "GET", headers: webhookHeaders() })).status).toBe(405); expect((await fetch(endpoint, { method: "POST", headers: { ...webhookHeaders(), "content-type": "text/plain" }, body: JSON.stringify(valid) })).status).toBe(415); expect((await fetch(endpoint, { method: "POST", headers: webhookHeaders(), body: "{" })).status).toBe(400); expect((await fetch(endpoint, { method: "POST", headers: webhookHeaders(), body: JSON.stringify({ update_id: 201, unexpected: "x".repeat(2_000) }) })).status).toBe(413); expect(await store.listEvents({ types: ["stream.thought.source.telegram.message"] })).toHaveLength(0);
const producerBatch = vi.spyOn(store, "appendProducerBatch").mockRejectedValueOnce(new Error("fixture durable write failure")); expect((await fetch(endpoint, { method: "POST", headers: webhookHeaders(), body: JSON.stringify(valid) })).status).toBe(503); producerBatch.mockRestore(); expect(await store.listEvents({ types: ["stream.thought.source.telegram.message"] })).toHaveLength(0);
expect((await fetch(endpoint, { method: "POST", headers: webhookHeaders(), body: JSON.stringify(valid) })).status).toBe(204); expect((await fetch(endpoint, { method: "POST", headers: webhookHeaders(), body: JSON.stringify(valid) })).status).toBe(204); expect(await store.listEvents({ types: ["stream.thought.source.telegram.message"] })).toHaveLength(1); await receiver.close(); });
test("serializes concurrent authenticated deliveries before touching Jazz", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const connector = new TelegramBotConnector({ id: "telegram:webhook-ordering", bot: telegramBotIdentity(), allowedChatIds: ["123456789"], }); const actualIngest = connector.ingest.bind(connector); let active = 0; let maximumActive = 0; vi.spyOn(connector, "ingest").mockImplementation(async (...arguments_) => { active += 1; maximumActive = Math.max(maximumActive, active); await new Promise((resolve) => setTimeout(resolve, 20)); try { return await actualIngest(...arguments_); } finally { active -= 1; } }); const receiver = await startTelegramWebhookServer({ connector, store, secretToken: "fixture-secret", path: "/webhooks/telegram", host: "127.0.0.1", port: 0, }); const endpoint = `http://${receiver.host}:${receiver.port}${receiver.path}`; const responses = await Promise.all([ fetch(endpoint, { method: "POST", headers: webhookHeaders(), body: JSON.stringify(telegramMessageUpdate(300, 30, "first")) }), fetch(endpoint, { method: "POST", headers: webhookHeaders(), body: JSON.stringify(telegramMessageUpdate(301, 31, "second")) }), ]); expect(responses.map((response) => response.status)).toEqual([204, 204]); expect(maximumActive).toBe(1); expect((await store.listEvents({ types: ["stream.thought.source.telegram.message"] })).map((event) => event.payload.text)).toEqual(["first", "second"]); await receiver.close(); });
test("registers and deletes the configured webhook explicitly with one upstream connection", async () => { const fixture = await telegramFixture(); const project = await temporaryProject(); roots.push(project); await fs.writeFile(path.join(project, "thoughtstream.yaml"), telegramManifest(false));
const registered = await runTelegramCli("telegram-webhook-register", project, fixture.baseUrl); expect(fixture.webhookRegistrations).toEqual([expect.objectContaining({ url: "https://thoughtstream.example/webhooks/telegram", secret_token: "fixture-webhook-secret", max_connections: 1, allowed_updates: ["message", "edited_message", "message_reaction"], drop_pending_updates: false, })]); expect(registered.stdout).not.toContain("fixture-webhook-secret"); expect(registered.stdout).not.toContain("fixture-token");
await runTelegramCli("telegram-webhook-delete", project, fixture.baseUrl); expect(fixture.webhookDeletions).toEqual([{ drop_pending_updates: false }]); });
test("registers and verifies the fixed command menu, then sends help through durable egress receipts", async () => { const fixture = await telegramFixture(); const project = await temporaryProject(); roots.push(project); await fs.writeFile(path.join(project, "thoughtstream.yaml"), telegramManifest(false));
const result = await runTelegramCli("telegram-menu-register", project, fixture.baseUrl);
expect(fixture.botCommands).toEqual([ { command: "help", description: "Show commands, feedback, and links" }, { command: "focus", description: "Propose a new focus from plain language" }, { command: "correct", description: "Reply to a message with /correct new text" }, ]); expect(fixture.botName).toBe("The Stream"); expect(fixture.menuButtons.get("123456789")).toBe("commands"); expect(fixture.sentMessages).toHaveLength(1); expect(fixture.sentMessages[0]?.text).toContain("The Stream"); expect(fixture.sentMessages[0]?.text).toContain("https://thought.stream/inspector/"); expect(result.stdout).toContain('"menuButton": "commands"'); expect(result.stdout).toContain('"botName": "The Stream"'); expect(result.stdout).toContain('"helpNotice": "sent"'); expect(result.stdout).not.toContain("fixture-token"); const store = await testStore(project); const delivered = await store.listEvents({ types: ["stream.thought.action.telegram.send.delivered"] }); expect(delivered).toHaveLength(1); expect(delivered[0]?.payload).toMatchObject({ messageKind: "operational-notice", runIds: [] }); // alpha55 storage takes an exclusive RocksDB lock, so the CLI subprocess // cannot open the same project root while this process holds a store open. // Close before spawning the second CLI invocation. await store.close();
const quiet = await runTelegramCli("telegram-menu-register", project, fixture.baseUrl, ["--skip-help-notice"]); expect(fixture.botName).toBe("The Stream"); expect(fixture.sentMessages).toHaveLength(1); expect(quiet.stdout).toContain('"helpNotice": "skipped"'); expect(quiet.stdout).toContain('"deliveries": []'); }, 30_000);});
describe("TelegramChannelDispatcher", () => { test.each([ { route: "resident-proposal", agentId: "resident-letta-conversation", triggerSource: "batch:stream-activity", triggerActor: "batch:stream-activity-ten-minute", }, { route: "social-proposal", agentId: "cameron-social-listener", triggerSource: "batch:cameron-social-observations", triggerActor: "batch:cameron-social-observation-window", }, ])("delivers only an exact high-importance proposal from the $route route", async ({ route, agentId, triggerSource, triggerActor, }) => { const fixture = await telegramFixture(); const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const connector = new JetstreamConnector({ id: "jetstream:cameron-bluesky", collections: ["app.bsky.feed.post"], dids: ["did:plc:gfrmhdmjvxn2sjedzboeudef"], }); await connector.ingestBatch(store, [postCommit(1786381200000000, "notification-member", "A source observation.")]); const [member] = await store.listEvents({ source: "jetstream:cameron-bluesky" }); expect(member).toBeDefined(); const trigger = (await store.appendEvent({ type: "stream.thought.derived.event.batch", schemaVersion: 1, source: triggerSource, sourceKind: "system", externalId: "stream-activity-batch-notification", idempotencyKey: "stream-activity-batch-notification", occurredAt: "2026-08-10T17:00:00.000Z", actor: triggerActor, correlationId: "stream-activity-batch-notification", privacy: member!.privacy, payload: { declaration: { id: "stream-activity-ten-minute", version: 1, fingerprint: "a".repeat(64) }, flushReason: "max-age", firstOccurredAt: member!.occurredAt, lastOccurredAt: member!.occurredAt, members: [{ eventId: member!.id, source: member!.source, sourceSequence: member!.sourceSequence, type: member!.type, schemaVersion: member!.schemaVersion, privacy: member!.privacy, occurredAt: member!.occurredAt, observedAt: member!.observedAt, payloadHash: member!.payloadHash, }], }, })).event; const exactResult = { summary: "Two source observations now point to the same release change.", tags: ["social-synthesis", "notify-cameron"], importance: "high", confidence: 0.9, recommendation: { target: "cameron-telegram", reason: "The connection is concrete and time-sensitive.", proposedAction: "notify", }, }; const appendProposalRun = async (id: string, result: typeof exactResult) => { const output = (await store.appendEvent({ type: "stream.thought.derived.message.observation", schemaVersion: 1, source: `agent:${agentId}`, sourceKind: "agent", externalId: id, idempotencyKey: `${id}:output`, occurredAt: "2026-08-10T17:00:01.000Z", actor: agentId, rootEventId: trigger.rootEventId, parentEventId: trigger.id, correlationId: trigger.correlationId, privacy: trigger.privacy, payload: { runId: id, executionKey: `execution:${id}`, inputEventId: trigger.id, inputSourceSequence: trigger.sourceSequence, ...result, }, traceId: id, })).event; await store.upsertRun({ id, executionKey: `execution:${id}`, triggerEventId: trigger.id, agentId, agentVersion: 1, status: "completed", inputEventIds: [trigger.id], outputEventIds: [output.id], attempt: 1, provider: "letta-cloud", model: "chatgpt-plus-pro/gpt-5.6-terra", promptHash: "fixture-resident-prompt", contextManifest: {}, result, createdAt: "2026-08-10T17:00:00.000Z", startedAt: "2026-08-10T17:00:00.000Z", completedAt: "2026-08-10T17:00:01.000Z", updatedAt: "2026-08-10T17:00:01.000Z", }); }; await appendProposalRun("run-resident-exact", exactResult); await appendProposalRun("run-resident-approximate", { ...exactResult, tags: ["social-synthesis"], summary: "This high observation omitted the exact proposal tag.", }); const dispatcher = new TelegramChannelDispatcher({ id: `telegram-notifier:${route}`, client: new TelegramBotClient({ token: "fixture-token", baseUrl: fixture.baseUrl }), chatId: "123456789", allowedSources: ["telegram:thoughtstream-bot", triggerSource], allowedActors: [triggerActor], directReplyAgentIds: ["resident-letta-conversation"], directReplySources: ["telegram:thoughtstream-bot"], notificationProposalAgentIds: [agentId], notificationProposalSources: [triggerSource], runStatuses: ["completed", "failed"], });
await expect(dispatcher.sendPending(store, { includeNormal: true, now: new Date("2026-08-10T17:00:02.000Z") })).resolves.toMatchObject({ pending: 2, eligible: 1, delivered: 1, failed: 0, }); expect(fixture.sentMessages).toEqual([{ chatId: "123456789", text: expect.stringContaining("The Stream · Observation\n\nTwo source observations now point to the same release change."), }]); const [receipt] = await store.listEvents({ source: `telegram-notifier:${route}`, types: ["stream.thought.action.telegram.send.delivered"], }); expect(receipt).toMatchObject({ payload: { messageKind: "stream-observation", runIds: ["run-resident-exact"] } }); expect(() => new TelegramChannelDispatcher({ id: "telegram-notifier:invalid-resident-overlap", client: new TelegramBotClient({ token: "fixture-token", baseUrl: fixture.baseUrl }), chatId: "123456789", allowedSources: [triggerSource], directReplyAgentIds: [agentId], directReplySources: [triggerSource], notificationProposalAgentIds: [agentId], notificationProposalSources: [triggerSource], })).toThrow("route tuples must be disjoint"); });
test("refreshes typing only for a recent running direct-reply run in the exact chat", async () => { const fixture = await telegramFixture(); const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const connector = new TelegramBotConnector({ id: "telegram:thoughtstream-bot", bot: telegramBotIdentity(), allowedChatIds: ["123456789"], }); const trigger = (await connector.ingest(store, parseTelegramBotUpdate(fixture.updates[0]))).events[0]!; const startedAt = "2026-07-28T05:00:00.000Z"; const run = runningTelegramRun(trigger.id, "telegram-conversation", startedAt); await store.upsertRun(run); const dispatcher = new TelegramChannelDispatcher({ id: "telegram-notifier:typing", client: new TelegramBotClient({ token: "fixture-token", baseUrl: fixture.baseUrl }), chatId: "123456789", allowedSources: ["telegram:thoughtstream-bot"], allowedActors: ["123456789"], directReplyAgentIds: ["telegram-conversation"], });
await expect(dispatcher.refreshTyping(store, { now: new Date("2026-07-28T05:00:01.000Z"), })).resolves.toEqual({ activeRuns: 1, sent: 1, failed: 0, deferred: 0 }); await expect(dispatcher.refreshTyping(store, { now: new Date("2026-07-28T05:00:04.000Z"), })).resolves.toEqual({ activeRuns: 1, sent: 0, failed: 0, deferred: 1 }); await expect(dispatcher.refreshTyping(store, { now: new Date("2026-07-28T05:00:05.000Z"), })).resolves.toEqual({ activeRuns: 1, sent: 1, failed: 0, deferred: 0 }); expect(fixture.typingActions).toEqual([ { chatId: "123456789", action: "typing" }, { chatId: "123456789", action: "typing" }, ]);
await store.upsertRun({ ...run, status: "completed", completedAt: "2026-07-28T05:00:06.000Z", updatedAt: "2026-07-28T05:00:06.000Z", result: { summary: "done", tags: ["conversation"], importance: "normal", confidence: 0.5 }, }); await store.upsertRun(runningTelegramRun( trigger.id, "other-agent", "2026-07-28T05:00:06.000Z", "run:wrong-agent", )); await store.upsertRun(runningTelegramRun( trigger.id, "telegram-conversation", "2026-07-28T04:50:00.000Z", "run:stale", )); const wrongChatTrigger = await store.appendEvent({ type: "stream.thought.source.telegram.message", schemaVersion: 1, source: "telegram:thoughtstream-bot", sourceKind: "telegram", externalId: "wrong-chat", idempotencyKey: "wrong-chat", occurredAt: "2026-07-28T05:00:06.000Z", actor: "123456789", correlationId: "wrong-chat", privacy: "sensitive", payload: { ...trigger.payload, chatId: "999", messageId: "11", text: "wrong route" }, }); await store.upsertRun(runningTelegramRun( wrongChatTrigger.event.id, "telegram-conversation", "2026-07-28T05:00:06.000Z", "run:wrong-chat", )); await expect(dispatcher.refreshTyping(store, { now: new Date("2026-07-28T05:00:10.000Z"), })).resolves.toEqual({ activeRuns: 0, sent: 0, failed: 0, deferred: 0 }); expect(fixture.typingActions).toHaveLength(2); expect(await store.listEvents({ types: ["stream.thought.action.telegram.send.started"] })).toHaveLength(0); });
test("typing failure is best-effort and cannot block the terminal reply", async () => { const warning = vi.spyOn(console, "warn").mockImplementation(() => {}); const fixture = await telegramFixture(); const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const connector = new TelegramBotConnector({ id: "telegram:thoughtstream-bot", bot: telegramBotIdentity(), allowedChatIds: ["123456789"], }); const trigger = (await connector.ingest(store, parseTelegramBotUpdate(fixture.updates[0]))).events[0]!; const run = runningTelegramRun(trigger.id, "telegram-conversation", "2026-07-28T05:00:00.000Z"); await store.upsertRun(run); const dispatcher = new TelegramChannelDispatcher({ id: "telegram-notifier:typing-failure", client: new TelegramBotClient({ token: "fixture-token", baseUrl: fixture.baseUrl }), chatId: "123456789", allowedSources: ["telegram:thoughtstream-bot"], allowedActors: ["123456789"], directReplyAgentIds: ["telegram-conversation"], runStatuses: ["completed"], }); fixture.setTypingFailure(true); await expect(dispatcher.refreshTyping(store, { now: new Date("2026-07-28T05:00:01.000Z"), })).resolves.toEqual({ activeRuns: 1, sent: 0, failed: 1, deferred: 0 }); expect(warning).toHaveBeenCalledWith("Telegram typing action unavailable");
await store.upsertRun({ ...run, status: "completed", completedAt: "2026-07-28T05:00:02.000Z", updatedAt: "2026-07-28T05:00:02.000Z", result: { summary: "The reply still arrives.", tags: ["conversation"], importance: "normal", confidence: 0.5 }, }); await expect(dispatcher.sendPending(store, { includeNormal: true })).resolves.toMatchObject({ delivered: 1, failed: 0, }); expect(fixture.sentMessages).toEqual([{ chatId: "123456789", text: "The reply still arrives." }]); });
test("carries an allowlisted Telegram blip through an agent run to a tied delivery receipt", async () => { const fixture = await telegramFixture(); const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const client = new TelegramBotClient({ token: "fixture-token", baseUrl: fixture.baseUrl }); const connector = new TelegramBotConnector({ id: "telegram:thoughtstream-bot", bot: telegramBotIdentity(), allowedChatIds: ["123456789"], }); const ingress = await connector.ingest(store, parseTelegramBotUpdate(fixture.updates[0])); const trigger = ingress.events[0]!; const declarations = await loadAgentDeclarations(path.join(process.cwd(), "agents"), testDeclarationEnvironment); const loadedDeclaration = declarations.find((candidate) => candidate.id === "telegram-conversation"); const declaration = loadedDeclaration ? structuredClone(loadedDeclaration) : undefined; if (!declaration) throw new Error("Missing Telegram conversation declaration"); declaration.enabled = true; declaration.initialReplay = "beginning"; declaration.sourcePatterns = [trigger.source]; delete declaration.contextDocumentMaxChars; delete declaration.contextDocumentSubscriptions; delete declaration.conversationCompaction; const runner: AgentRunner = { mode: "pi", run: async (input) => ({ summary: `Received Telegram blip: ${String(input.event.payload.text)} via ${String(input.event.payload.accountUsername)} account ${String(input.event.payload.accountId)}`, tags: ["telegram"], importance: "normal", confidence: 1, }), }; const runtime = new ThoughtAgentRuntime(store, [runner]); const processed = await runtime.consumeBacklog([declaration]); expect(processed).toHaveLength(1); const run = await store.getRun(processed[0]!.runId); expect(run).toMatchObject({ status: "completed", triggerEventId: trigger.id, agentId: "telegram-conversation", }); await store.upsertRun({ id: "run-internal-compactor-failure", executionKey: "execution-internal-compactor-failure", triggerEventId: trigger.id, agentId: "telegram-conversation-compactor", agentVersion: 1, status: "failed", inputEventIds: [trigger.id], outputEventIds: [], attempt: 1, provider: "tinker", model: "thinkingmachines/Inkling-Small", promptHash: "fixture-compactor", contextManifest: {}, result: { failureDiagnostic: { code: "timeout", stage: "provider", progressDisposition: "retry-delayed", }, }, createdAt: new Date().toISOString(), completedAt: new Date().toISOString(), updatedAt: new Date().toISOString(), });
const dispatcher = new TelegramChannelDispatcher({ id: "telegram-notifier:thoughtstream-feedback", client, chatId: "123456789", allowedSources: ["telegram:thoughtstream-bot"], allowedActors: ["123456789"], directReplyAgentIds: ["telegram-conversation"], runStatuses: ["completed", "failed"], maxMessagesPerWindow: 3, windowMs: 60_000, }); const result = await dispatcher.sendPending(store, { includeNormal: true }); expect(result).toMatchObject({ pending: 2, eligible: 1, delivered: 1, failed: 0 }); expect(fixture.sentMessages[0]?.text).not.toContain("The Stream · Telegram blip"); expect(fixture.sentMessages[0]?.text).toContain("Received Telegram blip: a thoughtstream blip"); expect(fixture.sentMessages[0]?.text).toContain("[private route]"); expect(fixture.sentMessages[0]?.text).not.toContain("CameronStreamBot"); expect(fixture.sentMessages[0]?.text).not.toContain("8765422491"); const [receipt] = await store.listEvents({ source: "telegram-notifier:thoughtstream-feedback", types: ["stream.thought.action.telegram.send.delivered"], }); expect(receipt).toMatchObject({ rootEventId: trigger.rootEventId, payload: { runIds: [run!.id], status: "delivered" }, }); });
test("sends posts immediately, coalesces likes, and enforces a durable window cap", async () => { const fixture = await telegramFixture(); const project = await temporaryProject(); roots.push(project); await fs.cp(path.join(process.cwd(), "agents"), path.join(project, "agents"), { recursive: true }); await fs.cp(path.join(process.cwd(), "prompts"), path.join(project, "prompts"), { recursive: true }); const store = await testStore(project); stores.push(store); const declarations = await loadAgentDeclarations(path.join(project, "agents"), testDeclarationEnvironment); const runtime = new ThoughtAgentRuntime(store); const connector = new JetstreamConnector({ id: "jetstream:cameron-bluesky", collections: ["app.bsky.feed.post", "app.bsky.feed.like"], dids: ["did:plc:gfrmhdmjvxn2sjedzboeudef"], }); await connector.ingestBatch(store, [ postCommit(1784042400000000, "post-one", "A bounded post"), likeCommit(1784042401000000, "like-one", "liked-one"), likeCommit(1784042402000000, "like-two", "liked-two"), ]); await runtime.consumeBacklog(declarations); const completedRuns = await store.listRuns({ status: "completed" }); const likeRuns = []; for (const run of completedRuns) { const trigger = await store.getEvent(run.triggerEventId); if (trigger?.payload.collection === "app.bsky.feed.like") likeRuns.push(run); } expect(likeRuns).toHaveLength(2); const firstLikeCompletedAt = new Date(Date.now() - 61_000); likeRuns[0]!.completedAt = firstLikeCompletedAt.toISOString(); likeRuns[0]!.updatedAt = firstLikeCompletedAt.toISOString(); likeRuns[1]!.completedAt = new Date(firstLikeCompletedAt.getTime() + 30_000).toISOString(); likeRuns[1]!.updatedAt = likeRuns[1]!.completedAt; await store.upsertRun(likeRuns[0]!); await store.upsertRun(likeRuns[1]!); const client = new TelegramBotClient({ token: "fixture-token", baseUrl: fixture.baseUrl }); const dispatcher = new TelegramChannelDispatcher({ id: "telegram-notifier:thoughtstream", client, chatId: "123456789", allowedSources: ["jetstream:cameron-bluesky"], allowedActors: ["did:plc:gfrmhdmjvxn2sjedzboeudef"], maxMessagesPerWindow: 2, windowMs: 60_000, likeDigestDelayMs: 60_000, maxLikesPerDigest: 10, }); const firstNow = new Date();
const firstActivation = await dispatcher.activate(store, new Date("2026-07-15T02:39:00.000Z")); const repeatedActivation = await dispatcher.activate(store, new Date("2026-07-15T03:39:00.000Z")); expect(firstActivation).toBe("2026-07-15T02:39:00.000Z"); expect(repeatedActivation).toBe("2026-07-15T03:39:00.000Z");
const first = await dispatcher.sendPending(store, { includeNormal: true, now: firstNow, }); expect(first).toMatchObject({ pending: 3, eligible: 3, delivered: 2, failed: 0, rateLimited: 0 }); expect(fixture.sentMessages).toHaveLength(2); expect(fixture.sentMessages.some((message) => message.text.includes("The Stream · Bluesky post"))).toBe(true); expect(fixture.sentMessages.some((message) => message.text.includes("The Stream · 2 Bluesky likes"))).toBe(true);
await connector.ingestBatch(store, [postCommit(1784042403000000, "post-two", "A rate-limited post")]); await runtime.consumeBacklog(declarations); const limited = await dispatcher.sendPending(store, { includeNormal: true, now: new Date(firstNow.getTime() + 30_000), }); expect(limited).toMatchObject({ pending: 1, eligible: 1, delivered: 0, rateLimited: 1, availableMessages: 0 }); expect(fixture.sentMessages).toHaveLength(2);
const resumed = await dispatcher.sendPending(store, { includeNormal: true, now: new Date(firstNow.getTime() + 61_000), }); expect(resumed).toMatchObject({ pending: 1, delivered: 1, rateLimited: 0, availableMessages: 2 }); expect(fixture.sentMessages).toHaveLength(3); const receipts = await store.listEvents({ source: "telegram-notifier:thoughtstream" }); expect(receipts.filter((event) => event.type.endsWith(".delivered"))).toHaveLength(3); });
test("dispatches only classified diagnostics for failed runs and ignores legacy raw traces", async () => { const fixture = await telegramFixture(); const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const trigger = await store.appendEvent({ type: "stream.thought.source.atproto.commit", schemaVersion: 1, source: "jetstream:cameron-bluesky", sourceKind: "jetstream", externalId: "failed-post", idempotencyKey: "failed-post", occurredAt: "2026-07-15T04:00:00.000Z", actor: "did:plc:gfrmhdmjvxn2sjedzboeudef", correlationId: "failed-post", privacy: "public-source", payload: { collection: "app.bsky.feed.post" }, }); await store.upsertRun({ id: "run_failed_transcript", executionKey: "failed-transcript", triggerEventId: trigger.event.id, agentId: "bluesky-watch", agentVersion: 1, status: "failed", inputEventIds: [trigger.event.id], outputEventIds: [], attempt: 1, provider: "tinker", model: "Qwen/Qwen3.5-4B", promptHash: "fixture", contextManifest: {}, errorText: "PRIVATE_PROVIDER_ERROR_BODY", result: { failureDiagnostic: { code: "invalid-final-output", stage: "final-output-validation", reason: "invalid-json", assistantMessages: 1, textParts: 1, textChars: 3219, textSha256: "fixture-text-hash", thinkingParts: 1, thinkingChars: 1373, thinkingSha256: "fixture-thinking-hash", thinkingRedacted: true, toolCallParts: 0, otherParts: 0, stopReason: "stop", }, }, createdAt: "2026-07-15T04:00:01.000Z", completedAt: "2026-07-15T04:00:02.000Z", updatedAt: "2026-07-15T04:00:02.000Z", }); await store.upsertRun({ id: "run_failed_transcript_retry", executionKey: "failed-transcript", triggerEventId: trigger.event.id, agentId: "bluesky-watch", agentVersion: 1, status: "failed", inputEventIds: [trigger.event.id], outputEventIds: [], attempt: 2, provider: "tinker", model: "Qwen/Qwen3.5-4B", promptHash: "fixture", contextManifest: {}, result: { failureDiagnostic: { code: "invalid-final-output", stage: "final-output-validation", reason: "invalid-json", assistantMessages: 1, textParts: 1, textChars: 3219, textSha256: "fixture-retry-text-hash", thinkingParts: 1, thinkingChars: 1373, thinkingSha256: "fixture-retry-thinking-hash", thinkingRedacted: true, toolCallParts: 0, otherParts: 0, stopReason: "stop", }, }, createdAt: "2026-07-15T04:00:03.000Z", completedAt: "2026-07-15T04:00:04.000Z", updatedAt: "2026-07-15T04:00:04.000Z", }); await store.appendTrace({ id: "trace_failed_transcript", runId: "run_failed_transcript", sequence: 1, type: "pi.message_end", payload: { data: { type: "message_end", message: { role: "assistant", content: [ { type: "thinking", thinking: "SECRET_PROVIDER_REASONING" }, { type: "text", text: "SECRET_MODEL_TEXT {almost valid json" }, { type: "toolCall", name: "fetch", arguments: { token: "SECRET_TOOL_ARGUMENT" } }, ], }, }, }, createdAt: "2026-07-15T04:00:02.000Z", }); const dispatcher = new TelegramChannelDispatcher({ id: "telegram-notifier:failures", client: new TelegramBotClient({ token: "fixture-token", baseUrl: fixture.baseUrl }), chatId: "123456789", allowedSources: ["jetstream:cameron-bluesky"], allowedActors: ["did:plc:gfrmhdmjvxn2sjedzboeudef"], runStatuses: ["failed"], });
const result = await dispatcher.sendPending(store, { includeNormal: false }); expect(result).toMatchObject({ pending: 1, eligible: 1, delivered: 1, failed: 0 }); const deliveredText = fixture.sentMessages[0]?.text ?? ""; expect(deliveredText).toContain("The Stream · Agent run failed"); expect(deliveredText).toContain("The final text was not one complete JSON value"); expect(deliveredText).toContain("Final text: 1 part, 3219 characters (content redacted)"); expect(deliveredText).toContain("Provider thinking: 1 part, 1373 characters (redacted)"); expect(deliveredText).toContain("No derived event was emitted"); expect(deliveredText).not.toContain("SECRET_MODEL_TEXT"); expect(deliveredText).not.toContain("SECRET_PROVIDER_REASONING"); expect(deliveredText).not.toContain("SECRET_TOOL_ARGUMENT"); expect(deliveredText).not.toContain("PRIVATE_PROVIDER_ERROR_BODY"); expect(deliveredText).not.toContain("fixture-text-hash"); const [receipt] = await store.listEvents({ types: ["stream.thought.action.telegram.send.delivered"] }); expect(receipt?.payload.messageKind).toBe("agent-failure"); expect(receipt?.payload.runIds).toEqual(["run_failed_transcript_retry"]);
const repeated = await dispatcher.sendPending(store, { includeNormal: false }); expect(repeated).toMatchObject({ pending: 0, delivered: 0, failed: 0 }); expect(fixture.sentMessages).toHaveLength(1); });});
async function telegramFixture(): Promise<{ baseUrl: string; sentMessages: Array<{ chatId: string; text: string }>; typingActions: Array<{ chatId: string; action: string }>; setTypingFailure: (failed: boolean) => void; sentMessageIds: string[]; updates: Array<Record<string, unknown>>; webhookRegistrations: Array<Record<string, unknown>>; webhookDeletions: Array<Record<string, unknown>>; botCommands: Array<{ command: string; description: string }>; botName: string; menuButtons: Map<string, string>; imageFileDownloads: () => number;}> { const sentMessages: Array<{ chatId: string; text: string }> = []; const typingActions: Array<{ chatId: string; action: string }> = []; const sentMessageIds: string[] = []; const webhookRegistrations: Array<Record<string, unknown>> = []; const webhookDeletions: Array<Record<string, unknown>> = []; let botCommands: Array<{ command: string; description: string }> = []; let botName = "Stream"; const menuButtons = new Map<string, string>(); let webhookUrl = ""; let webhookMaxConnections: number | undefined; let webhookAllowedUpdates: string[] | undefined; let typingFailure = false; let imageFileDownloads = 0; const updates: Array<Record<string, unknown>> = [ { update_id: 100, message: { message_id: 10, from: { id: 123456789, is_bot: false, first_name: "Cameron", username: "just_cameron" }, chat: { id: 123456789, type: "private", first_name: "Cameron", username: "just_cameron" }, date: 1784042400, text: "a thoughtstream blip", }, }, { update_id: 101, message: { message_id: 11, from: { id: 999, is_bot: false, first_name: "Wrong" }, chat: { id: 999, type: "private", first_name: "Wrong" }, date: 1784042401, text: "must be ignored", }, }, ]; const server = createServer(async (request, response) => { if (request.url?.includes("/file/botfixture-token/")) { imageFileDownloads += 1; response.writeHead(200, { "content-type": "image/png", "content-length": "8" }); response.end(Buffer.from([0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a])); return; } const body = await readJson(request); if (request.url?.endsWith("/getFile")) return json(response, { ok: true, result: { file_path: "photos/image.png" }, }); if (request.url?.endsWith("/getMe")) return json(response, { ok: true, result: { id: 8765422491, is_bot: true, first_name: "The Stream", username: "CameronStreamBot" }, }); if (request.url?.endsWith("/setWebhook")) { webhookRegistrations.push(body); webhookUrl = String(body.url ?? ""); webhookMaxConnections = typeof body.max_connections === "number" ? body.max_connections : undefined; webhookAllowedUpdates = Array.isArray(body.allowed_updates) ? body.allowed_updates.filter((value): value is string => typeof value === "string") : undefined; return json(response, { ok: true, result: true }); } if (request.url?.endsWith("/deleteWebhook")) { webhookDeletions.push(body); webhookUrl = ""; webhookMaxConnections = undefined; webhookAllowedUpdates = undefined; return json(response, { ok: true, result: true }); } if (request.url?.endsWith("/getWebhookInfo")) { return json(response, { ok: true, result: { url: webhookUrl, has_custom_certificate: false, pending_update_count: 0, ...(webhookMaxConnections === undefined ? {} : { max_connections: webhookMaxConnections }), ...(webhookAllowedUpdates === undefined ? {} : { allowed_updates: webhookAllowedUpdates }), }, }); } if (request.url?.endsWith("/setMyCommands")) { botCommands = Array.isArray(body.commands) ? body.commands.map((entry) => ({ command: String((entry as Record<string, unknown>).command), description: String((entry as Record<string, unknown>).description), })) : []; return json(response, { ok: true, result: true }); } if (request.url?.endsWith("/getMyCommands")) { return json(response, { ok: true, result: botCommands }); } if (request.url?.endsWith("/setMyName")) { botName = String(body.name); return json(response, { ok: true, result: true }); } if (request.url?.endsWith("/getMyName")) { return json(response, { ok: true, result: { name: botName } }); } if (request.url?.endsWith("/setChatMenuButton")) { menuButtons.set(String(body.chat_id), String((body.menu_button as Record<string, unknown>)?.type)); return json(response, { ok: true, result: true }); } if (request.url?.endsWith("/getChatMenuButton")) { return json(response, { ok: true, result: { type: menuButtons.get(String(body.chat_id)) ?? "default" } }); } if (request.url?.endsWith("/sendMessage")) { sentMessages.push({ chatId: String(body.chat_id), text: String(body.text) }); const messageId = 200 + sentMessages.length; sentMessageIds.push(String(messageId)); return json(response, { ok: true, result: { message_id: messageId, from: { id: 8765422491, is_bot: true, first_name: "The Stream", username: "CameronStreamBot" }, chat: { id: Number(body.chat_id), type: "private", first_name: "Cameron" }, date: 1784042500 + sentMessages.length, text: body.text, }, }); } if (request.url?.endsWith("/sendChatAction")) { typingActions.push({ chatId: String(body.chat_id), action: String(body.action) }); if (typingFailure) { return json(response, { ok: false, error_code: 503, description: "fixture unavailable" }, 503); } return json(response, { ok: true, result: true }); } return json(response, { ok: false, error_code: 404, description: "not found" }, 404); }); servers.push(server); await new Promise<void>((resolve, reject) => { server.listen(0, "127.0.0.1", resolve); server.once("error", reject); }); const address = server.address(); if (!address || typeof address === "string") throw new Error("Missing Telegram fixture address"); return { baseUrl: `http://127.0.0.1:${address.port}`, sentMessages, typingActions, setTypingFailure: (failed) => { typingFailure = failed; }, sentMessageIds, updates, webhookRegistrations, webhookDeletions, get botCommands() { return botCommands; }, get botName() { return botName; }, menuButtons, imageFileDownloads: () => imageFileDownloads, };}
function telegramBotIdentity(): TelegramBotUser { return { id: 8765422491, is_bot: true, first_name: "The Stream", username: "CameronStreamBot" };}
function runningTelegramRun( triggerEventId: string, agentId: string, startedAt: string, id = "run:telegram-typing",): AgentRun { return { id, executionKey: `${agentId}:${triggerEventId}`, triggerEventId, agentId, agentVersion: 4, status: "running", inputEventIds: [triggerEventId], outputEventIds: [], attempt: 1, provider: "letta-agent-sdk", model: "fixture", promptHash: "fixture-prompt", contextManifest: {}, createdAt: startedAt, startedAt, updatedAt: startedAt, };}
function telegramMessageUpdate(updateId: number, messageId: number, text: string): Record<string, unknown> { return { update_id: updateId, message: { message_id: messageId, from: { id: 123456789, is_bot: false, first_name: "Cameron", username: "just_cameron" }, chat: { id: 123456789, type: "private", first_name: "Cameron", username: "just_cameron" }, date: 1784042400 + updateId, text, }, };}
function correctionUpdate( updateId: number, messageId: number, text: string, replyToMessageId?: string, userId = "123456789",): Record<string, unknown> { return { update_id: updateId, message: { message_id: messageId, from: { id: Number(userId), is_bot: false, first_name: userId === "123456789" ? "Cameron" : "Wrong" }, chat: { id: 123456789, type: "private", first_name: "Cameron", username: "just_cameron" }, date: 1784042400 + updateId, text, ...(replyToMessageId ? { reply_to_message: { message_id: Number(replyToMessageId) } } : {}), }, };}
function webhookHeaders(): Record<string, string> { return { "content-type": "application/json", "x-telegram-bot-api-secret-token": "fixture-secret", };}
function reactionUpdate( updateId: number, messageId: string, oldReaction: Array<Record<string, unknown>>, newReaction: Array<Record<string, unknown>>, userId = "123456789",): Record<string, unknown> { return { update_id: updateId, message_reaction: { chat: { id: 123456789, type: "private", first_name: "Cameron", username: "just_cameron" }, message_id: Number(messageId), user: { id: Number(userId), is_bot: false, first_name: userId === "123456789" ? "Cameron" : "Wrong" }, date: 1784042500 + updateId, old_reaction: oldReaction, new_reaction: newReaction, }, };}
async function readJson(request: IncomingMessage): Promise<Record<string, unknown>> { const chunks: Buffer[] = []; for await (const chunk of request) chunks.push(Buffer.from(chunk)); return JSON.parse(Buffer.concat(chunks).toString("utf8")) as Record<string, unknown>;}
function json(response: ServerResponse, body: unknown, status = 200): void { response.writeHead(status, { "content-type": "application/json" }); response.end(JSON.stringify(body));}
function postCommit(timeUs: number, rkey: string, text: string): Record<string, unknown> { return { did: "did:plc:gfrmhdmjvxn2sjedzboeudef", time_us: timeUs, kind: "commit", commit: { rev: `rev-${rkey}`, operation: "create", collection: "app.bsky.feed.post", rkey, cid: `bafy-${rkey}`, record: { $type: "app.bsky.feed.post", text, createdAt: "2026-07-15T02:40:00.000Z" }, }, };}
function likeCommit(timeUs: number, rkey: string, targetRkey: string): Record<string, unknown> { return { did: "did:plc:gfrmhdmjvxn2sjedzboeudef", time_us: timeUs, kind: "commit", commit: { rev: `rev-${rkey}`, operation: "create", collection: "app.bsky.feed.like", rkey, cid: `bafy-${rkey}`, record: { $type: "app.bsky.feed.like", subject: { uri: `at://did:plc:target/app.bsky.feed.post/${targetRkey}`, cid: `bafy-${targetRkey}` }, createdAt: "2026-07-15T02:40:00.000Z", }, }, };}
async function runTelegramCli( command: "telegram-webhook-register" | "telegram-webhook-delete" | "telegram-menu-register" | "telegram-dispatcher", project: string, baseUrl: string, extraArgs: string[] = [],): Promise<{ stdout: string; stderr: string }> { return execFileAsync(process.execPath, [ "--import", "tsx", path.join(process.cwd(), "src/cli.ts"), command, "--once", "--max-runtime", "1", ...extraArgs, ], { cwd: process.cwd(), env: { HOME: process.env.HOME, PATH: process.env.PATH, NODE_PATH: process.env.NODE_PATH, THOUGHTSTREAM_ROOT: project, THOUGHTSTREAM_JAZZ_AUTO: "1", THOUGHTSTREAM_TELEGRAM_API_BASE_URL: baseUrl, FIXTURE_TELEGRAM_BOT_TOKEN: "fixture-token", FIXTURE_TELEGRAM_WEBHOOK_SECRET: "fixture-webhook-secret", }, maxBuffer: 1024 * 1024, });}
function telegramManifest(bootMessageEnabled: boolean): string { return `version: 1runtime: revision: telegram-fixturesources: - id: telegram:fixture kind: telegram-webhook enabled: true tokenEnv: FIXTURE_TELEGRAM_BOT_TOKEN webhookSecretEnv: FIXTURE_TELEGRAM_WEBHOOK_SECRET webhookUrl: https://thoughtstream.example/webhooks/telegram webhookPath: /webhooks/telegram listenHost: 127.0.0.1 listenPort: 4318 maxBodyBytes: 1048576 requestTimeoutMs: 5000 channels: - id: "123456789" enabled: true bootMessage: enabled: ${bootMessageEnabled} text: thought stream fixture is live. reactionFeedback: enabled: true allowedUserIds: - "123456789"`;}