From 46aaa59c49b7ab6cd3ab59737cb0210e2d2f52b3 Mon Sep 17 00:00:00 2001 From: Chad Miller Date: Sat, 8 Aug 2026 22:52:07 -0700 Subject: [PATCH] refactor(core): move preferences and the firehose into handler modules Two small route surfaces, each reaching three members of the instance. formatFirehoseEvent stays reachable as this._firehose.formatFirehoseEvent because the sync-event path and the relay notification both call it. The Firehose banner in pds.js had been titled Blobs until the blob handlers moved out, which is why subscribeRepos and the frame format sat under it. Coverage: the firehose 4.8% -> 100% of statements, preferences 13.3% -> 100%. The firehose tests read each frame back apart, so the header and body are asserted per event type rather than as an opaque buffer. A correction to how I read the coverage report earlier: a file missing from the text table is not necessarily complete. xrpc-preferences.js was absent and I took that for 100%, but it was at 13.3%; the glob threshold that grouped it with xrpc-blob.js is what showed otherwise. Read the json when it matters. Co-Authored-By: Claude Opus 5 (1M context) --- packages/core/src/handlers/xrpc-firehose.js | 155 ++++++++++ .../core/src/handlers/xrpc-preferences.js | 81 ++++++ packages/core/src/pds.js | 209 ++------------ test/xrpc-firehose.test.js | 264 ++++++++++++++++++ test/xrpc-preferences.test.js | 116 ++++++++ vitest.config.js | 7 +- 6 files changed, 641 insertions(+), 191 deletions(-) create mode 100644 packages/core/src/handlers/xrpc-firehose.js create mode 100644 packages/core/src/handlers/xrpc-preferences.js create mode 100644 test/xrpc-firehose.test.js create mode 100644 test/xrpc-preferences.test.js diff --git a/packages/core/src/handlers/xrpc-firehose.js b/packages/core/src/handlers/xrpc-firehose.js new file mode 100644 index 0000000..05cd9a4 --- /dev/null +++ b/packages/core/src/handlers/xrpc-firehose.js @@ -0,0 +1,155 @@ +// @pdsjs/core/handlers/xrpc-firehose - com.atproto.sync.subscribeRepos. +// +// A subscriber gets a websocket and every commit after the cursor it asked for, +// each as an AT Protocol frame: a header, then the event body, both dag-cbor. + +import { + CID, + cborDecode, + cborEncode, + cborEncodeDagCbor, + cidToBytes, +} from '../repo.js'; +import { parseNumericCursor } from '../validation.js'; + +/** + * @typedef {Object} FirehoseContext + * @property {import('../ports.js').ActorStoragePort} actorStorage + * @property {import('../ports.js').WebSocketPort} [webSocket] - Absent, subscribeRepos is not served + */ + +/** + * @param {FirehoseContext} ctx + * @returns {{ + * formatFirehoseEvent: (evt: {seq: number, commit_cid: string, evt: Uint8Array}) => Uint8Array, + * routes: import('../pds.js').Routes, + * }} + */ +export function createFirehoseHandlers(ctx) { + const { actorStorage, webSocket } = ctx; + + /** + * Format event for firehose (AT Protocol frame format) + * @param {{seq: number, commit_cid: string, evt: Uint8Array}} evt + * @returns {Uint8Array} + */ + function formatFirehoseEvent(evt) { + // Decode stored event to get ops, blocks, rev, and time + const evtData = cborDecode(evt.evt); + + // Typed (non-commit) events: #identity, #account, #sync + if (evtData.type && evtData.type !== 'commit') { + const typedHeader = cborEncode({ op: 1, t: `#${evtData.type}` }); + /** @type {Record} */ + const typedBody = { seq: evt.seq, did: evtData.did, time: evtData.time }; + if (evtData.type === 'identity') { + if (evtData.handle != null) typedBody.handle = evtData.handle; + } else if (evtData.type === 'account') { + typedBody.active = evtData.active; + if (evtData.status != null) typedBody.status = evtData.status; + } else if (evtData.type === 'sync') { + typedBody.blocks = evtData.blocks || new Uint8Array(0); + typedBody.rev = evtData.rev; + } + const encodedBody = cborEncodeDagCbor(typedBody); + const typedFrame = new Uint8Array( + typedHeader.length + encodedBody.length, + ); + typedFrame.set(typedHeader, 0); + typedFrame.set(encodedBody, typedHeader.length); + return typedFrame; + } + + // AT Protocol frame format: header + body + const header = cborEncode({ op: 1, t: '#commit' }); + + // Transform ops to wrap CIDs in CID class for DAG-CBOR tag 42 encoding + // Stored events have cid as { $link: cidString } or null + const ops = (evtData.ops || []).map( + ( + /** @type {{action: string, path: string, cid?: {$link: string}|null}} */ op, + ) => ({ + ...op, + cid: op.cid?.$link ? new CID(cidToBytes(op.cid.$link)) : null, + }), + ); + + // Build body with blocks included + // CIDs must be wrapped in CID class for proper DAG-CBOR tag 42 encoding + // Use commit_cid from database row (not from decoded event) for reliability + const body = cborEncodeDagCbor({ + seq: evt.seq, + rebase: false, + tooBig: false, + repo: evtData.did, + commit: new CID(cidToBytes(evt.commit_cid)), + rev: evtData.rev, + since: evtData.since ?? null, + ...(evtData.prevData + ? { prevData: new CID(cidToBytes(evtData.prevData)) } + : {}), + blocks: evtData.blocks || new Uint8Array(0), + ops, + blobs: [], + time: evtData.time, + }); + + // Combine header and body + const frame = new Uint8Array(header.length + body.length); + frame.set(header, 0); + frame.set(body, header.length); + return frame; + } + + /** + * com.atproto.sync.subscribeRepos - WebSocket firehose + * @param {Request} request + * @param {URL} url + * @returns {Response} + */ + function handleSubscribeRepos(request, url) { + // Check if WebSocket port is available + if (!webSocket) { + return Response.json( + { error: 'NotImplemented', message: 'WebSocket not supported' }, + { status: 501 }, + ); + } + + // Check for WebSocket upgrade + if (!webSocket.isUpgrade(request)) { + return new Response('Expected WebSocket upgrade', { status: 426 }); + } + + // Validate cursor before upgrading + const cursorParam = url.searchParams.get('cursor'); + const parsedCursor = parseNumericCursor(cursorParam); + if (parsedCursor && !parsedCursor.valid) { + return Response.json( + { error: 'InvalidRequest', message: parsedCursor.error }, + { status: 400 }, + ); + } + + // Upgrade connection and send historical events + return webSocket.upgrade(request, async (sender) => { + if (parsedCursor?.valid) { + const { events } = await actorStorage.getEvents( + parsedCursor.value, + 1000, + ); + for (const evt of events) { + sender.send(formatFirehoseEvent(evt)); + } + } + }); + } + return { + formatFirehoseEvent, + routes: { + '/xrpc/com.atproto.sync.subscribeRepos': { + handler: handleSubscribeRepos, + }, + }, + }; +} diff --git a/packages/core/src/handlers/xrpc-preferences.js b/packages/core/src/handlers/xrpc-preferences.js new file mode 100644 index 0000000..ed80de4 --- /dev/null +++ b/packages/core/src/handlers/xrpc-preferences.js @@ -0,0 +1,81 @@ +// @pdsjs/core/handlers/xrpc-preferences - app.bsky.actor.*Preferences. +// +// The server keeps the preferences array as the client sends it. Its contents +// are the AppView's concern, not this one's. + +/** + * @typedef {Object} PreferencesContext + * @property {import('../ports.js').ActorStoragePort} actorStorage + * @property {boolean} readOnly + * @property {() => Response} readOnlyError + */ + +/** + * @param {PreferencesContext} ctx + * @returns {{ routes: import('../pds.js').Routes }} + */ +export function createPreferencesHandlers(ctx) { + const { actorStorage, readOnly, readOnlyError } = ctx; + + /** + * app.bsky.actor.getPreferences + * @param {Request} _request + * @param {URL} _url + * @param {import('../pds.js').Auth|null} auth + * @returns {Promise} + */ + async function handleGetPreferences(_request, _url, auth) { + if (!auth) { + return Response.json( + { error: 'AuthenticationRequired', message: 'Authentication required' }, + { status: 401 }, + ); + } + + const preferences = await actorStorage.getPreferences(); + return Response.json({ preferences }); + } + + /** + * app.bsky.actor.putPreferences + * @param {Request} request + * @param {URL} _url + * @param {import('../pds.js').Auth|null} auth + * @returns {Promise} + */ + async function handlePutPreferences(request, _url, auth) { + if (readOnly) return readOnlyError(); + + if (!auth) { + return Response.json( + { error: 'AuthenticationRequired', message: 'Authentication required' }, + { status: 401 }, + ); + } + + const { preferences } = await request.json(); + if (!Array.isArray(preferences)) { + return Response.json( + { error: 'InvalidRequest', message: 'preferences must be an array' }, + { status: 400 }, + ); + } + + await actorStorage.setPreferences(preferences); + return Response.json({}); + } + + return { + routes: { + '/xrpc/app.bsky.actor.getPreferences': { + auth: 'required', + handler: handleGetPreferences, + }, + '/xrpc/app.bsky.actor.putPreferences': { + method: 'POST', + auth: 'required', + handler: handlePutPreferences, + }, + }, + }; +} diff --git a/packages/core/src/pds.js b/packages/core/src/pds.js index c7d91ef..c510168 100644 --- a/packages/core/src/pds.js +++ b/packages/core/src/pds.js @@ -63,7 +63,9 @@ import { } from './email.js'; import { createBackupHandlers } from './handlers/account-backup.js'; import { createBlobHandlers } from './handlers/xrpc-blob.js'; +import { createFirehoseHandlers } from './handlers/xrpc-firehose.js'; import { createIdentityHandlers } from './handlers/xrpc-identity.js'; +import { createPreferencesHandlers } from './handlers/xrpc-preferences.js'; import { createRepoHandlers } from './handlers/xrpc-repo.js'; import { createSyncHandlers } from './handlers/xrpc-sync.js'; import { loadRepositoryFromCar } from './loader.js'; @@ -750,6 +752,17 @@ export class PersonalDataServer { // The server builds this one time, so the engine has its own run lock. The // engine needs the repo CAR builder from the sync handlers. Thus it comes // after them. + this._firehose = createFirehoseHandlers({ + actorStorage: this.actorStorage, + webSocket: this.webSocket, + }); + + this._preferences = createPreferencesHandlers({ + actorStorage: this.actorStorage, + readOnly: this.readOnly, + readOnlyError: () => this.readOnlyError(), + }); + this._repo = createRepoHandlers({ actorStorage: this.actorStorage, lexiconResolver: this.lexiconResolver, @@ -800,6 +813,8 @@ export class PersonalDataServer { ...this._blob.routes, ...this._identity.routes, ...this._repo.routes, + ...this._preferences.routes, + ...this._firehose.routes, ...this._backup.routes, ...(spaceRoutes ?? {}), }; @@ -2370,58 +2385,6 @@ export class PersonalDataServer { } } - // ════════════════════════════════════════════════════════════════════════════ - // XRPC Handlers - Preferences - // ════════════════════════════════════════════════════════════════════════════ - - /** - * app.bsky.actor.getPreferences - * @param {Request} _request - * @param {URL} _url - * @param {Auth|null} auth - * @returns {Promise} - */ - async handleGetPreferences(_request, _url, auth) { - if (!auth) { - return Response.json( - { error: 'AuthenticationRequired', message: 'Authentication required' }, - { status: 401 }, - ); - } - - const preferences = await this.actorStorage.getPreferences(); - return Response.json({ preferences }); - } - - /** - * app.bsky.actor.putPreferences - * @param {Request} request - * @param {URL} _url - * @param {Auth|null} auth - * @returns {Promise} - */ - async handlePutPreferences(request, _url, auth) { - if (this.readOnly) return this.readOnlyError(); - - if (!auth) { - return Response.json( - { error: 'AuthenticationRequired', message: 'Authentication required' }, - { status: 401 }, - ); - } - - const { preferences } = await request.json(); - if (!Array.isArray(preferences)) { - return Response.json( - { error: 'InvalidRequest', message: 'preferences must be an array' }, - { status: 400 }, - ); - } - - await this.actorStorage.setPreferences(preferences); - return Response.json({}); - } - // ════════════════════════════════════════════════════════════════════════════ // OAuth Handlers // ════════════════════════════════════════════════════════════════════════════ @@ -6551,130 +6514,6 @@ export class PersonalDataServer { }); } - // ════════════════════════════════════════════════════════════════════════════ - // Firehose - // ════════════════════════════════════════════════════════════════════════════ - // - // com.atproto.sync.subscribeRepos and its frame format. Blob upload, download - // and listing are in handlers/xrpc-blob.js. - - /** - * Format event for firehose (AT Protocol frame format) - * @param {{seq: number, commit_cid: string, evt: Uint8Array}} evt - * @returns {Uint8Array} - */ - formatFirehoseEvent(evt) { - // Decode stored event to get ops, blocks, rev, and time - const evtData = cborDecode(evt.evt); - - // Typed (non-commit) events: #identity, #account, #sync - if (evtData.type && evtData.type !== 'commit') { - const typedHeader = cborEncode({ op: 1, t: `#${evtData.type}` }); - /** @type {Record} */ - const typedBody = { seq: evt.seq, did: evtData.did, time: evtData.time }; - if (evtData.type === 'identity') { - if (evtData.handle != null) typedBody.handle = evtData.handle; - } else if (evtData.type === 'account') { - typedBody.active = evtData.active; - if (evtData.status != null) typedBody.status = evtData.status; - } else if (evtData.type === 'sync') { - typedBody.blocks = evtData.blocks || new Uint8Array(0); - typedBody.rev = evtData.rev; - } - const encodedBody = cborEncodeDagCbor(typedBody); - const typedFrame = new Uint8Array( - typedHeader.length + encodedBody.length, - ); - typedFrame.set(typedHeader, 0); - typedFrame.set(encodedBody, typedHeader.length); - return typedFrame; - } - - // AT Protocol frame format: header + body - const header = cborEncode({ op: 1, t: '#commit' }); - - // Transform ops to wrap CIDs in CID class for DAG-CBOR tag 42 encoding - // Stored events have cid as { $link: cidString } or null - const ops = (evtData.ops || []).map( - ( - /** @type {{action: string, path: string, cid?: {$link: string}|null}} */ op, - ) => ({ - ...op, - cid: op.cid?.$link ? new CID(cidToBytes(op.cid.$link)) : null, - }), - ); - - // Build body with blocks included - // CIDs must be wrapped in CID class for proper DAG-CBOR tag 42 encoding - // Use commit_cid from database row (not from decoded event) for reliability - const body = cborEncodeDagCbor({ - seq: evt.seq, - rebase: false, - tooBig: false, - repo: evtData.did, - commit: new CID(cidToBytes(evt.commit_cid)), - rev: evtData.rev, - since: evtData.since ?? null, - ...(evtData.prevData - ? { prevData: new CID(cidToBytes(evtData.prevData)) } - : {}), - blocks: evtData.blocks || new Uint8Array(0), - ops, - blobs: [], - time: evtData.time, - }); - - // Combine header and body - const frame = new Uint8Array(header.length + body.length); - frame.set(header, 0); - frame.set(body, header.length); - return frame; - } - - /** - * com.atproto.sync.subscribeRepos - WebSocket firehose - * @param {Request} request - * @param {URL} url - * @returns {Response} - */ - handleSubscribeRepos(request, url) { - // Check if WebSocket port is available - if (!this.webSocket) { - return Response.json( - { error: 'NotImplemented', message: 'WebSocket not supported' }, - { status: 501 }, - ); - } - - // Check for WebSocket upgrade - if (!this.webSocket.isUpgrade(request)) { - return new Response('Expected WebSocket upgrade', { status: 426 }); - } - - // Validate cursor before upgrading - const cursorParam = url.searchParams.get('cursor'); - const parsedCursor = parseNumericCursor(cursorParam); - if (parsedCursor && !parsedCursor.valid) { - return Response.json( - { error: 'InvalidRequest', message: parsedCursor.error }, - { status: 400 }, - ); - } - - // Upgrade connection and send historical events - return this.webSocket.upgrade(request, async (sender) => { - if (parsedCursor?.valid) { - const { events } = await this.actorStorage.getEvents( - parsedCursor.value, - 1000, - ); - for (const evt of events) { - sender.send(this.formatFirehoseEvent(evt)); - } - } - }); - } - // ════════════════════════════════════════════════════════════════════════════ // Init and Well-Known Endpoints // ════════════════════════════════════════════════════════════════════════════ @@ -7049,7 +6888,7 @@ export class PersonalDataServer { // Broadcast to connected WebSocket subscribers if (this.webSocket?.broadcast) { - const frame = this.formatFirehoseEvent({ + const frame = this._firehose.formatFirehoseEvent({ seq, commit_cid: commitCid, evt: eventEncoded, @@ -7842,7 +7681,7 @@ export class PersonalDataServer { await this.actorStorage.putEvent(seq, did, '', eventEncoded); if (this.webSocket?.broadcast) { - const frame = this.formatFirehoseEvent({ + const frame = this._firehose.formatFirehoseEvent({ seq, commit_cid: '', evt: eventEncoded, @@ -8075,15 +7914,7 @@ const routes = { // resolveHandle and updateHandle come from handlers/xrpc-identity.js // Preferences - '/xrpc/app.bsky.actor.getPreferences': { - auth: 'required', - handler: PersonalDataServer.prototype.handleGetPreferences, - }, - '/xrpc/app.bsky.actor.putPreferences': { - method: 'POST', - auth: 'required', - handler: PersonalDataServer.prototype.handlePutPreferences, - }, + // The preferences endpoints come from handlers/xrpc-preferences.js // OAuth '/.well-known/oauth-authorization-server': { @@ -8327,7 +8158,5 @@ const routes = { // listRepos, getRepoStatus, getRepo, getRecord and getLatestCommit come from // handlers/xrpc-sync.js // getBlob and listBlobs come from handlers/xrpc-blob.js - '/xrpc/com.atproto.sync.subscribeRepos': { - handler: PersonalDataServer.prototype.handleSubscribeRepos, - }, + // subscribeRepos comes from handlers/xrpc-firehose.js }; diff --git a/test/xrpc-firehose.test.js b/test/xrpc-firehose.test.js new file mode 100644 index 0000000..cfb8550 --- /dev/null +++ b/test/xrpc-firehose.test.js @@ -0,0 +1,264 @@ +// Firehose tests - frame format per event type, and the subscribe handshake +import { describe, expect, it } from 'vitest'; +import { createFirehoseHandlers } from '../packages/core/src/handlers/xrpc-firehose.js'; +import { cborDecode, cborEncodeDagCbor } from '../packages/core/src/repo.js'; + +const DID = 'did:plc:firehosetestaccount2'; +const COMMIT_CID = + 'bafyreiald5e2ib3gzzgreu7vp2dsrjpphhehnlh46lrk53mw7z2na5eoiq'; + +/** + * @param {{events?: Array, withSocket?: boolean, isUpgrade?: boolean}} [options] + */ +function createFirehose({ + events = [], + withSocket = true, + isUpgrade = true, +} = {}) { + /** @type {Uint8Array[]} */ + const sent = []; + /** @type {Array<{cursor: number, limit: number}>} */ + const queries = []; + + const webSocket = { + isUpgrade: () => isUpgrade, + upgrade: ( + /** @type {Request} */ _request, + /** @type {(sender: {send: (b: Uint8Array) => void}) => Promise} */ onOpen, + ) => { + // A real port answers 101, which Response refuses to construct here. + onOpen({ send: (b) => sent.push(b) }); + return new Response(null, { headers: { 'x-upgraded': 'yes' } }); + }, + }; + + const firehose = createFirehoseHandlers({ + actorStorage: /** @type {any} */ ({ + async getEvents( + /** @type {number} */ cursor, + /** @type {number} */ limit, + ) { + queries.push({ cursor, limit }); + return { events }; + }, + }), + webSocket: withSocket ? /** @type {any} */ (webSocket) : undefined, + }); + return { firehose, sent, queries }; +} + +/** The two dag-cbor values a frame is made of, read back out. */ +function readFrame(/** @type {Uint8Array} */ frame) { + // cborDecode reads the first value; the body follows it in the same buffer. + const header = cborDecode(frame); + // Re-encode the header to learn its length, then decode the remainder. + let offset = 1; + while (offset < frame.length) { + try { + const body = cborDecode(frame.slice(offset)); + if (body && typeof body === 'object' && 'seq' in body) { + return { header, body }; + } + } catch { + // keep scanning + } + offset += 1; + } + throw new Error('no body found in frame'); +} + +/** @param {Object} evtData */ +function storedEvent(evtData, seq = 1) { + return { seq, commit_cid: COMMIT_CID, evt: cborEncodeDagCbor(evtData) }; +} + +const SUBSCRIBE = '/xrpc/com.atproto.sync.subscribeRepos'; + +/** + * @param {ReturnType} firehose + * @param {string} query + */ +function call(firehose, query = '') { + const route = firehose.routes[SUBSCRIBE]; + const url = new URL(`https://pds.example.com${SUBSCRIBE}${query}`); + const request = new Request(url); + const handler = + /** @type {(request: Request, url: URL, auth: null) => Response} */ ( + route.handler + ); + return handler(request, url, null); +} + +describe('formatFirehoseEvent', () => { + it('frames a commit with its ops and cid', () => { + const { firehose } = createFirehose(); + const frame = firehose.formatFirehoseEvent( + storedEvent({ + did: DID, + rev: 'rev1', + time: '2026-01-01T00:00:00.000Z', + blocks: new Uint8Array([1, 2, 3]), + ops: [ + { + action: 'create', + path: 'app.bsky.feed.post/abc', + cid: { $link: COMMIT_CID }, + }, + ], + }), + ); + + const { header, body } = readFrame(frame); + expect(header).toEqual({ op: 1, t: '#commit' }); + expect(body.seq).toBe(1); + expect(body.repo).toBe(DID); + expect(body.rev).toBe('rev1'); + expect(body.ops).toHaveLength(1); + expect(body.ops[0].action).toBe('create'); + expect(body.rebase).toBe(false); + expect(body.tooBig).toBe(false); + expect(body.blobs).toEqual([]); + }); + + it('carries a null cid for a delete', () => { + const { firehose } = createFirehose(); + const frame = firehose.formatFirehoseEvent( + storedEvent({ + did: DID, + rev: 'rev2', + time: '2026-01-01T00:00:00.000Z', + ops: [{ action: 'delete', path: 'app.bsky.feed.post/abc', cid: null }], + }), + ); + + const { body } = readFrame(frame); + expect(body.ops[0].cid).toBe(null); + }); + + it('frames an identity event with its handle', () => { + const { firehose } = createFirehose(); + const frame = firehose.formatFirehoseEvent( + storedEvent({ + type: 'identity', + did: DID, + time: '2026-01-01T00:00:00.000Z', + handle: 'alice.example.com', + }), + ); + + const { header, body } = readFrame(frame); + expect(header).toEqual({ op: 1, t: '#identity' }); + expect(body.did).toBe(DID); + expect(body.handle).toBe('alice.example.com'); + }); + + it('omits a handle an identity event does not carry', () => { + const { firehose } = createFirehose(); + const frame = firehose.formatFirehoseEvent( + storedEvent({ + type: 'identity', + did: DID, + time: '2026-01-01T00:00:00.000Z', + }), + ); + + const { body } = readFrame(frame); + expect('handle' in body).toBe(false); + }); + + it('frames an account event with its active flag and status', () => { + const { firehose } = createFirehose(); + const frame = firehose.formatFirehoseEvent( + storedEvent({ + type: 'account', + did: DID, + time: '2026-01-01T00:00:00.000Z', + active: false, + status: 'deactivated', + }), + ); + + const { header, body } = readFrame(frame); + expect(header).toEqual({ op: 1, t: '#account' }); + expect(body.active).toBe(false); + expect(body.status).toBe('deactivated'); + }); + + it('frames a sync event with its blocks and rev', () => { + const { firehose } = createFirehose(); + const frame = firehose.formatFirehoseEvent( + storedEvent({ + type: 'sync', + did: DID, + time: '2026-01-01T00:00:00.000Z', + blocks: new Uint8Array([9, 8]), + rev: 'rev3', + }), + ); + + const { header, body } = readFrame(frame); + expect(header).toEqual({ op: 1, t: '#sync' }); + expect(body.rev).toBe('rev3'); + expect(body.blocks).toEqual(new Uint8Array([9, 8])); + }); +}); + +describe('subscribeRepos', () => { + it('is a 501 with no websocket port', () => { + const { firehose } = createFirehose({ withSocket: false }); + const response = call(firehose); + + expect(response.status).toBe(501); + }); + + it('is a 426 without an upgrade header', () => { + const { firehose } = createFirehose({ isUpgrade: false }); + const response = call(firehose); + + expect(response.status).toBe(426); + }); + + it('refuses a cursor that is not a number', () => { + const { firehose } = createFirehose(); + const response = call(firehose, '?cursor=abc'); + + expect(response.status).toBe(400); + }); + + it('upgrades and replays the events after the cursor', async () => { + const { firehose, sent, queries } = createFirehose({ + events: [ + storedEvent( + { + did: DID, + rev: 'rev1', + time: '2026-01-01T00:00:00.000Z', + ops: [], + }, + 5, + ), + ], + }); + + const response = call(firehose, '?cursor=4'); + expect(response.headers.get('x-upgraded')).toBe('yes'); + // The upgrade callback runs on its own, so wait a turn for the replay. + await new Promise((r) => setTimeout(r, 0)); + + expect(queries).toEqual([{ cursor: 4, limit: 1000 }]); + expect(sent).toHaveLength(1); + expect(readFrame(sent[0]).body.seq).toBe(5); + }); + + it('sends no history without a cursor', async () => { + const { firehose, sent, queries } = createFirehose({ + events: [storedEvent({ did: DID, rev: 'r', time: 't', ops: [] })], + }); + + call(firehose); + await new Promise((r) => setTimeout(r, 0)); + + expect(queries).toHaveLength(0); + expect(sent).toHaveLength(0); + }); +}); diff --git a/test/xrpc-preferences.test.js b/test/xrpc-preferences.test.js new file mode 100644 index 0000000..c667906 --- /dev/null +++ b/test/xrpc-preferences.test.js @@ -0,0 +1,116 @@ +// Preferences handler tests - auth, read-only, and the array shape +import { describe, expect, it } from 'vitest'; +import { createPreferencesHandlers } from '../packages/core/src/handlers/xrpc-preferences.js'; + +const DID = 'did:plc:preferencestestaccount'; +const FULL = { did: DID, scope: 'com.atproto.access' }; + +/** @param {{readOnly?: boolean, stored?: unknown[]}} [options] */ +function createPreferences({ readOnly = false, stored = [] } = {}) { + let preferences = stored; + const handlers = createPreferencesHandlers({ + actorStorage: /** @type {any} */ ({ + async getPreferences() { + return preferences; + }, + async setPreferences(/** @type {unknown[]} */ next) { + preferences = next; + }, + }), + readOnly, + readOnlyError: () => + Response.json( + { error: 'AuthenticationRequired', message: 'This PDS is read-only' }, + { status: 401 }, + ), + }); + return { handlers, read: () => preferences }; +} + +const GET = '/xrpc/app.bsky.actor.getPreferences'; +const PUT = '/xrpc/app.bsky.actor.putPreferences'; + +/** + * @param {ReturnType} handlers + * @param {string} path + * @param {{body?: Object, auth?: {did: string, scope: string}|null}} [options] + */ +async function call(handlers, path, { body, auth = null } = {}) { + const route = handlers.routes[path]; + const url = new URL(`https://pds.example.com${path}`); + /** @type {RequestInit} */ + const init = { method: body ? 'POST' : 'GET' }; + if (body) init.body = JSON.stringify(body); + const handler = + /** @type {(request: Request, url: URL, auth: unknown) => Promise} */ ( + route.handler + ); + return handler(new Request(url, init), url, auth); +} + +describe('getPreferences', () => { + it('answers with what is stored', async () => { + const { handlers } = createPreferences({ + stored: [{ $type: 'app.bsky.actor.defs#adultContentPref' }], + }); + const body = await (await call(handlers, GET, { auth: FULL })).json(); + + expect(body.preferences).toHaveLength(1); + }); + + it('needs a session', async () => { + const { handlers } = createPreferences(); + const response = await call(handlers, GET); + + expect(response.status).toBe(401); + expect((await response.json()).error).toBe('AuthenticationRequired'); + }); +}); + +describe('putPreferences', () => { + it('replaces the stored array', async () => { + const { handlers, read } = createPreferences({ stored: ['old'] }); + const response = await call(handlers, PUT, { + auth: FULL, + body: { preferences: ['new'] }, + }); + + expect(response.status).toBe(200); + expect(await response.json()).toEqual({}); + expect(read()).toEqual(['new']); + }); + + it('refuses a read-only server', async () => { + const { handlers, read } = createPreferences({ stored: ['old'] }); + const readOnly = createPreferences({ readOnly: true, stored: ['old'] }); + const response = await call(readOnly.handlers, PUT, { + auth: FULL, + body: { preferences: ['new'] }, + }); + + expect(response.status).toBe(401); + expect(readOnly.read()).toEqual(['old']); + expect(read()).toEqual(['old']); + }); + + it('needs a session', async () => { + const { handlers } = createPreferences(); + const response = await call(handlers, PUT, { + body: { preferences: [] }, + }); + + expect(response.status).toBe(401); + }); + + it('refuses anything that is not an array', async () => { + const { handlers, read } = createPreferences({ stored: ['old'] }); + const response = await call(handlers, PUT, { + auth: FULL, + body: { preferences: { nope: true } }, + }); + + expect(response.status).toBe(400); + expect((await response.json()).message).toMatch(/must be an array/); + expect(read()).toEqual(['old']); + }); +}); diff --git a/vitest.config.js b/vitest.config.js index 2845cae..c1b279b 100644 --- a/vitest.config.js +++ b/vitest.config.js @@ -72,11 +72,16 @@ export default defineConfig({ branches: 84, functions: 100, }, - 'packages/core/src/handlers/xrpc-blob.js': { + 'packages/core/src/handlers/{xrpc-blob,xrpc-preferences}.js': { statements: 100, branches: 95, functions: 100, }, + 'packages/core/src/handlers/xrpc-firehose.js': { + statements: 100, + branches: 86, + functions: 100, + }, 'packages/spaces/src/*.js': { statements: 90, branches: 80, -- 2.51.2