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,