diff --git a/bun.lock b/bun.lock index 7f440f4..bce0e4f 100644 --- a/bun.lock +++ b/bun.lock @@ -6,12 +6,12 @@ "name": "atwiki", "dependencies": { "@atcute/cid": "^2.4.1", + "@atcute/jetstream": "^1.1.2", "@atcute/tid": "^1.1.2", "@atproto/api": "^0.13.0", "@atproto/identity": "^0.4.12", "@atproto/jwk-jose": "^0.1.0", "@atproto/oauth-client-node": "^0.1.0", - "@atproto/sync": "^0.1.40", "@codemirror/lang-markdown": "^6.5.0", "@codemirror/state": "^6.0.0", "@elysiajs/static": "^1.4.7", @@ -23,10 +23,10 @@ "katex": "^0.16.45", "markdown-it": "^14.1.1", "sharp": "^0.34.5", - "ws": "^8.0.0", }, "devDependencies": { "@atproto/dev-env": "0.3.213", + "@atproto/sync": "^0.1.40", "@biomejs/biome": "^2.4.4", "@tailwindcss/cli": "^4.2.1", "@tailwindcss/typography": "^0.5.16", @@ -38,6 +38,7 @@ "better-sqlite3": "12.8.0", "knip": "^6.0.2", "tailwindcss": "^4.2.1", + "ws": "^8.0.0", }, "peerDependencies": { "typescript": "^5", @@ -55,6 +56,10 @@ "packages": { "@atcute/cid": ["@atcute/cid@2.4.1", "", { "dependencies": { "@atcute/multibase": "^1.1.8", "@atcute/uint8array": "^1.1.1" } }, "sha512-bwhna69RCv7yetXudtj+2qrMPYvhhIQqvJz6YUpUS98v7OdF3X2dnye9Nig2NDrklZcuyOsu7sQo7GOykJXRLQ=="], + "@atcute/jetstream": ["@atcute/jetstream@1.1.2", "", { "dependencies": { "@atcute/lexicons": "^1.2.2", "@badrap/valita": "^0.4.6", "@mary-ext/event-iterator": "^1.0.0", "@mary-ext/simple-event-emitter": "^1.0.0", "partysocket": "^1.1.5", "type-fest": "^4.41.0", "yocto-queue": "^1.2.1" } }, "sha512-u6p/h2xppp7LE6W/9xErAJ6frfN60s8adZuCKtfAaaBBiiYbb1CfpzN8Uc+2qtJZNorqGvuuDb5572Jmh7yHBQ=="], + + "@atcute/lexicons": ["@atcute/lexicons@1.3.0", "", { "dependencies": { "@atcute/uint8array": "^1.1.1", "@atcute/util-text": "^1.2.0", "@standard-schema/spec": "^1.1.0", "esm-env": "^1.2.2" } }, "sha512-Eq5y+9onnCXNVUlNiMf31beSXHKqptB7lUo/68YbhlmxdaR7ooywHmahya9goP5AsmlYEA1z+dRPXIDAa9O7cg=="], + "@atcute/multibase": ["@atcute/multibase@1.2.0", "", { "dependencies": { "@atcute/uint8array": "^1.1.1" } }, "sha512-ZK2GRra+qIYq9nNuQB52m2ul0hOmCQEtPobGfTSUxm7pF0OGEkWGkWHugFhNEDVzHzTwPxHp6VGotdZFue4lYQ=="], "@atcute/tid": ["@atcute/tid@1.1.2", "", { "dependencies": { "@atcute/time-ms": "^1.2.2" } }, "sha512-bmPuOX/TOfcm/vsK9vM98spjkcx2wgd9S2PeK5oLgEr8IbNRPq7iMCAPzOL1nu5XAW3LlkOYQEbYRcw5vcQ37w=="], @@ -63,6 +68,8 @@ "@atcute/uint8array": ["@atcute/uint8array@1.1.1", "", {}, "sha512-3LsC8XB8TKe9q/5hOA5sFuzGaIFdJZJNewC5OKa3o/eU6+K7JR6see9Zy2JbQERNVnRl11EzbNov1efgLMAs4g=="], + "@atcute/util-text": ["@atcute/util-text@1.3.1", "", { "dependencies": { "unicode-segmenter": "^0.14.5" } }, "sha512-MRgJXkx67znuBXuoAYCJkBZyd3OApL7zZlNf5kXhuoCXcdiu1nblRDycYTADSkym4epBSQWxh26kmI9sewaq6A=="], + "@atproto-labs/did-resolver": ["@atproto-labs/did-resolver@0.1.4", "", { "dependencies": { "@atproto-labs/fetch": "0.1.1", "@atproto-labs/pipe": "0.1.0", "@atproto-labs/simple-store": "0.1.1", "@atproto-labs/simple-store-memory": "0.1.1", "@atproto/did": "0.1.2", "zod": "^3.23.8" } }, "sha512-5d+LHScS2ueYsFRjMOC3c1EwM2ui1yBVbBA0yY3MH7aydbljm5D28scsOVuymIhHwPFwcGvZbMON4PVSfpBbbQ=="], "@atproto-labs/fetch": ["@atproto-labs/fetch@0.1.1", "", { "dependencies": { "@atproto-labs/pipe": "0.1.0" }, "optionalDependencies": { "zod": "^3.23.8" } }, "sha512-X1zO1MDoJzEurbWXMAe1H8EZ995Xam/aXdxhGVrXmOMyPDuvBa1oxwh/kQNZRCKcMQUbiwkk+Jfq6ZkTuvGbww=="], @@ -245,6 +252,8 @@ "@aws/lambda-invoke-store": ["@aws/lambda-invoke-store@0.2.4", "", {}, "sha512-iY8yvjE0y651BixKNPgmv1WrQc+GZ142sb0z4gYnChDDY2YqI4P/jsSopBWrKfAt7LOJAkOXt7rC/hms+WclQQ=="], + "@badrap/valita": ["@badrap/valita@0.4.6", "", {}, "sha512-4kdqcjyxo/8RQ8ayjms47HCWZIF5981oE5nIenbfThKDxWXtEHKipAOWlflpPJzZx9y/JWYQkp18Awr7VuepFg=="], + "@biomejs/biome": ["@biomejs/biome@2.4.4", "", { "optionalDependencies": { "@biomejs/cli-darwin-arm64": "2.4.4", "@biomejs/cli-darwin-x64": "2.4.4", "@biomejs/cli-linux-arm64": "2.4.4", "@biomejs/cli-linux-arm64-musl": "2.4.4", "@biomejs/cli-linux-x64": "2.4.4", "@biomejs/cli-linux-x64-musl": "2.4.4", "@biomejs/cli-win32-arm64": "2.4.4", "@biomejs/cli-win32-x64": "2.4.4" }, "bin": { "biome": "bin/biome" } }, "sha512-tigwWS5KfJf0cABVd52NVaXyAVv4qpUXOWJ1rxFL8xF1RVoeS2q/LK+FHgYoKMclJCuRoCWAPy1IXaN9/mS61Q=="], "@biomejs/cli-darwin-arm64": ["@biomejs/cli-darwin-arm64@2.4.4", "", { "os": "darwin", "cpu": "arm64" }, "sha512-jZ+Xc6qvD6tTH5jM6eKX44dcbyNqJHssfl2nnwT6vma6B1sj7ZLTGIk6N5QwVBs5xGN52r3trk5fgd3sQ9We9A=="], @@ -411,6 +420,10 @@ "@marijn/find-cluster-break": ["@marijn/find-cluster-break@1.0.2", "", {}, "sha512-l0h88YhZFyKdXIFNfSWpyjStDjGHwZ/U7iobcK1cQQD8sejsONdQtTVU+1wVN1PBw40PiiHB1vA5S7VTfQiP9g=="], + "@mary-ext/event-iterator": ["@mary-ext/event-iterator@1.0.0", "", { "dependencies": { "yocto-queue": "^1.2.1" } }, "sha512-l6gCPsWJ8aRCe/s7/oCmero70kDHgIK5m4uJvYgwEYTqVxoBOIXbKr5tnkLqUHEg6mNduB4IWvms3h70Hp9ADQ=="], + + "@mary-ext/simple-event-emitter": ["@mary-ext/simple-event-emitter@1.0.1", "", {}, "sha512-9+VvZisxZ/gSg+JJH7hmXaA8Qj42Qjz3O58RSB+INYc8iLA0icATZxHB9vKbj59ojDGZjO3hCKzMXocx3L0H8w=="], + "@napi-rs/wasm-runtime": ["@napi-rs/wasm-runtime@1.1.1", "", { "dependencies": { "@emnapi/core": "^1.7.1", "@emnapi/runtime": "^1.7.1", "@tybys/wasm-util": "^0.10.1" } }, "sha512-p64ah1M1ld8xjWv3qbvFwHiFVWrq1yFvV4f7w+mzaqiR4IlSgkqhcRdHwsGgomwzBH51sRY4NEowLxnaBjcW/A=="], "@noble/curves": ["@noble/curves@1.9.7", "", { "dependencies": { "@noble/hashes": "1.8.0" } }, "sha512-gbKGcRUYIjA3/zCCNaWDciTMFI0dCkvou3TL8Zmy5Nc7sJ47a0jtOeZoTaMxkuqRo9cRhjOdZJXegxYE5FN/xw=="], @@ -665,6 +678,8 @@ "@smithy/uuid": ["@smithy/uuid@1.1.2", "", { "dependencies": { "tslib": "^2.6.2" } }, "sha512-O/IEdcCUKkubz60tFbGA7ceITTAJsty+lBjNoorP4Z6XRqaFb/OjQjZODophEcuq68nKm6/0r+6/lLQ+XVpk8g=="], + "@standard-schema/spec": ["@standard-schema/spec@1.1.0", "", {}, "sha512-l2aFy5jALhniG5HgqrD6jXLi/rUWrKvqN/qJx6yoJsgKhblVd+iqqU4RCXavm/jPityDo5TCvKMnpjKnOriy0w=="], + "@tailwindcss/cli": ["@tailwindcss/cli@4.2.1", "", { "dependencies": { "@parcel/watcher": "^2.5.1", "@tailwindcss/node": "4.2.1", "@tailwindcss/oxide": "4.2.1", "enhanced-resolve": "^5.19.0", "mri": "^1.2.0", "picocolors": "^1.1.1", "tailwindcss": "4.2.1" }, "bin": { "tailwindcss": "dist/index.mjs" } }, "sha512-b7MGn51IA80oSG+7fuAgzfQ+7pZBgjzbqwmiv6NO7/+a1sev32cGqnwhscT7h0EcAvMa9r7gjRylqOH8Xhc4DA=="], "@tailwindcss/node": ["@tailwindcss/node@4.2.1", "", { "dependencies": { "@jridgewell/remapping": "^2.3.5", "enhanced-resolve": "^5.19.0", "jiti": "^2.6.1", "lightningcss": "1.31.1", "magic-string": "^0.30.21", "source-map-js": "^1.2.1", "tailwindcss": "4.2.1" } }, "sha512-jlx6sLk4EOwO6hHe1oCGm1Q4AN/s0rSrTTPBGPM0/RQ6Uylwq17FuU8IeJJKEjtc6K6O07zsvP+gDO6MMWo7pg=="], @@ -975,10 +990,14 @@ "escape-string-regexp": ["escape-string-regexp@5.0.0", "", {}, "sha512-/veY75JbMK4j1yjvuUxuVsiS/hr/4iHs9FTT6cgTexxdE0Ly/glccBAkloH/DofkjRbZU3bnoj38mOmhkZ0lHw=="], + "esm-env": ["esm-env@1.2.2", "", {}, "sha512-Epxrv+Nr/CaL4ZcFGPJIYLWFom+YeV1DqMLHJoEd9SYRxNbaFruBwfEX/kkHUJf55j2+TUbmDcmuilbP1TmXHA=="], + "etag": ["etag@1.8.1", "", {}, "sha512-aIL5Fx7mawVa300al2BnEE4iNvo1qETxLrPI/o05L7z6go7fCw1J6EQmbK4FmJ2AS7kgVF/KEZWufBfdClMcPg=="], "etcd3": ["etcd3@1.1.2", "", { "dependencies": { "@grpc/grpc-js": "^1.8.20", "@grpc/proto-loader": "^0.7.8", "bignumber.js": "^9.1.1", "cockatiel": "^3.1.1" } }, "sha512-YIampCz1/OmrVo/tR3QltAVUtYCQQOSFoqmHKKeoHbalm+WdXe3l4rhLIylklu8EzR/I3PBiOF4dC847dDskKg=="], + "event-target-polyfill": ["event-target-polyfill@0.0.4", "", {}, "sha512-Gs6RLjzlLRdT8X9ZipJdIZI/Y6/HhRLyq9RdDlCsnpxr/+Nn6bU2EFGuC94GjxqhM+Nmij2Vcq98yoHrU8uNFQ=="], + "event-target-shim": ["event-target-shim@5.0.1", "", {}, "sha512-i/2XbnSz/uxRCU6+NdVJgKWDTM427+MqYbkQzD321DuCQJUqOuJKIA0IM2+W2xtYHdKOmZ4dR6fExsd4SXL+WQ=="], "eventemitter3": ["eventemitter3@4.0.7", "", {}, "sha512-8guHBZCwKnFhYdHr2ysuRWErTwhoN2X8XELRlrRwpmfeY2jjuUN4taQMsULKUVo1K4DvZl+0pgfyoysHxvmvEw=="], @@ -1253,6 +1272,8 @@ "parseurl": ["parseurl@1.3.3", "", {}, "sha512-CiyeOxFT/JZyN5m0z9PfXw4SCBJ6Sygz1Dpl0wqjlhDEGGBP1GnsUVEL0p63hoG1fcj3fHynXi9NYO4nWOL+qQ=="], + "partysocket": ["partysocket@1.1.18", "", { "dependencies": { "event-target-polyfill": "^0.0.4" }, "peerDependencies": { "react": ">=17" }, "optionalPeers": ["react"] }, "sha512-SyuvH9VavWOSa14v6dYdp3yfSUDII4BQB1+TkGOFBkjfZKjnDBiba4fhdhwBlqGBkqw4ea3gTA1DYhSffX24Wg=="], + "path-expression-matcher": ["path-expression-matcher@1.1.3", "", {}, "sha512-qdVgY8KXmVdJZRSS1JdEPOKPdTiEK/pi0RkcT2sw1RhXxohdujUlJFPuS1TSkevZ9vzd3ZlL7ULl1MHGTApKzQ=="], "path-key": ["path-key@3.1.1", "", {}, "sha512-ojmeN0qd+y0jszEtoY48r0Peq5dwMEkIlCOu6Q5f41lfkswXuKtYrhgoTpLnyIcHm24Uhqx+5Tqm2InSwLhE6Q=="], @@ -1457,7 +1478,7 @@ "tunnel-agent": ["tunnel-agent@0.6.0", "", { "dependencies": { "safe-buffer": "^5.0.1" } }, "sha512-McnNiV1l8RYeY8tBgEpuodCC1mLUdbSN+CYBL7kJsJNInOP8UjDDEwdk6Mw60vdLLrr5NHKZhMAOSrR2NZuQ+w=="], - "type-fest": ["type-fest@2.19.0", "", {}, "sha512-RAH822pAdBgcNMAfWnCBU3CFZcfZ/i1eZjwFU/dsLKumyuuP3niueg2UAukXYF0E2AAoc82ZSSf9J0WQBinzHA=="], + "type-fest": ["type-fest@4.41.0", "", {}, "sha512-TeTSQ6H5YHvpqVwBRcnLDCBnDOHWYu7IvGbHT6N8AOymcr9PJGjc1GTtiWZTYg0NCgYwvnYWEkVChQAr9bjfwA=="], "type-is": ["type-is@1.6.18", "", { "dependencies": { "media-typer": "0.3.0", "mime-types": "~2.1.24" } }, "sha512-TkRKr9sUTxEH8MdfuCSP7VizJyzRNMjj2J2do2Jr3Kym598JVdEksuzPQCnlFPW4ky9Q+iA+ma9BGm06XQBy8g=="], @@ -1517,6 +1538,8 @@ "yargs-parser": ["yargs-parser@21.1.1", "", {}, "sha512-tVpsJW7DdjecAiFpbIB1e3qxIQsE6NoPc5/eTdrbbIC4h0LVsWhnoa3g+m2HclBIujHzsxZ4VJVA+GUuc2/LBw=="], + "yocto-queue": ["yocto-queue@1.2.2", "", {}, "sha512-4LCcse/U2MHZ63HAJVE+v71o7yOdIe4cZ70Wpf8D/IyjDKYQLV5GD46B+hSTjJsvV5PztjvHoU580EftxjDZFQ=="], + "zod": ["zod@4.3.6", "", {}, "sha512-rftlrkhHZOcjDwkGlnUtZZkvaPHCsDATp4pGpuOOMDaTdDDXF91wuVDJoWoPsKX/3YPQ5fHuF3STjcYyKr+Qhg=="], "@atproto-labs/did-resolver/@atproto-labs/simple-store-memory": ["@atproto-labs/simple-store-memory@0.1.1", "", { "dependencies": { "@atproto-labs/simple-store": "0.1.1", "lru-cache": "^10.2.0" } }, "sha512-PCRqhnZ8NBNBvLku53O56T0lsVOtclfIrQU/rwLCc4+p45/SBPrRYNBi6YFq5rxZbK6Njos9MCmILV/KLQxrWA=="], @@ -1699,6 +1722,8 @@ "htmlparser2/entities": ["entities@2.2.0", "", {}, "sha512-p92if5Nz619I0w+akJrLZH0MX0Pb5DX39XOwQTtXSdQQOaYH03S1uIQp4mhOZtAXrxq4ViO67YTiLBo2638o9A=="], + "http-terminator/type-fest": ["type-fest@2.19.0", "", {}, "sha512-RAH822pAdBgcNMAfWnCBU3CFZcfZ/i1eZjwFU/dsLKumyuuP3niueg2UAukXYF0E2AAoc82ZSSf9J0WQBinzHA=="], + "ioredis/debug": ["debug@4.4.3", "", { "dependencies": { "ms": "^2.1.3" } }, "sha512-RGwwWnwQvkVfavKVt22FGLw+xYSdzARwm0ru6DhTVA3umU5hZc28V3kO4stgYryrTlLpuvgI9GiijltAjNbcqA=="], "micromatch/picomatch": ["picomatch@2.3.1", "", {}, "sha512-JU3teHTNjmE2VCGFzuY8EXzCDVwEqB2a8fsIvwaStHhAWJEeVd1o1QD80CU6+ZdEXXSLbSsuLwJjkCBWqRQUVA=="], diff --git a/package.json b/package.json index 05ee9b1..725244e 100644 --- a/package.json +++ b/package.json @@ -25,6 +25,7 @@ }, "devDependencies": { "@atproto/dev-env": "0.3.213", + "@atproto/sync": "^0.1.40", "@biomejs/biome": "^2.4.4", "@tailwindcss/cli": "^4.2.1", "@tailwindcss/typography": "^0.5.16", @@ -35,7 +36,8 @@ "@typescript/native-preview": "^7.0.0-dev.20260322.1", "better-sqlite3": "12.8.0", "knip": "^6.0.2", - "tailwindcss": "^4.2.1" + "tailwindcss": "^4.2.1", + "ws": "^8.0.0" }, "peerDependencies": { "typescript": "^5" @@ -50,12 +52,12 @@ }, "dependencies": { "@atcute/cid": "^2.4.1", + "@atcute/jetstream": "^1.1.2", "@atcute/tid": "^1.1.2", "@atproto/api": "^0.13.0", "@atproto/identity": "^0.4.12", "@atproto/jwk-jose": "^0.1.0", "@atproto/oauth-client-node": "^0.1.0", - "@atproto/sync": "^0.1.40", "@codemirror/lang-markdown": "^6.5.0", "@codemirror/state": "^6.0.0", "@elysiajs/static": "^1.4.7", @@ -66,7 +68,6 @@ "fflate": "^0.8.2", "katex": "^0.16.45", "markdown-it": "^14.1.1", - "sharp": "^0.34.5", - "ws": "^8.0.0" + "sharp": "^0.34.5" } } diff --git a/src/atproto/env.ts b/src/atproto/env.ts index 3ce355d..02ac03e 100644 --- a/src/atproto/env.ts +++ b/src/atproto/env.ts @@ -18,8 +18,10 @@ export function isAuthEnabled(): boolean { return getAtprotoEnv() !== null; } -export function getRelayUrl(): string { - return process.env["RELAY_URL"] ?? "wss://bsky.network"; +export function getJetstreamUrl(): string { + return ( + process.env["JETSTREAM_URL"] ?? "wss://jetstream2.us-east.bsky.network" + ); } export function getHandleResolverUrl(): string { diff --git a/src/firehose/handlers.ts b/src/firehose/handlers.ts index 15a8877..3c4e8c3 100644 --- a/src/firehose/handlers.ts +++ b/src/firehose/handlers.ts @@ -1,5 +1,3 @@ -import type { CommitEvt } from "@atproto/sync"; -import { didOwnsUri } from "../lib/at-uri.ts"; import { COLLECTIONS, normalizeRole } from "../lib/constants.ts"; import { LIMITS } from "../lib/limits.ts"; import { @@ -197,11 +195,24 @@ function isBookmarkRecord(r: Rec): r is Rec & BookmarkRecord { // --- Event dispatch --- -export function handleCommitEvent(evt: CommitEvt): void { - const atUri = evt.uri.toString(); - if (!didOwnsUri(evt.did, atUri)) return; +/** + * Normalized commit event consumed by the handler. Shaped to be easy to build + * from either a jetstream message (production) or an `@atproto/sync` commit + * (integration tests against a local PDS). + */ +export interface FirehoseCommit { + did: string; + collection: string; + rkey: string; + operation: "create" | "update" | "delete"; + /** The record body. Required for create/update; ignored for delete. */ + record?: unknown; +} + +export function handleCommitEvent(evt: FirehoseCommit): void { + const atUri = `at://${evt.did}/${evt.collection}/${evt.rkey}`; - if (evt.event === "delete") { + if (evt.operation === "delete") { handleDelete(atUri, evt.collection); return; } @@ -236,7 +247,7 @@ export function handleCommitEvent(evt: CommitEvt): void { } } -// --- Handlers (validation already done by type guards + didOwnsUri) --- +// --- Handlers (validation already done by type guards) --- function handleWiki( did: string, diff --git a/src/firehose/index.ts b/src/firehose/index.ts index 32b1663..6677feb 100644 --- a/src/firehose/index.ts +++ b/src/firehose/index.ts @@ -1,66 +1,77 @@ -import "../lib/ws-polyfill.ts"; - -import { type Event, Firehose, MemoryRunner } from "@atproto/sync"; -import { getDevPdsUrl, getRelayUrl } from "../atproto/env.ts"; +import { JetstreamSubscription } from "@atcute/jetstream"; +import { getDevPdsUrl, getJetstreamUrl } from "../atproto/env.ts"; import { COLLECTIONS } from "../lib/constants.ts"; -import { getIdResolver } from "../lib/identity.ts"; import { getCursor, setCursor } from "../server/db/queries/index.ts"; -import { handleCommitEvent } from "./handlers.ts"; +import { type FirehoseCommit, handleCommitEvent } from "./handlers.ts"; -const relayUrl = getRelayUrl(); -// In dev mode the PDS is ephemeral -- each run starts fresh from seq 0. -// Persisting the cursor across sessions causes "FutureCursor" errors. +const jetstreamUrl = getJetstreamUrl(); +// In dev mode the PDS is ephemeral -- each run starts fresh. +// Persisting the cursor across sessions causes the subscriber to ask jetstream +// for events older than its retention window. const isDevMode = !!getDevPdsUrl(); const savedCursor = isDevMode ? null : getCursor(); -console.log(`Firehose connecting to ${relayUrl}`); +console.log(`Jetstream connecting to ${jetstreamUrl}`); if (savedCursor) { console.log(`Resuming from cursor ${savedCursor}`); } -const idResolver = getIdResolver(); - -const runner = new MemoryRunner({ - ...(savedCursor != null && { startCursor: savedCursor }), - setCursor: async (cursor: number) => { - if (!isDevMode) setCursor(cursor); - }, +const subscription = new JetstreamSubscription({ + url: jetstreamUrl, + wantedCollections: Object.values(COLLECTIONS), + ...(savedCursor != null && { cursor: savedCursor }), }); -const firehose = new Firehose({ - idResolver, - runner, - service: relayUrl, - filterCollections: Object.values(COLLECTIONS), - excludeIdentity: true, - excludeAccount: true, - handleEvent: (evt: Event) => { - if ( - evt.event === "create" || - evt.event === "update" || - evt.event === "delete" - ) { - try { - handleCommitEvent(evt); - } catch (err) { - console.error(`Error handling ${evt.collection} ${evt.event}:`, err); - } +// Persist the cursor every 5s so a restart resumes near where we left off. +// Jetstream cursors are microsecond timestamps. +const cursorInterval = setInterval(() => { + if (!isDevMode && subscription.cursor != null) { + setCursor(subscription.cursor); + } +}, 5_000); + +let shuttingDown = false; + +async function run(): Promise { + for await (const event of subscription) { + if (event.kind !== "commit") continue; + + const commit: FirehoseCommit = { + did: event.did, + collection: event.commit.collection, + rkey: event.commit.rkey, + operation: event.commit.operation, + record: + event.commit.operation !== "delete" ? event.commit.record : undefined, + }; + + try { + handleCommitEvent(commit); + } catch (err) { + console.error( + `Error handling ${commit.collection} ${commit.operation}:`, + err, + ); } - }, - onError: (err: Error) => { - console.error("Firehose error:", err); - }, -}); + } +} -firehose.start(); -console.log("Firehose subscriber running"); +run().catch((err) => { + if (!shuttingDown) console.error("Jetstream error:", err); +}); +console.log("Jetstream subscriber running"); -function shutdown() { - console.log("Shutting down firehose..."); - firehose.destroy().then(() => { - console.log("Firehose stopped"); - process.exit(0); - }); +function shutdown(): void { + if (shuttingDown) return; + shuttingDown = true; + console.log("Shutting down jetstream..."); + clearInterval(cursorInterval); + if (!isDevMode && subscription.cursor != null) { + setCursor(subscription.cursor); + } + // JetstreamSubscription closes its WebSocket when the async iterator exits; + // process.exit triggers that via teardown. + process.exit(0); } process.on("SIGINT", shutdown); diff --git a/tests/atproto/env.test.ts b/tests/atproto/env.test.ts index 9bae126..633c506 100644 --- a/tests/atproto/env.test.ts +++ b/tests/atproto/env.test.ts @@ -5,7 +5,7 @@ import { getDevPdsUrl, getDevPlcUrl, getHandleResolverUrl, - getRelayUrl, + getJetstreamUrl, isAuthEnabled, } from "../../src/atproto/env.ts"; @@ -14,7 +14,7 @@ const origEnv: Record = {}; const ENV_KEYS = [ "PUBLIC_URL", "OAUTH_PRIVATE_KEY_PATH", - "RELAY_URL", + "JETSTREAM_URL", "HANDLE_RESOLVER_URL", "DEV_PDS_URL", "DEV_PLC_URL", @@ -75,15 +75,15 @@ describe("isAuthEnabled", () => { }); }); -describe("getRelayUrl", () => { - test("defaults to bsky.network", () => { - delete process.env["RELAY_URL"]; - expect(getRelayUrl()).toBe("wss://bsky.network"); +describe("getJetstreamUrl", () => { + test("defaults to a Bluesky jetstream instance", () => { + delete process.env["JETSTREAM_URL"]; + expect(getJetstreamUrl()).toBe("wss://jetstream2.us-east.bsky.network"); }); test("reads from env", () => { - process.env["RELAY_URL"] = "wss://custom.relay"; - expect(getRelayUrl()).toBe("wss://custom.relay"); + process.env["JETSTREAM_URL"] = "wss://custom.jetstream"; + expect(getJetstreamUrl()).toBe("wss://custom.jetstream"); }); }); diff --git a/tests/firehose/handlers.test.ts b/tests/firehose/handlers.test.ts index 865eb31..7f9baea 100644 --- a/tests/firehose/handlers.test.ts +++ b/tests/firehose/handlers.test.ts @@ -1,6 +1,8 @@ import { afterAll, beforeAll, describe, expect, test } from "bun:test"; -import type { CommitEvt } from "@atproto/sync"; -import { handleCommitEvent } from "../../src/firehose/handlers.ts"; +import { + type FirehoseCommit, + handleCommitEvent, +} from "../../src/firehose/handlers.ts"; import { getDb } from "../../src/server/db/index.ts"; import { getCurrentNote, @@ -16,33 +18,21 @@ const BOB_DID = "did:plc:bob"; const WIKI_AT_URI = `at://${ALICE_DID}/wiki.lichen.wiki/test-wiki`; interface CommitEvtInput { - event: string; + event: "create" | "update" | "delete"; collection: string; rkey: string; did?: string; - uri?: { toString(): string }; record?: Record; } -function makeCommitEvt(input: CommitEvtInput): CommitEvt { - const did = input.did ?? ALICE_DID; - const atUri = - input.uri?.toString() ?? `at://${did}/${input.collection}/${input.rkey}`; - +function makeCommitEvt(input: CommitEvtInput): FirehoseCommit { return { - seq: 1, - time: new Date().toISOString(), - commit: {} as never, - blocks: {} as never, - rev: "rev1", - did, + did: input.did ?? ALICE_DID, collection: input.collection, rkey: input.rkey, - uri: { toString: () => atUri } as never, - event: input.event, + operation: input.event, record: input.record, - cid: {} as never, - } as CommitEvt; + }; } const HANDLER_TEST_WIKIS = [ @@ -127,28 +117,6 @@ describe("wiki handler", () => { expect(wiki?.visibility).toBe("private"); }); - test("rejects event when DID does not match AT-URI", () => { - handleCommitEvent( - makeCommitEvt({ - event: "create", - collection: "wiki.lichen.wiki", - rkey: "evil-wiki", - did: BOB_DID, - uri: { - toString: () => `at://${ALICE_DID}/wiki.lichen.wiki/evil-wiki`, - }, - record: { - name: "Evil Wiki", - visibility: "public", - createdAt: "2026-01-01T00:00:00.000Z", - }, - }), - ); - - const wiki = getWiki("evil-wiki"); - expect(wiki).toBeNull(); - }); - test("parses language field from wiki record", () => { handleCommitEvent( makeCommitEvt({ @@ -605,7 +573,6 @@ describe("bookmark handler", () => { }); test("deletes bookmark on delete event", () => { - const bookmarkUri = `at://${BOB_DID}/wiki.lichen.bookmark/bk1`; expect(isBookmarked(BOB_DID, WIKI_AT_URI)).toBe(true); handleCommitEvent( @@ -614,7 +581,6 @@ describe("bookmark handler", () => { collection: "wiki.lichen.bookmark", rkey: "bk1", did: BOB_DID, - uri: { toString: () => bookmarkUri }, }), ); diff --git a/tests/integration/helpers.ts b/tests/integration/helpers.ts index ffde55c..c4edf53 100644 --- a/tests/integration/helpers.ts +++ b/tests/integration/helpers.ts @@ -105,7 +105,13 @@ export function createTestFirehose(net: TestNetwork): Firehose { evt.event === "delete" ) { try { - handleCommitEvent(evt); + handleCommitEvent({ + did: evt.did, + collection: evt.collection, + rkey: evt.rkey, + operation: evt.event, + record: evt.event !== "delete" ? evt.record : undefined, + }); } catch (err) { console.error( `[test firehose] Error handling ${evt.collection} ${evt.event}:`,