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 { WebSocketServer } from "ws"; import { temporaryProject, testDeclarationEnvironment } from "./helpers.js"; const roots: string[] = []; const servers: WebSocketServer[] = []; afterEach(async () => { await Promise.all(servers.splice(0).map(async (server) => { for (const client of server.clients) client.terminate(); await new Promise((resolve) => server.close(() => resolve())); })); await Promise.all(roots.splice(0).map((root) => fs.rm(root, { recursive: true, force: true }))); }); describe("thought stream jetstream command", () => { test("runs against an explicit local fixture server, persists events, and exits at its message bound", async () => { const observedUrls: string[] = []; const server = new WebSocketServer({ host: "127.0.0.1", port: 0 }); servers.push(server); await new Promise((resolve, reject) => { server.once("listening", resolve); server.once("error", reject); }); server.on("connection", (socket, request) => { observedUrls.push(request.url ?? ""); socket.send(JSON.stringify(commitMessage(1784042400000000, "one", "rev-one"))); socket.send(JSON.stringify(commitMessage(1784042401000000, "two", "rev-two"))); }); const address = server.address(); if (!address || typeof address === "string") throw new Error("Missing Jetstream CLI fixture server address"); const project = await temporaryProject(); roots.push(project); await fs.cp(path.join(process.cwd(), "agents"), path.join(project, "agents"), { recursive: true }); await fs.cp(path.join(process.cwd(), "prompts"), path.join(project, "prompts"), { recursive: true }); const result = await run([ "--import", "tsx", "src/cli.ts", "jetstream", "--source", "jetstream:cli-fixture", "--collections", "app.bsky.feed.post", "--max-runtime", "5", "--max-messages", "2", "--max-reconnects", "0", ], { THOUGHTSTREAM_ROOT: project, THOUGHTSTREAM_JETSTREAM_URL: `ws://127.0.0.1:${address.port}/subscribe`, }); expect(result).toMatchObject({ code: 0, stderr: "" }); const output = JSON.parse(result.stdout) as { subscription: { reason: string; messages: number; inserted: number; reconnects: number }; producerOnly: boolean; consumerRuns: number; }; expect(output).toEqual({ subscription: expect.objectContaining({ reason: "message-limit", messages: 2, inserted: 2, reconnects: 0 }), producerOnly: false, consumerRuns: expect.any(Number), }); expect(output.consumerRuns).toBeGreaterThanOrEqual(2); const requestUrl = new URL(observedUrls[0]!, `ws://127.0.0.1:${address.port}`); expect(requestUrl.pathname).toBe("/subscribe"); expect(requestUrl.searchParams.getAll("wantedCollections")).toEqual(["app.bsky.feed.post"]); expect(requestUrl.searchParams.has("cursor")).toBe(false); }, 15_000); test("producer-only mode persists Jetstream events without loading or running declarations", async () => { const server = new WebSocketServer({ host: "127.0.0.1", port: 0 }); servers.push(server); await new Promise((resolve, reject) => { server.once("listening", resolve); server.once("error", reject); }); server.on("connection", (socket) => { socket.send(JSON.stringify(commitMessage(1784042400000000, "producer-only", "rev-producer-only"))); }); const address = server.address(); if (!address || typeof address === "string") throw new Error("Missing producer-only fixture server address"); const project = await temporaryProject(); roots.push(project); await fs.mkdir(path.join(project, "agents")); await fs.writeFile(path.join(project, "agents", "must-not-load.yaml"), "invalid: [declaration\n"); const result = await run([ "--import", "tsx", "src/cli.ts", "jetstream", "--producer-only", "--source", "jetstream:producer-only-fixture", "--collections", "app.bsky.feed.post", "--max-runtime", "5", "--max-messages", "1", "--max-reconnects", "0", ], { THOUGHTSTREAM_ROOT: project, THOUGHTSTREAM_JETSTREAM_URL: `ws://127.0.0.1:${address.port}/subscribe`, }); expect(result).toMatchObject({ code: 0, stderr: "" }); expect(JSON.parse(result.stdout)).toEqual({ subscription: expect.objectContaining({ source: "jetstream:producer-only-fixture", reason: "message-limit", messages: 1, inserted: 1, }), producerOnly: true, consumerRuns: 0, }); }, 15_000); }); function commitMessage(timeUs: number, rkey: string, rev: string): Record { return { did: "did:plc:alicefixture", time_us: timeUs, kind: "commit", commit: { rev, operation: "create", collection: "app.bsky.feed.post", rkey, cid: `bafyreifixture-${rkey}`, record: { $type: "app.bsky.feed.post", text: rkey, createdAt: "2026-07-14T15:20:00.000Z" }, }, }; } async function run(args: string[], extraEnv: Record): Promise<{ code: number | null; stdout: string; stderr: string }> { return new Promise((resolve, reject) => { const child = spawn(process.execPath, args, { cwd: path.resolve(import.meta.dirname, ".."), env: { ...process.env, ...testDeclarationEnvironment, ...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 })); }); }