import { spawn } from "node:child_process"; import fs from "node:fs/promises"; import path from "node:path"; import { afterEach, describe, expect, test } from "vitest"; import { temporaryProject } from "./helpers.js"; const roots: string[] = []; afterEach(async () => { await Promise.all(roots.splice(0).map((root) => fs.rm(root, { recursive: true, force: true }))); }); describe("thought stream telegram-spool command", () => { test("runs one explicit fixture ingestion, prints receipts, and exits", async () => { const project = await temporaryProject(); roots.push(project); const spool = path.join(project, "telegram.ndjson"); await fs.copyFile(path.join(process.cwd(), "fixtures", "telegram", "messages.ndjson"), spool); const result = await run([ "--silent", "thought", "stream", "telegram-spool", "--file", spool, "--source", "telegram:cli-fixture", "--max-records", "10", ], { THOUGHTSTREAM_ROOT: project }); expect(result.code).toBe(0); expect(result.stderr).toBe(""); const output = JSON.parse(result.stdout) as { ingest: { status: string; inserted: number; consumedRecords: number; events: unknown[] }; }; expect(output.ingest).toMatchObject({ status: "updated", inserted: 2, consumedRecords: 2 }); expect(output.ingest.events).toHaveLength(2); }, 10_000); }); async function run(args: string[], extraEnv: Record): Promise<{ code: number | null; stdout: string; stderr: string }> { return new Promise((resolve, reject) => { const child = spawn("pnpm", args, { cwd: process.cwd(), env: { ...process.env, ...extraEnv }, stdio: ["ignore", "pipe", "pipe"], }); let stdout = ""; let stderr = ""; child.stdout.on("data", (chunk) => { stdout += String(chunk); }); child.stderr.on("data", (chunk) => { stderr += String(chunk); }); child.once("error", reject); child.once("exit", (code) => resolve({ code, stdout, stderr })); }); }