Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
4.9 kB · 95 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596import 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); });});