Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
17 kB · 425 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426import { createServer, type IncomingMessage, type Server, type ServerResponse } from "node:http";import { execFile } from "node:child_process";import fs from "node:fs/promises";import path from "node:path";import { promisify } from "node:util";import { afterEach, describe, expect, test } from "vitest";import { sha256 } from "../src/core/json.js";import { XApiClient, assertBoundedXReplayWindow, planXSubscriptions, type XSubscriptionRecord,} from "../src/connectors/x-client.js";import { JazzThoughtStore } from "../src/jazz/store.js";import { temporaryProject } from "./helpers.js";
const roots: string[] = [];const servers: Server[] = [];const execFileAsync = promisify(execFile);
afterEach(async () => { await Promise.all(servers.splice(0).map((server) => new Promise<void>((resolve) => server.close(() => resolve())))); await Promise.all(roots.splice(0).map((root) => fs.rm(root, { recursive: true, force: true })));});
describe("X management client", () => { test("validates exact UTC-minute replay boundaries against the 24-hour horizon", () => { const now = Date.UTC(2026, 7, 10, 12, 0, 59); expect(() => assertBoundedXReplayWindow("202608091200", "202608101200", now)).not.toThrow(); expect(() => assertBoundedXReplayWindow("202608091159", "202608101159", now)).toThrow("replay horizon"); expect(() => assertBoundedXReplayWindow("202608091200", "202608101201", now)).toThrow("future"); expect(() => assertBoundedXReplayWindow("202608091159", "202608101200", now)).toThrow("cannot exceed 24 hours"); expect(() => assertBoundedXReplayWindow("202602300000", "202603010000", now)).toThrow("not a valid UTC minute"); expect(() => assertBoundedXReplayWindow("202608101100", "202608101100", now)).toThrow("must precede"); });
test("uses the documented webhook, subscription, replay, and user endpoints", async () => { const fixture = await xApiFixture(); const client = new XApiClient({ bearerToken: "fixture-bearer", baseUrl: fixture.baseUrl }); expect(await client.listWebhooks()).toEqual([]); const webhook = await client.createWebhook("https://thoughtstream.example/webhooks/x/fixture"); expect(webhook).toMatchObject({ id: "7001", valid: true }); expect(await client.validateWebhook(webhook.id)).toBe(true); expect(await client.lookupUser("@karpathy")).toEqual({ id: "33836629", username: "karpathy", name: "Andrej Karpathy" });
const subscription = await client.createSubscription({ eventType: "post.create", userId: "33836629", tag: "thoughtstream:x:fixture:post-create", webhookId: webhook.id, }); expect(subscription).toMatchObject({ eventType: "post.create", userId: "33836629", webhookId: "7001" }); expect(await client.listSubscriptions()).toHaveLength(1); expect((await client.updateSubscription(subscription.id, { tag: "thoughtstream:x:fixture:post-create-v2", webhookId: webhook.id, })).tag).toBe("thoughtstream:x:fixture:post-create-v2"); expect(await client.createReplay(webhook.id, "202608101000", "202608101100")).toEqual({ jobId: "replay-1", createdAt: "2026-08-10T11:01:00.000Z", }); expect(await client.deleteSubscription(subscription.id)).toBe(true); expect(await client.deleteWebhook(webhook.id)).toBe(true);
expect(fixture.calls.map((call) => `${call.method} ${call.path}`)).toEqual([ "GET /2/webhooks", "POST /2/webhooks", "PUT /2/webhooks/7001", "GET /2/users/by/username/karpathy", "POST /2/activity/subscriptions", "GET /2/activity/subscriptions", "GET /2/activity/subscriptions", `PUT /2/activity/subscriptions/${subscription.id}`, "POST /2/webhooks/replay", `DELETE /2/activity/subscriptions/${subscription.id}`, "DELETE /2/webhooks/7001", ]); });
test("creates and reads back an outbound user-context like subscription with direction in its identity", async () => { const fixture = await xApiFixture(); const client = new XApiClient({ bearerToken: "fixture-bearer", baseUrl: fixture.baseUrl }); const subscription = await client.createSubscription({ eventType: "like.create", userId: "1232326955652931584", direction: "outbound", tag: "thoughtstream:x:cameron-private:just_cameron:like-create-outbound", webhookId: "7001", });
expect(subscription).toMatchObject({ eventType: "like.create", userId: "1232326955652931584", direction: "outbound", webhookId: "7001", }); expect(fixture.calls.find((call) => call.method === "POST" && call.path === "/2/activity/subscriptions")?.body).toMatchObject({ event_type: "like.create", filter: { user_id: "1232326955652931584", direction: "outbound" }, }); });
test("plans only source-owned subscription mutations and fails closed on unmanaged conflicts", () => { const desired = [ { eventType: "post.create" as const, userId: "1", tag: "thoughtstream:x:fixture:create" }, { eventType: "post.delete" as const, userId: "1", tag: "thoughtstream:x:fixture:delete" }, ]; const live: XSubscriptionRecord[] = [ { id: "10", eventType: "post.create", userId: "1", tag: "thoughtstream:x:fixture:create", webhookId: "70" }, { id: "11", eventType: "post.delete", userId: "2", tag: "thoughtstream:x:fixture:stale", webhookId: "70" }, { id: "12", eventType: "post.create", userId: "999", tag: "unmanaged", webhookId: "80" }, ]; const plan = planXSubscriptions("x:fixture", "70", desired, live); expect(plan).toMatchObject({ create: [{ eventType: "post.delete", userId: "1" }], update: [], delete: [{ subscriptionId: "11" }], unchanged: ["10"], conflicts: [], unmanagedCount: 1, }); expect(plan.planHash).toMatch(/^[a-f0-9]{64}$/);
const conflict = planXSubscriptions("x:fixture", "70", desired, [ ...live, { id: "13", eventType: "post.delete", userId: "1", tag: "someone-else", webhookId: "80" }, ]); expect(conflict.conflicts).toEqual(["unmanaged-existing:post.delete:1:none"]); expect(conflict.create).toHaveLength(0);
const directionChange = planXSubscriptions("x:fixture", "70", [{ eventType: "like.create", userId: "1", direction: "outbound", tag: "thoughtstream:x:fixture:like", }], [{ id: "14", eventType: "like.create", userId: "1", direction: "inbound", tag: "thoughtstream:x:fixture:like", webhookId: "70", }]); expect(directionChange).toMatchObject({ create: [{ eventType: "like.create", userId: "1", direction: "outbound" }], delete: [{ subscriptionId: "14" }], unchanged: [], conflicts: [], }); });
test("keeps status and plan read-only and requires exact mutation confirmations in the CLI", async () => { const fixture = await xApiFixture(); const root = await temporaryProject("thoughtstream-x-cli-"); roots.push(root); await fs.writeFile(path.join(root, "thoughtstream.yaml"), xManifest());
await expect(runXCli("x-webhook-register", root, fixture.baseUrl)).rejects.toThrow("--confirm-url-hash"); expect(fixture.calls.filter((call) => call.method !== "GET")).toHaveLength(0);
const register = await runXCli("x-webhook-register", root, fixture.baseUrl, [ "--confirm-url-hash", sha256("https://thoughtstream.example/webhooks/x/fixture"), ]); expect(register.stdout).toContain('"status": "created"'); expect(register.stdout).not.toContain("fixture-bearer");
const mutationCount = fixture.calls.filter((call) => call.method !== "GET").length; const planResult = await runXCli("x-subscriptions-plan", root, fixture.baseUrl); const plan = JSON.parse(planResult.stdout) as { xSubscriptionPlan: { planHash: string; create: unknown[] } }; expect(plan.xSubscriptionPlan.create).toHaveLength(2); expect(fixture.calls.filter((call) => call.method !== "GET")).toHaveLength(mutationCount);
await expect(runXCli("x-subscriptions-apply", root, fixture.baseUrl, [ "--confirm-plan-hash", "wrong", ])).rejects.toThrow("--confirm-plan-hash"); expect(fixture.subscriptions).toHaveLength(0);
const applied = await runXCli("x-subscriptions-apply", root, fixture.baseUrl, [ "--confirm-plan-hash", plan.xSubscriptionPlan.planHash, ]); expect(applied.stdout).toContain('"created": 2'); expect(fixture.subscriptions).toHaveLength(2);
const beforeStatus = fixture.calls.filter((call) => call.method !== "GET").length; const status = await runXCli("x-webhook-status", root, fixture.baseUrl, ["--record-runtime-ready"]); expect(status.stdout).toContain('"valid": true'); expect(fixture.calls.filter((call) => call.method !== "GET")).toHaveLength(beforeStatus); const controlStore = await JazzThoughtStore.open({ projectRoot: root }); try { const source = (await controlStore.listSources()).find((candidate) => candidate.id === "x:fixture"); expect(source).toEqual(expect.objectContaining({ id: "x:fixture", kind: "x-webhook", enabled: true, lastSequence: 0, })); expect(source?.config.control).toMatchObject({ version: 1, source: "x:fixture", kind: "x-webhook", upstream: { registered: true, valid: true, desiredSubscriptionCount: 2, liveSubscriptionCount: 2, subscriptionsConverged: true, }, runtime: { state: "ready", revision: "x-activity-v2-webhook-v2", }, }); } finally { await controlStore.close(); }
for (const subscription of fixture.subscriptions) subscription.tag = `someone-else:${subscription.subscription_id}`; await expect(runXCli("x-webhook-delete", root, fixture.baseUrl, [ "--confirm-webhook-id", "7001", ])).rejects.toThrow("linked subscriptions remain"); expect(fixture.webhooks).toHaveLength(1);
const [fromDate, toDate] = recentReplayWindow(); const replay = await runXCli("x-webhook-replay", root, fixture.baseUrl, [ "--from", fromDate, "--to", toDate, "--confirm-webhook-id", "7001", ]); expect(replay.stdout).toContain('"jobId": "replay-1"'); }, 30_000);});
interface FixtureCall { method: string; path: string; body?: unknown;}
async function xApiFixture() { const calls: FixtureCall[] = []; const webhooks: Array<{ id: string; url: string; valid: boolean; created_at: string }> = []; const subscriptions: Array<{ subscription_id: string; event_type: string; filter: { user_id: string; direction?: string }; tag: string; webhook_id: string; created_at: string; updated_at: string; }> = []; let nextSubscription = 8000; const server = createServer(async (request, response) => { const url = new URL(request.url ?? "/", "http://127.0.0.1"); const body = await jsonBody(request); calls.push({ method: request.method ?? "GET", path: url.pathname, ...(body !== undefined ? { body } : {}) }); if (request.headers.authorization !== "Bearer fixture-bearer") { send(response, 401, { title: "unauthorized" }); return; } if (request.method === "GET" && url.pathname === "/2/webhooks") { send(response, 200, { data: webhooks }); return; } if (request.method === "POST" && url.pathname === "/2/webhooks") { const record = { id: "7001", url: String((body as { url: string }).url), valid: true, created_at: "2026-08-10T10:00:00.000Z", }; webhooks.splice(0, webhooks.length, record); send(response, 200, { data: record }); return; } const webhookMatch = /^\/2\/webhooks\/([0-9]+)$/.exec(url.pathname); if (request.method === "PUT" && webhookMatch) { const record = webhooks.find((item) => item.id === webhookMatch[1]); if (!record) return send(response, 404, {}); record.valid = true; send(response, 200, { data: record }); return; } if (request.method === "DELETE" && webhookMatch) { const index = webhooks.findIndex((item) => item.id === webhookMatch[1]); const deleted = index >= 0; if (deleted) webhooks.splice(index, 1); send(response, 200, { data: { deleted } }); return; } if (request.method === "GET" && url.pathname === "/2/activity/subscriptions") { send(response, 200, { data: subscriptions, meta: { total_subscriptions: subscriptions.length } }); return; } if (request.method === "POST" && url.pathname === "/2/activity/subscriptions") { const input = body as { event_type: string; filter: { user_id: string }; tag: string; webhook_id: string }; const record = { subscription_id: String(++nextSubscription), event_type: input.event_type, filter: input.filter, tag: input.tag, webhook_id: input.webhook_id, created_at: "2026-08-10T10:00:00.000Z", updated_at: "2026-08-10T10:00:00.000Z", }; subscriptions.push(record); send(response, 200, { data: { accepted: true } }); return; } const subscriptionMatch = /^\/2\/activity\/subscriptions\/([0-9]+)$/.exec(url.pathname); if (request.method === "PUT" && subscriptionMatch) { const record = subscriptions.find((item) => item.subscription_id === subscriptionMatch[1]); if (!record) return send(response, 404, {}); const input = body as { tag: string; webhook_id: string }; record.tag = input.tag; record.webhook_id = input.webhook_id; record.updated_at = "2026-08-10T10:01:00.000Z"; send(response, 200, { data: record }); return; } if (request.method === "DELETE" && subscriptionMatch) { const index = subscriptions.findIndex((item) => item.subscription_id === subscriptionMatch[1]); const deleted = index >= 0; if (deleted) subscriptions.splice(index, 1); send(response, 200, { data: { deleted } }); return; } if (request.method === "POST" && url.pathname === "/2/webhooks/replay") { send(response, 200, { data: { job_id: "replay-1", created_at: "2026-08-10T11:01:00.000Z" } }); return; } if (request.method === "GET" && url.pathname === "/2/users/by/username/karpathy") { send(response, 200, { data: { id: "33836629", username: "karpathy", name: "Andrej Karpathy" } }); return; } send(response, 404, { title: "not found" }); }); servers.push(server); await new Promise<void>((resolve) => server.listen(0, "127.0.0.1", resolve)); const address = server.address(); if (!address || typeof address === "string") throw new Error("X API fixture did not bind"); return { baseUrl: `http://127.0.0.1:${address.port}`, calls, webhooks, subscriptions };}
async function jsonBody(request: IncomingMessage): Promise<unknown> { const chunks: Buffer[] = []; for await (const chunk of request) chunks.push(Buffer.from(chunk)); if (chunks.length === 0) return undefined; return JSON.parse(Buffer.concat(chunks).toString("utf8")) as unknown;}
function send(response: ServerResponse, status: number, value: unknown): void { const body = JSON.stringify(value); response.writeHead(status, { "content-type": "application/json", "content-length": Buffer.byteLength(body) }); response.end(body);}
async function runXCli(command: string, root: string, baseUrl: string, extraArgs: string[] = []) { return execFileAsync(process.execPath, [ "--import", "tsx", path.join(process.cwd(), "src/cli.ts"), command, "--source", "x:fixture", ...extraArgs, ], { cwd: process.cwd(), env: { HOME: process.env.HOME, PATH: process.env.PATH, NODE_PATH: process.env.NODE_PATH, THOUGHTSTREAM_ROOT: root, THOUGHTSTREAM_JAZZ_AUTO: "1", THOUGHTSTREAM_X_API_BASE_URL: baseUrl, FIXTURE_X_BEARER: "fixture-bearer", }, maxBuffer: 1024 * 1024, });}
function xManifest(): string { return `version: 1runtime: revision: x-fixturesources: - id: x:fixture kind: x-webhook enabled: true lane: public-watchlist consumerSecretEnv: FIXTURE_X_SECRET managementBearerTokenEnv: FIXTURE_X_BEARER webhookUrl: https://thoughtstream.example/webhooks/x/fixture webhookPath: /webhooks/x/fixture listenHost: 127.0.0.1 listenPort: 4320 maxBodyBytes: 1048576 requestTimeoutMs: 5000 expectedSubscriptions: - eventType: post.create userId: "33836629" tag: thoughtstream:x:fixture:post-create - eventType: post.delete userId: "33836629" tag: thoughtstream:x:fixture:post-delete`;}
function recentReplayWindow(): [string, string] { const to = new Date(Date.now() - 10 * 60_000); const from = new Date(to.getTime() - 60 * 60_000); return [formatXTimestamp(from), formatXTimestamp(to)];}
function formatXTimestamp(value: Date): string { return [ value.getUTCFullYear().toString().padStart(4, "0"), (value.getUTCMonth() + 1).toString().padStart(2, "0"), value.getUTCDate().toString().padStart(2, "0"), value.getUTCHours().toString().padStart(2, "0"), value.getUTCMinutes().toString().padStart(2, "0"), ].join("");}