From 56d56bf9df5be8f8c42e5179a852036d36017fd1 Mon Sep 17 00:00:00 2001 From: Steve Date: Wed, 25 Feb 2026 16:21:59 -0500 Subject: [PATCH] feat: added new jetstream-consumer package Replaces need for tap to index documents --- bun.lock | 25 ++- packages/jetstream-consumer/.gitignore | 34 ++++ packages/jetstream-consumer/CLAUDE.md | 106 ++++++++++ packages/jetstream-consumer/Dockerfile | 13 ++ packages/jetstream-consumer/README.md | 15 ++ packages/jetstream-consumer/cursor.txt | 1 + .../jetstream-consumer/docker-compose.yml | 14 ++ packages/jetstream-consumer/package.json | 18 ++ packages/jetstream-consumer/src/index.ts | 190 ++++++++++++++++++ packages/jetstream-consumer/tsconfig.json | 13 ++ 10 files changed, 427 insertions(+), 2 deletions(-) create mode 100644 packages/jetstream-consumer/.gitignore create mode 100644 packages/jetstream-consumer/CLAUDE.md create mode 100644 packages/jetstream-consumer/Dockerfile create mode 100644 packages/jetstream-consumer/README.md create mode 100644 packages/jetstream-consumer/cursor.txt create mode 100644 packages/jetstream-consumer/docker-compose.yml create mode 100644 packages/jetstream-consumer/package.json create mode 100644 packages/jetstream-consumer/src/index.ts create mode 100644 packages/jetstream-consumer/tsconfig.json diff --git a/bun.lock b/bun.lock index 74f7fc4..e6112a3 100644 --- a/bun.lock +++ b/bun.lock @@ -21,6 +21,17 @@ "vite": "^5.0.0", }, }, + "packages/jetstream-consumer": { + "name": "jetstream-consumer", + "version": "1.0.0", + "devDependencies": { + "@types/bun": "^1.3.9", + "@types/node": "^25.3.0", + }, + "peerDependencies": { + "typescript": "^5", + }, + }, "packages/server": { "name": "@document-feeds/server", "version": "1.0.0", @@ -257,9 +268,11 @@ "@types/babel__traverse": ["@types/babel__traverse@7.28.0", "", { "dependencies": { "@babel/types": "^7.28.2" } }, "sha512-8PvcXf70gTDZBgt9ptxJ8elBeBjcLOAcOtoO/mPJjtji1+CdGbHgm77om1GrsPxsiE+uXIpNSK64UYaIwQXd4Q=="], + "@types/bun": ["@types/bun@1.3.9", "", { "dependencies": { "bun-types": "1.3.9" } }, "sha512-KQ571yULOdWJiMH+RIWIOZ7B2RXQGpL1YQrBtLIV3FqDcCu6FsbFUBwhdKUlCKUpS3PJDsHlJ1QKlpxoVR+xtw=="], + "@types/estree": ["@types/estree@1.0.8", "", {}, "sha512-dWHzHa2WqEXI/O1E9OjrocMTKJl2mSrEolh1Iomrv6U+JuNwaHXsXx9bLu5gG7BUWFIN0skIQJQ/L1rIex4X6w=="], - "@types/node": ["@types/node@25.0.6", "", { "dependencies": { "undici-types": "~7.16.0" } }, "sha512-NNu0sjyNxpoiW3YuVFfNz7mxSQ+S4X2G28uqg2s+CzoqoQjLPsWSbsFFyztIAqt2vb8kfEAsJNepMGPTxFDx3Q=="], + "@types/node": ["@types/node@25.3.0", "", { "dependencies": { "undici-types": "~7.18.0" } }, "sha512-4K3bqJpXpqfg2XKGK9bpDTc6xO/xoUP/RBWS7AtRMug6zZFaRekiLzjVtAoZMquxoAbzBvy5nxQ7veS5eYzf8A=="], "@types/prop-types": ["@types/prop-types@15.7.15", "", {}, "sha512-F6bEyamV9jKGAFBEmlQnesRPGOQqS2+Uwi0Em15xenOxHaf2hv6L8YCVn3rPdPJOiJfPiCnLIRyvwVaqMY3MIw=="], @@ -281,6 +294,8 @@ "browserslist": ["browserslist@4.28.1", "", { "dependencies": { "baseline-browser-mapping": "^2.9.0", "caniuse-lite": "^1.0.30001759", "electron-to-chromium": "^1.5.263", "node-releases": "^2.0.27", "update-browserslist-db": "^1.2.0" }, "bin": { "browserslist": "cli.js" } }, "sha512-ZC5Bd0LgJXgwGqUknZY/vkUQ04r8NXnJZ3yYi4vDmSiZmC/pdSN0NbNRPxZpbtO4uAfDUAFffO8IZoM3Gj8IkA=="], + "bun-types": ["bun-types@1.3.9", "", { "dependencies": { "@types/node": "*" } }, "sha512-+UBWWOakIP4Tswh0Bt0QD0alpTY8cb5hvgiYeWCMet9YukHbzuruIEeXC2D7nMJPB12kbh8C7XJykSexEqGKJg=="], + "caniuse-lite": ["caniuse-lite@1.0.30001764", "", {}, "sha512-9JGuzl2M+vPL+pz70gtMF9sHdMFbY9FJaQBi186cHKH3pSzDvzoUJUPV6fqiKIMyXbud9ZLg4F3Yza1vJ1+93g=="], "color": ["color@4.2.3", "", { "dependencies": { "color-convert": "^2.0.1", "color-string": "^1.9.0" } }, "sha512-1rXeuUUiGGrykh+CeBdu5Ie7OJwinCgQY0bc7GCRxy5xVHy+moaqkpL/jqQq0MtQOeYcrqEz4abc5f0KtU7W4A=="], @@ -333,6 +348,8 @@ "is-arrayish": ["is-arrayish@0.3.4", "", {}, "sha512-m6UrgzFVUYawGBh1dUsWR5M2Clqic9RVXC/9f8ceNlv2IcO9j9J/z8UoCLPqtsPBFNzEpfR3xftohbfqDx8EQA=="], + "jetstream-consumer": ["jetstream-consumer@workspace:packages/jetstream-consumer"], + "js-tokens": ["js-tokens@4.0.0", "", {}, "sha512-RdJUflcE3cUzKiMqQgsCu06FPu9UdIJO0beYbPhHN4k6apgJtifcoCtT9bcxOpYBtpD2kCM6Sbzg4CausW/PKQ=="], "jsesc": ["jsesc@3.1.0", "", { "bin": { "jsesc": "bin/jsesc" } }, "sha512-/sM3dO2FOzXjKQhJuo0Q173wf2KOo8t4I8vHy6lF9poUp7bKT0/NHE8fPX23PwfhnykfqnC2xRxOnVw5XuGIaA=="], @@ -411,7 +428,7 @@ "undici": ["undici@5.29.0", "", { "dependencies": { "@fastify/busboy": "^2.0.0" } }, "sha512-raqeBD6NQK4SkWhQzeYKd1KmIG6dllBOTt55Rmkt4HtI9mwdWtJljnrXjAFUBLTSN67HWrOIZ3EPF4kjUw80Bg=="], - "undici-types": ["undici-types@7.16.0", "", {}, "sha512-Zz+aZWSj8LE6zoxD+xrjh4VfkIG8Ya6LvYkZqtUQGJPZjYl53ypCaUwWqo7eI0x66KBGeRo+mlBEkMSeSZ38Nw=="], + "undici-types": ["undici-types@7.18.2", "", {}, "sha512-AsuCzffGHJybSaRrmr5eHr81mwJU3kjw6M+uprWvCXiNeN9SOGwQ3Jn8jb8m3Z6izVgknn1R0FTCEAP2QrLY/w=="], "unenv": ["unenv@2.0.0-rc.14", "", { "dependencies": { "defu": "^6.1.4", "exsolve": "^1.0.1", "ohash": "^2.0.10", "pathe": "^2.0.3", "ufo": "^1.5.4" } }, "sha512-od496pShMen7nOy5VmVJCnq8rptd45vh6Nx/r2iPbrba6pa6p+tS2ywuIHRZ/OBvSbQZB0kWvpO9XBNVFXHD3Q=="], @@ -437,10 +454,14 @@ "@cspotcode/source-map-support/@jridgewell/trace-mapping": ["@jridgewell/trace-mapping@0.3.9", "", { "dependencies": { "@jridgewell/resolve-uri": "^3.0.3", "@jridgewell/sourcemap-codec": "^1.4.10" } }, "sha512-3Belt6tdc8bPgAtbcmdtNJlirVoTmEb5e2gC94PnkwEW9jI6CAHUeoG85tjWP5WquqfavoMtMwiG4P926ZKKuQ=="], + "bun-types/@types/node": ["@types/node@25.0.6", "", { "dependencies": { "undici-types": "~7.16.0" } }, "sha512-NNu0sjyNxpoiW3YuVFfNz7mxSQ+S4X2G28uqg2s+CzoqoQjLPsWSbsFFyztIAqt2vb8kfEAsJNepMGPTxFDx3Q=="], + "sharp/semver": ["semver@7.7.3", "", { "bin": { "semver": "bin/semver.js" } }, "sha512-SdsKMrI9TdgjdweUSR9MweHA4EJ8YxHn8DFaDisvhVlUOe4BF1tLD7GAj0lIqWVl+dPb/rExr0Btby5loQm20Q=="], "wrangler/esbuild": ["esbuild@0.17.19", "", { "optionalDependencies": { "@esbuild/android-arm": "0.17.19", "@esbuild/android-arm64": "0.17.19", "@esbuild/android-x64": "0.17.19", "@esbuild/darwin-arm64": "0.17.19", "@esbuild/darwin-x64": "0.17.19", "@esbuild/freebsd-arm64": "0.17.19", "@esbuild/freebsd-x64": "0.17.19", "@esbuild/linux-arm": "0.17.19", "@esbuild/linux-arm64": "0.17.19", "@esbuild/linux-ia32": "0.17.19", "@esbuild/linux-loong64": "0.17.19", "@esbuild/linux-mips64el": "0.17.19", "@esbuild/linux-ppc64": "0.17.19", "@esbuild/linux-riscv64": "0.17.19", "@esbuild/linux-s390x": "0.17.19", "@esbuild/linux-x64": "0.17.19", "@esbuild/netbsd-x64": "0.17.19", "@esbuild/openbsd-x64": "0.17.19", "@esbuild/sunos-x64": "0.17.19", "@esbuild/win32-arm64": "0.17.19", "@esbuild/win32-ia32": "0.17.19", "@esbuild/win32-x64": "0.17.19" }, "bin": { "esbuild": "bin/esbuild" } }, "sha512-XQ0jAPFkK/u3LcVRcvVHQcTIqD6E2H1fvZMA5dQPSOWb3suUbWbfbRf94pjc0bNzRYLfIrDRQXr7X+LHIm5oHw=="], + "bun-types/@types/node/undici-types": ["undici-types@7.16.0", "", {}, "sha512-Zz+aZWSj8LE6zoxD+xrjh4VfkIG8Ya6LvYkZqtUQGJPZjYl53ypCaUwWqo7eI0x66KBGeRo+mlBEkMSeSZ38Nw=="], + "wrangler/esbuild/@esbuild/android-arm": ["@esbuild/android-arm@0.17.19", "", { "os": "android", "cpu": "arm" }, "sha512-rIKddzqhmav7MSmoFCmDIb6e2W57geRsM94gV2l38fzhXMwq7hZoClug9USI2pFRGL06f4IOPHHpFNOkWieR8A=="], "wrangler/esbuild/@esbuild/android-arm64": ["@esbuild/android-arm64@0.17.19", "", { "os": "android", "cpu": "arm64" }, "sha512-KBMWvEZooR7+kzY0BtbTQn0OAYY7CsiydT63pVEaPtVYF0hXbUaOyZog37DKxK7NF3XacBJOpYT4adIJh+avxA=="], diff --git a/packages/jetstream-consumer/.gitignore b/packages/jetstream-consumer/.gitignore new file mode 100644 index 0000000..a14702c --- /dev/null +++ b/packages/jetstream-consumer/.gitignore @@ -0,0 +1,34 @@ +# dependencies (bun install) +node_modules + +# output +out +dist +*.tgz + +# code coverage +coverage +*.lcov + +# logs +logs +_.log +report.[0-9]_.[0-9]_.[0-9]_.[0-9]_.json + +# dotenv environment variable files +.env +.env.development.local +.env.test.local +.env.production.local +.env.local + +# caches +.eslintcache +.cache +*.tsbuildinfo + +# IntelliJ based IDEs +.idea + +# Finder (MacOS) folder config +.DS_Store diff --git a/packages/jetstream-consumer/CLAUDE.md b/packages/jetstream-consumer/CLAUDE.md new file mode 100644 index 0000000..764c1dd --- /dev/null +++ b/packages/jetstream-consumer/CLAUDE.md @@ -0,0 +1,106 @@ + +Default to using Bun instead of Node.js. + +- Use `bun ` instead of `node ` or `ts-node ` +- Use `bun test` instead of `jest` or `vitest` +- Use `bun build ` instead of `webpack` or `esbuild` +- Use `bun install` instead of `npm install` or `yarn install` or `pnpm install` +- Use `bun run + + +``` + +With the following `frontend.tsx`: + +```tsx#frontend.tsx +import React from "react"; +import { createRoot } from "react-dom/client"; + +// import .css files directly and it works +import './index.css'; + +const root = createRoot(document.body); + +export default function Frontend() { + return

Hello, world!

; +} + +root.render(); +``` + +Then, run index.ts + +```sh +bun --hot ./index.ts +``` + +For more information, read the Bun API docs in `node_modules/bun-types/docs/**.mdx`. diff --git a/packages/jetstream-consumer/Dockerfile b/packages/jetstream-consumer/Dockerfile new file mode 100644 index 0000000..63a409e --- /dev/null +++ b/packages/jetstream-consumer/Dockerfile @@ -0,0 +1,13 @@ +FROM oven/bun:latest + +WORKDIR /app + +COPY package.json bun.lockb* ./ +RUN bun install --frozen-lockfile || bun install + +COPY src ./src +COPY tsconfig.json ./ + +RUN mkdir -p /app/data + +CMD ["bun", "run", "src/index.ts"] diff --git a/packages/jetstream-consumer/README.md b/packages/jetstream-consumer/README.md new file mode 100644 index 0000000..ea9ec5a --- /dev/null +++ b/packages/jetstream-consumer/README.md @@ -0,0 +1,15 @@ +# jetstream-consumer + +To install dependencies: + +```bash +bun install +``` + +To run: + +```bash +bun run src/index.ts +``` + +This project was created using `bun init` in bun v1.3.5. [Bun](https://bun.com) is a fast all-in-one JavaScript runtime. diff --git a/packages/jetstream-consumer/cursor.txt b/packages/jetstream-consumer/cursor.txt new file mode 100644 index 0000000..2e44e93 --- /dev/null +++ b/packages/jetstream-consumer/cursor.txt @@ -0,0 +1 @@ +1772039237880654 \ No newline at end of file diff --git a/packages/jetstream-consumer/docker-compose.yml b/packages/jetstream-consumer/docker-compose.yml new file mode 100644 index 0000000..302918c --- /dev/null +++ b/packages/jetstream-consumer/docker-compose.yml @@ -0,0 +1,14 @@ +services: + jetstream-consumer: + build: . + environment: + - WEBHOOK_URL=https://example.com/webhook/tap/batch + - WEBHOOK_SECRET= + - JETSTREAM_INSTANCES=jetstream1.us-east.bsky.network,jetstream2.us-east.bsky.network,jetstream1.us-west.bsky.network,jetstream2.us-west.bsky.network + - WANTED_COLLECTIONS=site.standard.document + - CURSOR_FILE=/app/data/cursor.txt + - BATCH_INTERVAL_MS=5000 + - INACTIVITY_TIMEOUT_MS=300000 + volumes: + - ./data:/app/data + restart: unless-stopped diff --git a/packages/jetstream-consumer/package.json b/packages/jetstream-consumer/package.json new file mode 100644 index 0000000..8a88286 --- /dev/null +++ b/packages/jetstream-consumer/package.json @@ -0,0 +1,18 @@ +{ + "name": "jetstream-consumer", + "version": "1.0.0", + "private": true, + "scripts": { + "start": "bun run src/index.ts", + "dev": "bun run src/index.ts" + }, + "devDependencies": { + "@types/bun": "^1.3.9", + "@types/node": "^25.3.0" + }, + "module": "src/index.ts", + "type": "module", + "peerDependencies": { + "typescript": "^5" + } +} diff --git a/packages/jetstream-consumer/src/index.ts b/packages/jetstream-consumer/src/index.ts new file mode 100644 index 0000000..06d02f1 --- /dev/null +++ b/packages/jetstream-consumer/src/index.ts @@ -0,0 +1,190 @@ +export {}; + +const WEBHOOK_URL = process.env.WEBHOOK_URL; +const WEBHOOK_SECRET = process.env.WEBHOOK_SECRET; +const JETSTREAM_INSTANCES = ( + process.env.JETSTREAM_INSTANCES ?? + "jetstream1.us-east.bsky.network,jetstream2.us-east.bsky.network,jetstream1.us-west.bsky.network,jetstream2.us-west.bsky.network" +) + .split(",") + .map((h) => h.trim()); +const WANTED_COLLECTIONS = ( + process.env.WANTED_COLLECTIONS ?? "site.standard.document" +) + .split(",") + .map((c) => c.trim()); +const CURSOR_FILE = process.env.CURSOR_FILE ?? "./cursor.txt"; +const BATCH_INTERVAL_MS = Number(process.env.BATCH_INTERVAL_MS) || 5000; +const INACTIVITY_TIMEOUT_MS = + Number(process.env.INACTIVITY_TIMEOUT_MS) || 300_000; + +const DEFAULT_CURSOR = "1772036746"; + +if (!WEBHOOK_URL) { + console.error("WEBHOOK_URL environment variable is required"); + process.exit(1); +} + +let cursor = DEFAULT_CURSOR; +let currentInstanceIndex = 0; +let ws: WebSocket | null = null; +let lastMessageTime = Date.now(); +let batch: Array<{ + type: string; + did: string; + collection: string; + rkey: string; + cid: string; + record?: Record; +}> = []; + +async function loadCursor(): Promise { + try { + const file = Bun.file(CURSOR_FILE); + const text = await file.text(); + const trimmed = text.trim(); + if (trimmed) { + cursor = trimmed; + console.log(`Loaded cursor from file: ${cursor}`); + } + } catch { + console.log(`No cursor file found, using default: ${cursor}`); + } +} + +async function saveCursor(): Promise { + try { + await Bun.write(CURSOR_FILE, cursor); + } catch (err) { + console.error("Failed to save cursor:", err); + } +} + +async function flushBatch(): Promise { + if (batch.length === 0) return; + + const events = batch.splice(0, batch.length); + console.log(`Flushing batch of ${events.length} events`); + + try { + const headers: Record = { + "Content-Type": "application/json", + }; + if (WEBHOOK_SECRET) { + headers["Authorization"] = `Bearer ${WEBHOOK_SECRET}`; + } + + const res = await fetch(WEBHOOK_URL!, { + method: "POST", + headers, + body: JSON.stringify(events), + }); + + if (!res.ok) { + console.error( + `Webhook responded with ${res.status}: ${await res.text()}`, + ); + return; + } + + const result = (await res.json()) as Record; + console.log("Batch result:", JSON.stringify(result)); + await saveCursor(); + } catch (err) { + console.error("Failed to send batch:", err); + } +} + +function connect(): void { + const host = JETSTREAM_INSTANCES[currentInstanceIndex]; + const collections = WANTED_COLLECTIONS.map( + (c) => `wantedCollections=${c}`, + ).join("&"); + const url = `wss://${host}/subscribe?${collections}&cursor=${cursor}`; + + console.log(`Connecting to ${host} with cursor ${cursor}`); + + ws = new WebSocket(url); + + ws.addEventListener("open", () => { + console.log(`Connected to ${host}`); + lastMessageTime = Date.now(); + }); + + ws.addEventListener("message", (event) => { + lastMessageTime = Date.now(); + + try { + const data = JSON.parse(String(event.data)); + + if (data.kind !== "commit") return; + + cursor = String(data.time_us); + + batch.push({ + type: data.commit.operation, + did: data.did, + collection: data.commit.collection, + rkey: data.commit.rkey, + cid: data.commit.cid, + record: data.commit.record, + }); + } catch (err) { + console.error("Failed to parse message:", err); + } + }); + + ws.addEventListener("close", (event) => { + console.log(`WebSocket closed: code=${event.code} reason=${event.reason}`); + }); + + ws.addEventListener("error", (event) => { + console.error("WebSocket error:", event); + }); +} + +function switchInstance(): void { + currentInstanceIndex = + (currentInstanceIndex + 1) % JETSTREAM_INSTANCES.length; + const newHost = JETSTREAM_INSTANCES[currentInstanceIndex]; + console.log(`Switching to instance: ${newHost}`); + + if (ws) { + try { + ws.close(); + } catch { + // ignore close errors + } + ws = null; + } + + connect(); +} + +// --- Main --- + +await loadCursor(); +connect(); + +// Batch flush interval +setInterval(() => { + flushBatch(); +}, BATCH_INTERVAL_MS); + +// Inactivity check every 30s +setInterval(() => { + const elapsed = Date.now() - lastMessageTime; + if (elapsed > INACTIVITY_TIMEOUT_MS) { + console.log( + `No messages for ${Math.round(elapsed / 1000)}s, triggering failover`, + ); + switchInstance(); + } +}, 30_000); + +console.log("Jetstream consumer started"); +console.log(` Webhook URL: ${WEBHOOK_URL}`); +console.log(` Instances: ${JETSTREAM_INSTANCES.join(", ")}`); +console.log(` Collections: ${WANTED_COLLECTIONS.join(", ")}`); +console.log(` Batch interval: ${BATCH_INTERVAL_MS}ms`); +console.log(` Inactivity timeout: ${INACTIVITY_TIMEOUT_MS}ms`); diff --git a/packages/jetstream-consumer/tsconfig.json b/packages/jetstream-consumer/tsconfig.json new file mode 100644 index 0000000..c54f693 --- /dev/null +++ b/packages/jetstream-consumer/tsconfig.json @@ -0,0 +1,13 @@ +{ + "compilerOptions": { + "target": "ES2022", + "module": "ESNext", + "moduleResolution": "bundler", + "types": ["bun-types"], + "strict": true, + "esModuleInterop": true, + "skipLibCheck": true, + "noEmit": true + }, + "include": ["src"] +} -- 2.51.2