Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
2.0 kB · 56 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657import { 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<string, string>): 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 })); });}