import fs from "node:fs/promises"; import path from "node:path"; import { afterEach, describe, expect, test } from "vitest"; import { parseTelegramSpoolRecord, TelegramSpoolConnector } from "../src/connectors/telegram-spool.js"; import type { JazzThoughtStore } from "../src/jazz/store.js"; import { temporaryProject, testStore } from "./helpers.js"; const stores: JazzThoughtStore[] = []; const roots: string[] = []; afterEach(async () => { 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 }))); }); describe("TelegramSpoolConnector", () => { test("validates the captured route and attachment-reference contract", async () => { const lines = (await fs.readFile(path.join(process.cwd(), "fixtures", "telegram", "messages.ndjson"), "utf8")).trim().split("\n"); expect(parseTelegramSpoolRecord(lines[1]!)).toMatchObject({ accountId: "account-fixture", chatId: "chat-fixture", messageId: "message-2", route: { agentId: "agent-fixture", conversationId: "conversation-fixture" }, attachments: [{ kind: "audio", reference: "telegram-file:fixture-voice-1" }], }); expect(() => parseTelegramSpoolRecord(JSON.stringify({ ...JSON.parse(lines[0]!) as object, token: "must-not-pass-through", }))).toThrow("Invalid Telegram spool record"); }); test("ingests bounded records across restarts and waits for a complete trailing line", async () => { const project = await temporaryProject(); roots.push(project); const spool = path.join(project, "telegram.ndjson"); const fixture = await fs.readFile(path.join(process.cwd(), "fixtures", "telegram", "messages.ndjson"), "utf8"); const [firstLine, secondLine] = fixture.trimEnd().split("\n"); await fs.writeFile(spool, `${firstLine}\n${secondLine}`); const connector = new TelegramSpoolConnector({ id: "telegram:fixture-spool", file: spool, maxRecords: 1 }); const firstStore = await testStore(project); const first = await connector.ingest(firstStore); expect(first).toMatchObject({ status: "updated", inserted: 1, consumedRecords: 1 }); expect(first.partialTrailingBytes).toBe(Buffer.byteLength(secondLine!)); expect(first.events[0]).toMatchObject({ type: "stream.thought.source.telegram.message", privacy: "sensitive", externalId: "account-fixture:chat-fixture:message-1", payload: { text: "A synthetic captured message", route: { agentId: "agent-fixture", conversationId: "conversation-fixture" }, }, }); await firstStore.close(); const restartedStore = await testStore(project); stores.push(restartedStore); const partial = await connector.ingest(restartedStore); expect(partial).toMatchObject({ status: "unchanged", inserted: 0, consumedRecords: 0 }); expect(partial.partialTrailingBytes).toBe(Buffer.byteLength(secondLine!)); expect(await restartedStore.listEvents({ types: ["stream.thought.source.telegram.message"] })).toHaveLength(1); await fs.appendFile(spool, "\n"); const completed = await connector.ingest(restartedStore); expect(completed).toMatchObject({ status: "updated", inserted: 1, consumedRecords: 1, partialTrailingBytes: 0 }); expect(completed.events[0]?.payload.attachments).toEqual([expect.objectContaining({ reference: "telegram-file:fixture-voice-1" })]); expect((await restartedStore.getSourceCursor("cursor:telegram:fixture-spool"))?.cursor).toMatchObject({ spoolRevision: "telegram-spool-ndjson-v1", recordCount: 2, byteOffset: Buffer.byteLength(fixture), }); }); test("fails closed when consumed spool content changes", async () => { const project = await temporaryProject(); roots.push(project); const spool = path.join(project, "telegram.ndjson"); const fixture = await fs.readFile(path.join(process.cwd(), "fixtures", "telegram", "messages.ndjson"), "utf8"); await fs.writeFile(spool, fixture); const connector = new TelegramSpoolConnector({ id: "telegram:continuity", file: spool }); const store = await testStore(project); stores.push(store); const first = await connector.ingest(store); expect(first.inserted).toBe(2); const before = await store.getSourceCursor("cursor:telegram:continuity"); const changed = fixture.replace("A synthetic captured message", "A changed synthetic message"); await fs.writeFile(spool, changed); await expect(connector.ingest(store)).rejects.toThrow(/changed behind its durable cursor|truncated behind its durable cursor/); const after = await store.getSourceCursor("cursor:telegram:continuity"); expect(after?.cursor).toEqual(before?.cursor); expect(after?.lastFailureAt).toBeDefined(); expect(await store.listEvents({ types: ["stream.thought.source.telegram.message"] })).toHaveLength(2); }); });