Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
12 kB · 300 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301import 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<void> { 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<string, unknown> { 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<string, Set<(event: any) => 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); }}