Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
4.4 kB · 108 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109import { spawn, type ChildProcess } 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[] = [];const children: ChildProcess[] = [];
afterEach(async () => { for (const child of children.splice(0)) { if (child.exitCode === null) child.kill("SIGKILL"); } await Promise.all(roots.splice(0).map((root) => fs.rm(root, { recursive: true, force: true })));});
describe("thought stream watch", () => { test("runs producer-only without declarations, coalesces mutations, and exits cleanly on SIGTERM", async () => { const project = await temporaryProject(); roots.push(project); const vault = path.join(project, "vault"); await fs.mkdir(vault); await fs.writeFile(path.join(vault, "watched.md"), "# Initial\n");
const child = spawn(process.execPath, [ "--import", "tsx", "src/cli.ts", "watch", "--root", vault, "--source", "filesystem:watch-test", "--debounce", "100", "--producer-only", ], { cwd: path.resolve(import.meta.dirname, ".."), env: { ...process.env, THOUGHTSTREAM_ROOT: project }, stdio: ["ignore", "pipe", "pipe"], }); children.push(child); let stdout = ""; let stderr = ""; child.stdout!.on("data", (chunk) => { stdout += String(chunk); }); child.stderr!.on("data", (chunk) => { stderr += String(chunk); });
await waitUntil(() => stdout.includes("thought stream watching"), () => stderr, 8_000); await fs.writeFile(path.join(vault, "watched.md"), "# One\n"); await fs.writeFile(path.join(vault, "watched.md"), "# Two\n"); await fs.writeFile(path.join(vault, "watched.md"), "# Three\n"); await waitUntil(() => stdout.includes("# Three"), () => stderr, 8_000); expect(stdout.match(/"changed": 1/g)).toHaveLength(1);
child.kill("SIGTERM"); const exit = await waitForExit(child, 8_000); expect(exit).toEqual({ code: 0, signal: null }); expect(stderr).toContain("thought stream watcher stopping."); expect(stderr).toContain("thought stream watcher stopped."); }, 20_000);
test("producer-only mode persists files without loading agent declarations", async () => { const project = await temporaryProject(); roots.push(project); const vault = path.join(project, "vault"); await fs.mkdir(vault); await fs.mkdir(path.join(project, "agents")); await fs.writeFile(path.join(vault, "coil.md"), "# Coil source\n"); await fs.writeFile(path.join(project, "agents", "must-not-load.yaml"), "not: a valid declaration\n"); const child = spawn(process.execPath, [ "--import", "tsx", "src/cli.ts", "watch", "--producer-only", "--root", vault, "--source", "filesystem:coil", ], { cwd: path.resolve(import.meta.dirname, ".."), env: { ...process.env, THOUGHTSTREAM_ROOT: project }, stdio: ["ignore", "pipe", "pipe"], }); children.push(child); let stdout = ""; let stderr = ""; child.stdout!.on("data", (chunk) => { stdout += String(chunk); }); child.stderr!.on("data", (chunk) => { stderr += String(chunk); });
await waitUntil(() => stdout.includes("producer-only mode"), () => stderr, 8_000); expect(stdout).toContain('"added": 1'); child.kill("SIGTERM"); expect(await waitForExit(child, 8_000)).toEqual({ code: 0, signal: null }); }, 20_000);});
async function waitUntil(predicate: () => boolean, diagnostics: () => string, timeoutMs: number): Promise<void> { const deadline = Date.now() + timeoutMs; while (!predicate()) { if (Date.now() >= deadline) throw new Error(`Timed out waiting for watcher output. stderr:\n${diagnostics()}`); await new Promise((resolve) => setTimeout(resolve, 25)); }}
async function waitForExit(child: ChildProcess, timeoutMs: number): Promise<{ code: number | null; signal: NodeJS.Signals | null }> { if (child.exitCode !== null) return { code: child.exitCode, signal: child.signalCode }; return await new Promise((resolve, reject) => { const timeout = setTimeout(() => reject(new Error("Watcher did not exit after SIGTERM")), timeoutMs); child.once("exit", (code, signal) => { clearTimeout(timeout); resolve({ code, signal }); }); });}