Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
9.6 kB · 213 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214import fs from "node:fs/promises";import path from "node:path";import { afterEach, describe, expect, test } from "vitest";import { JetstreamConnector, parseJetstreamNdjson } from "../src/connectors/jetstream.js";import type { JazzThoughtStore } from "../src/jazz/store.js";import { loadThoughtStreamManifest } from "../src/runtime/manifest.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("JetstreamConnector", () => { test("normalizes captured create, update, and delete commits with strong references", async () => { const messages = await fixtureMessages(); const connector = new JetstreamConnector({ id: "jetstream:fixture", collections: ["app.bsky.feed.post"] }); const candidates = messages.map((message) => connector.normalize(message, "fixture-batch")).filter(Boolean);
expect(candidates).toHaveLength(3); expect(candidates[0]).toMatchObject({ externalId: "at://did:plc:alicefixture/app.bsky.feed.post/post-one", idempotencyKey: "did:plc:alicefixture:app.bsky.feed.post:post-one:3mfixture001:create", actor: "did:plc:alicefixture", occurredAt: "2026-07-14T15:20:00.000Z", payload: { operation: "create", cid: "bafyreifixturecreate", atUri: "at://did:plc:alicefixture/app.bsky.feed.post/post-one", record: expect.objectContaining({ text: "Evidence before inference." }), }, }); expect(candidates[2]?.payload).toMatchObject({ operation: "delete", rev: "3mfixture003" }); expect(candidates[2]?.payload).not.toHaveProperty("record"); expect(candidates[2]?.payload).not.toHaveProperty("cid"); expect(connector.describe()).toMatchObject({ transport: "captured-batch-and-live-websocket" }); });
test("advances a filter-bound cursor after durable commits and skips replay after restart", async () => { const messages = await fixtureMessages(); const project = await temporaryProject(); roots.push(project); const connector = new JetstreamConnector({ id: "jetstream:fixture", collections: ["app.bsky.feed.post"] });
const firstStore = await testStore(project); const first = await connector.ingestBatch(firstStore, messages); expect(first).toMatchObject({ inserted: 3, unchanged: 0, ignored: 2, stale: 0 }); expect(first.cursor?.cursor).toEqual({ filterRevision: connector.filterRevision, timeUs: 1784042404000000 }); expect(await firstStore.listEvents({ types: ["stream.thought.source.atproto.commit"] })).toHaveLength(3); await firstStore.close();
const restartedStore = await testStore(project); stores.push(restartedStore); const replay = await connector.ingestBatch(restartedStore, messages); expect(replay).toMatchObject({ inserted: 0, unchanged: 0, ignored: 0, stale: 5 }); expect(await restartedStore.listEvents({ types: ["stream.thought.source.atproto.commit"] })).toHaveLength(3);
const firstMessage = messages[0] as Record<string, unknown>; const redelivery = { ...firstMessage, time_us: 1784042405000000 }; const duplicate = await connector.ingestBatch(restartedStore, [redelivery]); expect(duplicate).toMatchObject({ inserted: 0, unchanged: 1, ignored: 0, stale: 0 }); expect(duplicate.cursor?.cursor).toMatchObject({ timeUs: 1784042405000000 }); expect(await restartedStore.listEvents({ types: ["stream.thought.source.atproto.commit"] })).toHaveLength(3); expect(await restartedStore.listEvents({ types: ["stream.thought.connector.cursor.advanced"] })).toHaveLength(2); });
test("fails closed on malformed batches and changed filters without advancing the cursor", async () => { const messages = await fixtureMessages(); const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const connector = new JetstreamConnector({ id: "jetstream:fixture", collections: ["app.bsky.feed.post"] }); await connector.ingestBatch(store, messages.slice(0, 1)); const before = await store.getSourceCursor("cursor:jetstream:fixture");
await expect(connector.ingestBatch(store, [{ did: "not-a-did", time_us: 1784023205000000, kind: "commit" }])) .rejects.toThrow(); expect((await store.getSourceCursor("cursor:jetstream:fixture"))?.cursor).toEqual(before?.cursor);
const changed = new JetstreamConnector({ id: "jetstream:fixture", collections: ["app.bsky.feed.like"] }); await expect(changed.ingestBatch(store, [])).rejects.toThrow("filter revision"); expect((await store.getSourceCursor("cursor:jetstream:fixture"))?.cursor).toEqual(before?.cursor); expect(await store.listEvents({ types: ["stream.thought.connector.failed"] })).toHaveLength(2); });
test("tracks Semble collection links without separately ingesting cards or collections", async () => { const loaded = await loadThoughtStreamManifest(process.cwd()); const tracked = loaded.manifest.sources.find((source) => source.id === "jetstream:cameron-bluesky"); expect(tracked).toMatchObject({ kind: "jetstream", enabled: false, collections: [ "app.bsky.feed.post", "app.bsky.feed.like", "network.cosmik.collectionLink", ], }); if (!tracked || tracked.kind !== "jetstream") throw new Error("Missing tracked Jetstream fixture"); expect(tracked.collections).not.toContain("network.cosmik.card"); expect(tracked.collections).not.toContain("network.cosmik.collection");
const connector = new JetstreamConnector({ id: "jetstream:semble-fixture", collections: tracked.collections, dids: ["did:plc:fixtureowner"], }); expect(connector.describe()).toMatchObject({ createOnlyCollections: ["network.cosmik.collectionLink"], filterRevision: expect.any(String), }); const link = connector.normalize({ did: "did:plc:fixtureowner", time_us: 1784042407000000, kind: "commit", commit: { rev: "rev-link-fixture", operation: "create", collection: "network.cosmik.collectionLink", rkey: "link-fixture", cid: "bafy-fixture-link", record: { $type: "network.cosmik.collectionLink", card: { uri: "at://did:plc:fixturecard/network.cosmik.card/card-fixture", cid: "bafy-fixture-card", }, collection: { uri: "at://did:plc:fixturecollection/network.cosmik.collection/collection-fixture", cid: "bafy-fixture-collection", }, }, }, }, "semble-fixture"); const deleteMessage = { did: "did:plc:fixtureowner", time_us: 1784042407500000, kind: "commit", commit: { rev: "rev-link-delete-fixture", operation: "delete", collection: "network.cosmik.collectionLink", rkey: "link-fixture", }, }; const deletion = connector.normalize(deleteMessage, "semble-fixture"); const card = connector.normalize({ did: "did:plc:fixtureowner", time_us: 1784042408000000, kind: "commit", commit: { rev: "rev-card-fixture", operation: "create", collection: "network.cosmik.card", rkey: "card-fixture", cid: "bafy-fixture-card", record: { $type: "network.cosmik.card" }, }, }, "semble-fixture");
expect(link).toMatchObject({ externalId: "at://did:plc:fixtureowner/network.cosmik.collectionLink/link-fixture", payload: { cid: "bafy-fixture-link", collection: "network.cosmik.collectionLink", operation: "create", record: { card: { uri: "at://did:plc:fixturecard/network.cosmik.card/card-fixture", cid: "bafy-fixture-card", }, collection: { uri: "at://did:plc:fixturecollection/network.cosmik.collection/collection-fixture", cid: "bafy-fixture-collection", }, }, }, }); expect(deletion).toBeUndefined(); expect(card).toBeUndefined();
const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); expect(await connector.ingestBatch(store, [deleteMessage])).toMatchObject({ inserted: 0, ignored: 1, cursor: { cursor: { timeUs: 1784042407500000 } }, }); expect(await store.listEvents({ types: ["stream.thought.source.atproto.commit"] })).toEqual([]); });
test("requires a bounded collection filter", () => { expect(() => new JetstreamConnector({ id: "jetstream:global", collections: [] })).toThrow("collection filter"); expect(() => new JetstreamConnector({ id: "jetstream:bad-prefix", collections: ["app.bsky.fo*"] })).toThrow("exact NSIDs"); const prefix = new JetstreamConnector({ id: "jetstream:prefix", collections: ["app.bsky.feed.*"] }); expect(prefix.normalize({ did: "did:plc:alicefixture", time_us: 1784042406000000, kind: "commit", commit: { rev: "3mfixture005", operation: "create", collection: "app.bsky.feed.post", rkey: "post-two" }, }, "prefix-batch")).toBeDefined(); });});
async function fixtureMessages(): Promise<unknown[]> { const fixture = await fs.readFile(path.join(process.cwd(), "fixtures", "jetstream", "commits.ndjson"), "utf8"); return parseJetstreamNdjson(fixture);}