import fs from "node:fs/promises"; import { afterEach, describe, expect, test } from "vitest"; import { buildJetstreamSubscriptionUrl, subscribeJetstream, type JetstreamWebSocket, type JetstreamWebSocketFactory, } from "../src/connectors/jetstream-live.js"; import { JetstreamConnector } from "../src/connectors/jetstream.js"; import type { JazzThoughtStore } from "../src/jazz/store.js"; import { temporaryProject, testStore } from "./helpers.js"; const stores: JazzThoughtStore[] = []; const roots: string[] = []; afterEach(async () => { await Promise.all(stores.splice(0).map((store) => store.close())); await Promise.all(roots.splice(0).map((root) => fs.rm(root, { recursive: true, force: true }))); }); describe("live Jetstream subscription", () => { test("constructs a bounded filter URL with an optional rewind cursor", () => { const connector = new JetstreamConnector({ id: "jetstream:live-test", collections: ["app.bsky.feed.post", "app.bsky.graph.*"], dids: ["did:plc:alicefixture"], }); const value = buildJetstreamSubscriptionUrl(connector, { endpoint: "ws://127.0.0.1:6008/subscribe?existing=kept", cursor: 1784042400000000, maxMessageSizeBytes: 100_000, }); const url = new URL(value); expect(url.protocol).toBe("ws:"); expect(url.searchParams.getAll("wantedCollections")).toEqual(["app.bsky.feed.post", "app.bsky.graph.*"]); expect(url.searchParams.getAll("wantedDids")).toEqual(["did:plc:alicefixture"]); expect(url.searchParams.get("cursor")).toBe("1784042400000000"); expect(url.searchParams.get("maxMessageSizeBytes")).toBe("100000"); expect(url.searchParams.get("existing")).toBe("kept"); expect(() => buildJetstreamSubscriptionUrl(connector, { endpoint: "https://example.com/subscribe" })) .toThrow("ws: or wss:"); }); test("serially persists live messages, exposes them to a Jazz consumer subscription, and stops at the message limit", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const connector = new JetstreamConnector({ id: "jetstream:live-test", collections: ["app.bsky.feed.post"] }); const receivedEventIds: string[] = []; const unsubscribe = store.subscribeConsumerEvents({ consumerId: "fixture-observer", consumerVersion: 1, source: "jetstream:live-test", eventTypes: ["stream.thought.source.atproto.commit"], acceptedPrivacy: ["public-source"], afterSequence: 0, }, (events) => { receivedEventIds.push(...events.map((event) => event.id)); }); const socketFactory = scriptedSockets([ (socket) => { socket.serverMessage(commitMessage(1784042400000000, "one", "rev-one")); socket.serverMessage({ did: "did:plc:alicefixture", time_us: 1784042401000000, kind: "identity" }); socket.serverMessage(commitMessage(1784042402000000, "two", "rev-two")); }, ]); const result = await subscribeJetstream(store, connector, { endpoint: "ws://fixture.test/subscribe", maxMessages: 3, maxReconnects: 0, webSocketFactory: socketFactory.factory, }); await waitFor(() => receivedEventIds.length === 2); unsubscribe(); expect(result).toMatchObject({ reason: "message-limit", messages: 3, inserted: 2, ignored: 1, unchanged: 0, reconnects: 0, }); expect(receivedEventIds).toHaveLength(2); expect(await store.listEvents({ types: ["stream.thought.source.atproto.commit"] })).toHaveLength(2); expect((await store.getSourceCursor("cursor:jetstream:live-test"))?.cursor).toMatchObject({ timeUs: 1784042402000000 }); expect(await store.listEvents({ types: ["stream.thought.connector.subscription.started"] })).toHaveLength(1); expect(await store.listEvents({ types: ["stream.thought.connector.subscription.connected"] })).toHaveLength(1); expect(await store.listEvents({ types: ["stream.thought.connector.subscription.stopped"] })).toHaveLength(1); expect(new URL(socketFactory.urls[0]!).searchParams.has("cursor")).toBe(false); }); test("rewinds after disconnect and offers overlap to idempotency without losing an equal-timestamp event", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const connector = new JetstreamConnector({ id: "jetstream:reconnect-test", collections: ["app.bsky.feed.post"] }); const timeUs = 1784042400000000; const first = commitMessage(timeUs, "one", "rev-one"); const equalTimestampButUnseen = commitMessage(timeUs, "two", "rev-two"); const socketFactory = scriptedSockets([ (socket) => { socket.serverMessage(first); socket.serverClose(1012, "fixture restart"); }, (socket) => { socket.serverMessage(first); socket.serverMessage(equalTimestampButUnseen); }, ]); const delays: number[] = []; const result = await subscribeJetstream(store, connector, { endpoint: "ws://fixture.test/subscribe", rewindUs: 2_000_000, maxMessages: 3, maxReconnects: 1, initialBackoffMs: 10, maxBackoffMs: 10, random: () => 0.5, sleep: async (milliseconds) => { delays.push(milliseconds); }, webSocketFactory: socketFactory.factory, }); expect(result).toMatchObject({ reason: "message-limit", messages: 3, inserted: 2, unchanged: 1, reconnects: 1, }); expect(delays).toEqual([10]); expect(await store.listEvents({ types: ["stream.thought.source.atproto.commit"] })).toHaveLength(2); expect(await store.listEvents({ types: ["stream.thought.connector.failed"] })).toHaveLength(1); expect(await store.listEvents({ types: ["stream.thought.connector.recovered"] })).toHaveLength(1); expect(new URL(socketFactory.urls[1]!).searchParams.get("cursor")).toBe(String(timeUs - 2_000_000)); expect((await store.getSourceCursor("cursor:jetstream:reconnect-test"))?.cursor).toMatchObject({ timeUs }); }); test("records malformed transport data and reconnects without advancing the cursor", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const connector = new JetstreamConnector({ id: "jetstream:malformed-test", collections: ["app.bsky.feed.post"] }); const socketFactory = scriptedSockets([ (socket) => { socket.serverRaw("{broken"); }, (socket) => { socket.serverMessage(commitMessage(1784042400000000, "one", "rev-one")); }, ]); const result = await subscribeJetstream(store, connector, { endpoint: "ws://fixture.test/subscribe", maxMessages: 2, maxReconnects: 1, initialBackoffMs: 0, maxBackoffMs: 0, random: () => 0.5, sleep: async () => {}, webSocketFactory: socketFactory.factory, }); expect(result).toMatchObject({ reason: "message-limit", messages: 2, inserted: 1, reconnects: 1 }); expect(await store.listEvents({ types: ["stream.thought.connector.failed"] })).toHaveLength(1); expect(await store.listEvents({ types: ["stream.thought.source.atproto.commit"] })).toHaveLength(1); }); test("finishes the current durable message before aborting", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const connector = new JetstreamConnector({ id: "jetstream:abort-test", collections: ["app.bsky.feed.post"] }); const controller = new AbortController(); const socketFactory = scriptedSockets([ (socket) => { socket.serverMessage(commitMessage(1784042400000000, "one", "rev-one")); socket.serverMessage(commitMessage(1784042401000000, "two", "rev-two")); }, ]); // Abort deterministically after the first durable ingest. (alpha53 delivered // consumer subscriptions synchronously, so a subscription callback could // trigger the abort between the two queued messages; alpha55 delivery is // asynchronous through the sync transport, so the trigger wraps ingest.) const originalIngest = connector.ingestBatch.bind(connector); let ingests = 0; connector.ingestBatch = async (...arguments_) => { const batch = await originalIngest(...arguments_); ingests += 1; if (ingests === 1) controller.abort(); return batch; }; const result = await subscribeJetstream(store, connector, { endpoint: "ws://fixture.test/subscribe", signal: controller.signal, maxMessages: 10, maxReconnects: 0, webSocketFactory: socketFactory.factory, }); expect(result).toMatchObject({ reason: "aborted", messages: 1, inserted: 1 }); expect(await store.listEvents({ types: ["stream.thought.source.atproto.commit"] })).toHaveLength(1); expect((await store.getSourceCursor("cursor:jetstream:abort-test"))?.cursor).toMatchObject({ timeUs: 1784042400000000 }); }); }); async function waitFor(predicate: () => boolean, timeoutMs = 2_000): Promise { const deadline = Date.now() + timeoutMs; while (Date.now() < deadline) { if (predicate()) return; await new Promise((resolve) => setTimeout(resolve, 10)); } throw new Error("Timed out waiting for Jazz subscription delivery"); } 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" }, }, }; } function scriptedSockets(scripts: Array<(socket: FakeWebSocket) => void>): { factory: JetstreamWebSocketFactory; urls: string[]; } { const urls: string[] = []; let index = 0; return { urls, factory: (url) => { urls.push(url); const script = scripts[index++]; if (!script) throw new Error(`Unexpected fixture WebSocket connection ${index}`); return new FakeWebSocket(script); }, }; } class FakeWebSocket implements JetstreamWebSocket { readyState = 0; private readonly listeners = new Map void>>(); constructor(script: (socket: FakeWebSocket) => void) { queueMicrotask(() => { this.readyState = 1; this.emit("open", {}); queueMicrotask(() => script(this)); }); } addEventListener(type: string, listener: (event: any) => void): void { const listeners = this.listeners.get(type) ?? new Set(); listeners.add(listener); this.listeners.set(type, listeners); } removeEventListener(type: string, listener: (event: any) => void): void { this.listeners.get(type)?.delete(listener); } close(code = 1000, reason = ""): void { if (this.readyState >= 2) return; this.readyState = 2; queueMicrotask(() => { this.readyState = 3; this.emit("close", { code, reason }); }); } serverMessage(value: unknown): void { this.serverRaw(JSON.stringify(value)); } serverRaw(value: string): void { if (this.readyState !== 1) return; this.emit("message", { data: value }); } serverClose(code: number, reason: string): void { if (this.readyState !== 1) return; this.readyState = 3; this.emit("close", { code, reason }); } private emit(type: string, event: unknown): void { for (const listener of [...(this.listeners.get(type) ?? [])]) listener(event); } }