To read the Bluesky firehose in Node.js, open a WebSocket to Jetstream and parse each message as JSON. On Node 18 or 20 use the ws package; on Node 22 and later the built-in WebSocket global works with no dependencies at all.
The basic connection is ten lines, and our firehose streaming guide already has a simple keyword monitor. This tutorial is about the part that guide leaves as an exercise: a consumer you can leave running, which resumes from a cursor, notices a stalled socket, backs off politely and writes what it finds to disk.
Quick Answer
| Need | Answer |
|---|---|
| Node version | 18+ with ws; 22+ can use the built-in WebSocket |
| Endpoint | wss://jetstream2.us-east.bsky.network/subscribe (or jetstream1/2 at us-east or us-west) |
| Posts only | ?wantedCollections=app.bsky.feed.post |
| Resume after a crash | &cursor=<time_us>, rewound a few seconds |
| Auth | None for reading; an app password only if you post back |
The Ten-Line Version
Create a folder, run npm init -y and npm install ws, and save this as stream.mjs (the .mjs extension lets you use import without changing package.json):
// stream.mjs (npm install ws)
import WebSocket from 'ws';
const url =
'wss://jetstream2.us-east.bsky.network/subscribe?wantedCollections=app.bsky.feed.post';
const ws = new WebSocket(url);
ws.on('message', (data) => {
const event = JSON.parse(data.toString());
if (event.kind !== 'commit' || event.commit.operation !== 'create') return;
console.log(event.did, event.commit.record.text?.slice(0, 80));
});
Run it with node stream.mjs and posts start scrolling. Each message is a JSON object with did, time_us, kind and, for record changes, a commit holding operation, collection, rkey and record. Deletes have no record, which is why the operation check comes first.
Without any dependency on Node 22+
Node 22 ships a browser-compatible WebSocket. The API is the browser's, so you use addEventListener and e.data is already a string for Jetstream's text frames:
// Node 22+: no dependency
const ws = new WebSocket(
'wss://jetstream2.us-east.bsky.network/subscribe?wantedCollections=app.bsky.feed.post'
);
ws.addEventListener('message', (e) => {
const event = JSON.parse(e.data);
if (event.kind !== 'commit' || event.commit.operation !== 'create') return;
console.log(event.did, event.commit.record.text?.slice(0, 80));
});
ws.addEventListener('close', () => console.log('closed'));
That is fine for scripts. For a long-running service, the rest of this post uses ws, mainly for terminate() and the handshake timeout. The same logic ports to the built-in client with close() and your own timer.
A Consumer You Can Leave Running
Here is the full program. It filters posts by keyword, writes matches to a JSON Lines file, saves its position, and reconnects with exponential backoff when anything goes wrong.
// stream.mjs (npm install ws)
import WebSocket from 'ws';
import { appendFile, readFile, writeFile } from 'node:fs/promises';
const BASE = 'wss://jetstream2.us-east.bsky.network/subscribe';
const CURSOR_FILE = './cursor.txt';
const OUT_FILE = './matches.jsonl';
const REWIND_US = 5_000_000; // five seconds, in microseconds
const STALL_MS = 30_000;
const PATTERN = /\b(skyscraper|jetstream|atproto)\b/i;
let lastTimeUs = null;
let attempt = 0;
async function loadCursor() {
try {
const saved = Number(await readFile(CURSOR_FILE, 'utf8'));
return Number.isFinite(saved) && saved > 0 ? saved : null;
} catch {
return null;
}
}
const postUrl = (did, rkey) => `https://bsky.app/profile/${did}/post/${rkey}`;
async function handle(event) {
lastTimeUs = event.time_us;
if (event.kind !== 'commit') return;
const { operation, record, rkey } = event.commit;
if (operation !== 'create' || typeof record?.text !== 'string') return;
if (!PATTERN.test(record.text)) return;
const row = {
url: postUrl(event.did, rkey),
did: event.did,
text: record.text,
createdAt: record.createdAt,
};
await appendFile(OUT_FILE, JSON.stringify(row) + '\n');
console.log(row.url);
}
async function connect() {
const saved = lastTimeUs ?? (await loadCursor());
const params = new URLSearchParams({ wantedCollections: 'app.bsky.feed.post' });
if (saved) params.set('cursor', String(saved - REWIND_US));
const ws = new WebSocket(`${BASE}?${params}`, { handshakeTimeout: 10_000 });
let watchdog;
const resetWatchdog = () => {
clearTimeout(watchdog);
watchdog = setTimeout(() => ws.terminate(), STALL_MS);
};
ws.on('open', () => {
attempt = 0;
resetWatchdog();
console.log('connected');
});
ws.on('message', (data) => {
resetWatchdog();
handle(JSON.parse(data.toString())).catch(console.error);
});
ws.on('error', (err) => console.error('socket error:', err.message));
ws.on('close', () => {
clearTimeout(watchdog);
const delay = Math.min(60_000, 1000 * 2 ** attempt) + Math.random() * 1000;
attempt += 1;
console.log(`closed; reconnecting in ${(delay / 1000).toFixed(1)}s`);
setTimeout(connect, delay);
});
}
// Persist the cursor every five seconds, and once more on Ctrl-C.
setInterval(() => {
if (lastTimeUs) writeFile(CURSOR_FILE, String(lastTimeUs)).catch(console.error);
}, 5_000);
process.on('SIGINT', async () => {
if (lastTimeUs) await writeFile(CURSOR_FILE, String(lastTimeUs));
process.exit(0);
});
connect();
What Each Piece Is For
The keyword filter
Jetstream filters by collection and by DID (wantedDids), never by text, so matching happens in your process. The regex uses \b word boundaries; a bare includes('art') would match "start" and "party". Each match gets a shareable link built from the DID and record key: https://bsky.app/profile/{did}/post/{rkey}. bsky.app accepts the DID in place of a handle, so there is no lookup.
The cursor
Jetstream's cursor is just the time_us of an event: a Unix timestamp in microseconds. The program keeps the latest one in memory, flushes it to cursor.txt every five seconds and on Ctrl-C, and on reconnect asks Jetstream to start five seconds before it. Rewinding means a handful of duplicates after every reconnect. That is the trade for never silently dropping an event, and it is why your downstream writes should be idempotent: key rows on did + rkey.
Because the cursor is a timestamp, not a server-specific sequence, you can usually point the same cursor at a different Jetstream host if your usual one is down. Replay only covers a bounded window, so a consumer that has been down for a long time should expect a gap.
The watchdog
The nastiest failure is not a disconnect. It is a socket that stays open and delivers nothing. The watchdog resets on every message; if 30 seconds pass in silence, it calls ws.terminate(), which fires close, which runs the normal reconnect path. On the posts collection, 30 seconds without a single event means something is wrong. If you subscribe to a quiet collection or a handful of DIDs, raise the threshold.
Backoff with jitter
Delays go 1s, 2s, 4s and so on, capped at a minute, plus up to a second of random jitter. The counter resets on a successful open. When a Jetstream instance restarts, every client reconnects at the same moment; jitter spreads you out of that herd. With ws, a failed connection emits error and then close, so the single close handler covers both cases.
The JSON Lines writer
One JSON object per line is the simplest durable format there is: append-only, greppable, and easy to load later with jq or pandas. If you need queries, swap appendFile for an insert into SQLite or Postgres. Keep it asynchronous and off the hot path; at peak, the posts collection alone is busy enough that a slow synchronous write will make you fall behind.
Deploying It
A consumer like this is a natural fit for a small always-on container. Two things to get right:
- Put the cursor on persistent storage. On most container hosts the filesystem is wiped on every deploy. Mount a volume for
cursor.txt, or store it in your database, or every deploy loses the events from the gap. - Handle
SIGTERMas well asSIGINT. Platforms sendSIGTERMon redeploy. Add the same handler for it so the final cursor is flushed.
Our Railway and Docker deployment guide covers the Dockerfile, health checks and volumes.
The Raw Firehose and the @atproto Packages
Jetstream is a convenience layer: it decodes the relay's CBOR stream once and re-emits JSON, at the cost of signature verification. If you need the signed stream itself (com.atproto.sync.subscribeRepos), use the official @atproto packages on npm rather than hand-rolling CBOR and CAR decoding. Bluesky's own TypeScript stack includes tooling for consuming and verifying the repo stream; check the package READMEs for the current API, since it has moved between packages over time.
For writing back – replying to a match, liking it, following the author – use @atproto/api. Our TypeScript bot guide covers logging in with an app password and posting. Remember that writes are rate-limited; queue them rather than firing one per event. The rate limits guide has the numbers.
Frequently Asked Questions
Do I need the ws package on Node 22?
No. Node 22 and later have a browser-style WebSocket global, which is enough for Jetstream. The ws package is still the better choice on Node 18 and 20, and it gives you terminate() and a handshake timeout, which are handy for reconnect logic.
How do I resume the Bluesky firehose after a restart in Node.js?
Store the time_us of the last event you handled and reconnect with a cursor query parameter set to that value minus a few seconds. Jetstream replays everything after the cursor, within its retention window. Expect a few duplicates and handle them idempotently.
Why does my Jetstream connection stop receiving events without closing?
A connection can stall without a close frame, especially through proxies and NAT. Keep a watchdog timer that resets on every message and force-closes the socket if nothing arrives for 30 seconds, so your normal reconnect path takes over.
Can I use this to build a Bluesky bot?
Yes. Use the stream to find what to react to, and @atproto/api to reply, like or follow. Writes are rate-limited, so queue them rather than firing one per matching event.
Which npm package reads the raw firehose?
The official @atproto packages include tooling for the raw subscribeRepos stream, and community libraries wrap Jetstream. For most projects, the plain WebSocket code in this tutorial is all you need.
Skip the Plumbing
- Keyword Alerts – an email when a keyword appears on Bluesky, without running a consumer
- Trending hashtags leaderboard – hashtag counts already aggregated from the stream
- Skyscraper – a free native Bluesky client for iPhone, iPad and Mac