From 05c327b4589cfeefb3dd0d61976f4255da0ce040 Mon Sep 17 00:00:00 2001 From: Chad Miller Date: Mon, 5 Jan 2026 13:23:53 -0800 Subject: [PATCH] refactor: extract PersonalDataServer route table MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-Authored-By: Claude Opus 4.5 --- src/pds.js | 424 +++++++++++++++++++++++++++++------------------------ 1 file changed, 235 insertions(+), 189 deletions(-) diff --git a/src/pds.js b/src/pds.js index d5f07d8..9917da8 100644 --- a/src/pds.js +++ b/src/pds.js @@ -814,6 +814,75 @@ export function buildCarFile(rootCid, blocks) { return car } +/** + * Route handler function type + * @callback RouteHandler + * @param {PersonalDataServer} pds - PDS instance + * @param {Request} request - HTTP request + * @param {URL} url - Parsed URL + * @returns {Promise} HTTP response + */ + +/** + * @typedef {Object} Route + * @property {string} [method] - Required HTTP method (default: any) + * @property {RouteHandler} handler - Handler function + */ + +/** @type {Record} */ +const pdsRoutes = { + '/.well-known/atproto-did': { + handler: (pds, req, url) => pds.handleAtprotoDid() + }, + '/init': { + method: 'POST', + handler: (pds, req, url) => pds.handleInit(req) + }, + '/status': { + handler: (pds, req, url) => pds.handleStatus() + }, + '/reset-repo': { + handler: (pds, req, url) => pds.handleResetRepo() + }, + '/forward-event': { + handler: (pds, req, url) => pds.handleForwardEvent(req) + }, + '/register-did': { + handler: (pds, req, url) => pds.handleRegisterDid(req) + }, + '/get-registered-dids': { + handler: (pds, req, url) => pds.handleGetRegisteredDids() + }, + '/repo-info': { + handler: (pds, req, url) => pds.handleRepoInfo() + }, + '/xrpc/com.atproto.server.describeServer': { + handler: (pds, req, url) => pds.handleDescribeServer(req) + }, + '/xrpc/com.atproto.sync.listRepos': { + handler: (pds, req, url) => pds.handleListRepos() + }, + '/xrpc/com.atproto.repo.createRecord': { + method: 'POST', + handler: (pds, req, url) => pds.handleCreateRecord(req) + }, + '/xrpc/com.atproto.repo.getRecord': { + handler: (pds, req, url) => pds.handleGetRecord(url) + }, + '/xrpc/com.atproto.sync.getLatestCommit': { + handler: (pds, req, url) => pds.handleGetLatestCommit() + }, + '/xrpc/com.atproto.sync.getRepoStatus': { + handler: (pds, req, url) => pds.handleGetRepoStatus() + }, + '/xrpc/com.atproto.sync.getRepo': { + handler: (pds, req, url) => pds.handleGetRepo() + }, + '/xrpc/com.atproto.sync.subscribeRepos': { + handler: (pds, req, url) => pds.handleSubscribeRepos(req, url) + } +} + export class PersonalDataServer { constructor(state, env) { this.state = state @@ -1107,221 +1176,198 @@ export class PersonalDataServer { } } - async fetch(request) { - const url = new URL(request.url) - - // Handle resolution - doesn't require ?did= param - if (url.pathname === '/.well-known/atproto-did') { - let did = await this.getDid() - // If no DID on this instance, check registered DIDs (default instance) - if (!did) { - const registeredDids = await this.state.storage.get('registeredDids') || [] - did = registeredDids[0] - } - if (!did) { - return new Response('User not found', { status: 404 }) - } - return new Response(did, { - headers: { 'Content-Type': 'text/plain' } - }) - } - - if (url.pathname === '/init') { - const body = await request.json() - if (!body.did || !body.privateKey) { - return Response.json({ error: 'missing did or privateKey' }, { status: 400 }) - } - await this.initIdentity(body.did, body.privateKey, body.handle || null) - return Response.json({ ok: true, did: body.did, handle: body.handle || null }) - } - if (url.pathname === '/status') { - const did = await this.getDid() - return Response.json({ - initialized: !!did, - did: did || null - }) - } - // Reset endpoint - clears all repo data but keeps identity - if (url.pathname === '/reset-repo') { - this.sql.exec(`DELETE FROM blocks`) - this.sql.exec(`DELETE FROM records`) - this.sql.exec(`DELETE FROM commits`) - this.sql.exec(`DELETE FROM seq_events`) - await this.state.storage.delete('head') - await this.state.storage.delete('rev') - return Response.json({ ok: true, message: 'repo data cleared' }) - } - // Internal endpoint to forward events for relay broadcasting - if (url.pathname === '/forward-event') { - const evt = await request.json() - // Convert evt back to proper format and broadcast - const numSockets = [...this.state.getWebSockets()].length - console.log(`forward-event: received event seq=${evt.seq}, ${numSockets} connected sockets`) - this.broadcastEvent({ - seq: evt.seq, - did: evt.did, - commit_cid: evt.commit_cid, - evt: new Uint8Array(Object.values(evt.evt)) - }) - return Response.json({ ok: true, sockets: numSockets }) - } - // Internal endpoint to register DIDs for discovery - if (url.pathname === '/register-did') { - const body = await request.json() - const registeredDids = await this.state.storage.get('registeredDids') || [] - if (!registeredDids.includes(body.did)) { - registeredDids.push(body.did) - await this.state.storage.put('registeredDids', registeredDids) - } - return Response.json({ ok: true }) - } - // Internal endpoint to get registered DIDs - if (url.pathname === '/get-registered-dids') { + async handleAtprotoDid() { + let did = await this.getDid() + if (!did) { const registeredDids = await this.state.storage.get('registeredDids') || [] - return Response.json({ dids: registeredDids }) + did = registeredDids[0] } - // Internal endpoint to get repo info (head/rev) - if (url.pathname === '/repo-info') { - const head = await this.state.storage.get('head') - const rev = await this.state.storage.get('rev') - return Response.json({ head: head || null, rev: rev || null }) - } - if (url.pathname === '/xrpc/com.atproto.server.describeServer') { - // Server DID should be did:web based on hostname, passed via header - const hostname = request.headers.get('x-hostname') || 'localhost' - return Response.json({ - did: `did:web:${hostname}`, - availableUserDomains: [`.${hostname}`], - inviteCodeRequired: false, - phoneVerificationRequired: false, - links: {}, - contact: {} - }) + if (!did) { + return new Response('User not found', { status: 404 }) } - if (url.pathname === '/xrpc/com.atproto.sync.listRepos') { - const registeredDids = await this.state.storage.get('registeredDids') || [] - // If this is the default instance, return registered DIDs - // If this is a user instance, return its own DID - const did = await this.getDid() - const repos = did ? [{ did, head: null, rev: null }] : - registeredDids.map(d => ({ did: d, head: null, rev: null })) - return Response.json({ repos }) + return new Response(did, { headers: { 'Content-Type': 'text/plain' } }) + } + + async handleInit(request) { + const body = await request.json() + if (!body.did || !body.privateKey) { + return Response.json({ error: 'missing did or privateKey' }, { status: 400 }) } - if (url.pathname === '/xrpc/com.atproto.repo.createRecord') { - if (request.method !== 'POST') { - return Response.json({ error: 'method not allowed' }, { status: 405 }) - } + await this.initIdentity(body.did, body.privateKey, body.handle || null) + return Response.json({ ok: true, did: body.did, handle: body.handle || null }) + } - const body = await request.json() - if (!body.collection || !body.record) { - return Response.json({ error: 'missing collection or record' }, { status: 400 }) - } + async handleStatus() { + const did = await this.getDid() + return Response.json({ initialized: !!did, did: did || null }) + } - try { - const result = await this.createRecord(body.collection, body.record, body.rkey) - return Response.json(result) - } catch (err) { - return Response.json({ error: err.message }, { status: 500 }) - } + async handleResetRepo() { + this.sql.exec(`DELETE FROM blocks`) + this.sql.exec(`DELETE FROM records`) + this.sql.exec(`DELETE FROM commits`) + this.sql.exec(`DELETE FROM seq_events`) + await this.state.storage.delete('head') + await this.state.storage.delete('rev') + return Response.json({ ok: true, message: 'repo data cleared' }) + } + + async handleForwardEvent(request) { + const evt = await request.json() + const numSockets = [...this.state.getWebSockets()].length + console.log(`forward-event: received event seq=${evt.seq}, ${numSockets} connected sockets`) + this.broadcastEvent({ + seq: evt.seq, + did: evt.did, + commit_cid: evt.commit_cid, + evt: new Uint8Array(Object.values(evt.evt)) + }) + return Response.json({ ok: true, sockets: numSockets }) + } + + async handleRegisterDid(request) { + const body = await request.json() + const registeredDids = await this.state.storage.get('registeredDids') || [] + if (!registeredDids.includes(body.did)) { + registeredDids.push(body.did) + await this.state.storage.put('registeredDids', registeredDids) } - if (url.pathname === '/xrpc/com.atproto.repo.getRecord') { - const collection = url.searchParams.get('collection') - const rkey = url.searchParams.get('rkey') + return Response.json({ ok: true }) + } - if (!collection || !rkey) { - return Response.json({ error: 'missing collection or rkey' }, { status: 400 }) - } + async handleGetRegisteredDids() { + const registeredDids = await this.state.storage.get('registeredDids') || [] + return Response.json({ dids: registeredDids }) + } - const did = await this.getDid() - const uri = `at://${did}/${collection}/${rkey}` + async handleRepoInfo() { + const head = await this.state.storage.get('head') + const rev = await this.state.storage.get('rev') + return Response.json({ head: head || null, rev: rev || null }) + } - const rows = this.sql.exec( - `SELECT cid, value FROM records WHERE uri = ?`, uri - ).toArray() + handleDescribeServer(request) { + const hostname = request.headers.get('x-hostname') || 'localhost' + return Response.json({ + did: `did:web:${hostname}`, + availableUserDomains: [`.${hostname}`], + inviteCodeRequired: false, + phoneVerificationRequired: false, + links: {}, + contact: {} + }) + } - if (rows.length === 0) { - return Response.json({ error: 'record not found' }, { status: 404 }) - } + async handleListRepos() { + const registeredDids = await this.state.storage.get('registeredDids') || [] + const did = await this.getDid() + const repos = did ? [{ did, head: null, rev: null }] : + registeredDids.map(d => ({ did: d, head: null, rev: null })) + return Response.json({ repos }) + } - const row = rows[0] - // Decode CBOR for response (convert ArrayBuffer to Uint8Array) - const value = cborDecode(new Uint8Array(row.value)) + async handleCreateRecord(request) { + const body = await request.json() + if (!body.collection || !body.record) { + return Response.json({ error: 'missing collection or record' }, { status: 400 }) + } + try { + const result = await this.createRecord(body.collection, body.record, body.rkey) + return Response.json(result) + } catch (err) { + return Response.json({ error: err.message }, { status: 500 }) + } + } - return Response.json({ uri, cid: row.cid, value }) + async handleGetRecord(url) { + const collection = url.searchParams.get('collection') + const rkey = url.searchParams.get('rkey') + if (!collection || !rkey) { + return Response.json({ error: 'missing collection or rkey' }, { status: 400 }) } - if (url.pathname === '/xrpc/com.atproto.sync.getLatestCommit') { - const commits = this.sql.exec( - `SELECT cid, rev FROM commits ORDER BY seq DESC LIMIT 1` - ).toArray() + const did = await this.getDid() + const uri = `at://${did}/${collection}/${rkey}` + const rows = this.sql.exec( + `SELECT cid, value FROM records WHERE uri = ?`, uri + ).toArray() + if (rows.length === 0) { + return Response.json({ error: 'record not found' }, { status: 404 }) + } + const row = rows[0] + const value = cborDecode(new Uint8Array(row.value)) + return Response.json({ uri, cid: row.cid, value }) + } - if (commits.length === 0) { - return Response.json({ error: 'RepoNotFound', message: 'repo not found' }, { status: 404 }) - } + handleGetLatestCommit() { + const commits = this.sql.exec( + `SELECT cid, rev FROM commits ORDER BY seq DESC LIMIT 1` + ).toArray() + if (commits.length === 0) { + return Response.json({ error: 'RepoNotFound', message: 'repo not found' }, { status: 404 }) + } + return Response.json({ cid: commits[0].cid, rev: commits[0].rev }) + } - return Response.json({ cid: commits[0].cid, rev: commits[0].rev }) + async handleGetRepoStatus() { + const did = await this.getDid() + const commits = this.sql.exec( + `SELECT cid, rev FROM commits ORDER BY seq DESC LIMIT 1` + ).toArray() + if (commits.length === 0 || !did) { + return Response.json({ error: 'RepoNotFound', message: 'repo not found' }, { status: 404 }) } - if (url.pathname === '/xrpc/com.atproto.sync.getRepoStatus') { - const did = await this.getDid() - const commits = this.sql.exec( - `SELECT cid, rev FROM commits ORDER BY seq DESC LIMIT 1` - ).toArray() + return Response.json({ did, active: true, status: 'active', rev: commits[0].rev }) + } - if (commits.length === 0 || !did) { - return Response.json({ error: 'RepoNotFound', message: 'repo not found' }, { status: 404 }) - } + handleGetRepo() { + const commits = this.sql.exec( + `SELECT cid FROM commits ORDER BY seq DESC LIMIT 1` + ).toArray() + if (commits.length === 0) { + return Response.json({ error: 'repo not found' }, { status: 404 }) + } + const blocks = this.sql.exec(`SELECT cid, data FROM blocks`).toArray() + const blocksForCar = blocks.map(b => ({ + cid: b.cid, + data: new Uint8Array(b.data) + })) + const car = buildCarFile(commits[0].cid, blocksForCar) + return new Response(car, { + headers: { 'content-type': 'application/vnd.ipld.car' } + }) + } - return Response.json({ - did, - active: true, - status: 'active', - rev: commits[0].rev - }) + handleSubscribeRepos(request, url) { + const upgradeHeader = request.headers.get('Upgrade') + if (upgradeHeader !== 'websocket') { + return new Response('expected websocket', { status: 426 }) } - if (url.pathname === '/xrpc/com.atproto.sync.getRepo') { - const commits = this.sql.exec( - `SELECT cid FROM commits ORDER BY seq DESC LIMIT 1` + const { 0: client, 1: server } = new WebSocketPair() + this.state.acceptWebSocket(server) + const cursor = url.searchParams.get('cursor') + if (cursor) { + const events = this.sql.exec( + `SELECT * FROM seq_events WHERE seq > ? ORDER BY seq`, + parseInt(cursor) ).toArray() - - if (commits.length === 0) { - return Response.json({ error: 'repo not found' }, { status: 404 }) + for (const evt of events) { + server.send(this.formatEvent(evt)) } - - const blocks = this.sql.exec(`SELECT cid, data FROM blocks`).toArray() - const blocksForCar = blocks.map(b => ({ - cid: b.cid, - data: new Uint8Array(b.data) - })) - const car = buildCarFile(commits[0].cid, blocksForCar) - - return new Response(car, { - headers: { 'content-type': 'application/vnd.ipld.car' } - }) } - if (url.pathname === '/xrpc/com.atproto.sync.subscribeRepos') { - const upgradeHeader = request.headers.get('Upgrade') - if (upgradeHeader !== 'websocket') { - return new Response('expected websocket', { status: 426 }) - } - - const { 0: client, 1: server } = new WebSocketPair() - this.state.acceptWebSocket(server) - - // Send backlog if cursor provided - const cursor = url.searchParams.get('cursor') - if (cursor) { - const events = this.sql.exec( - `SELECT * FROM seq_events WHERE seq > ? ORDER BY seq`, - parseInt(cursor) - ).toArray() + return new Response(null, { status: 101, webSocket: client }) + } - for (const evt of events) { - server.send(this.formatEvent(evt)) - } - } + async fetch(request) { + const url = new URL(request.url) + const route = pdsRoutes[url.pathname] - return new Response(null, { status: 101, webSocket: client }) + if (!route) { + return Response.json({ error: 'not found' }, { status: 404 }) + } + if (route.method && request.method !== route.method) { + return Response.json({ error: 'method not allowed' }, { status: 405 }) } - return Response.json({ error: 'not found' }, { status: 404 }) + return route.handler(this, request, url) } } -- 2.51.2