diff --git a/src/pds.js b/src/pds.js index 826870d..e02e1ee 100644 --- a/src/pds.js +++ b/src/pds.js @@ -909,6 +909,10 @@ const pdsRoutes = { method: 'POST', handler: (pds, req, url) => pds.handleCreateRecord(req) }, + '/xrpc/com.atproto.repo.deleteRecord': { + method: 'POST', + handler: (pds, req, url) => pds.handleDeleteRecord(req) + }, '/xrpc/com.atproto.repo.getRecord': { handler: (pds, req, url) => pds.handleGetRecord(url) }, @@ -1154,6 +1158,112 @@ export class PersonalDataServer { return { uri, cid: recordCidStr, commit: commitCidStr } } + async deleteRecord(collection, rkey) { + const did = await this.getDid() + if (!did) throw new Error('PDS not initialized') + + const uri = `at://${did}/${collection}/${rkey}` + + // Check if record exists + const existing = this.sql.exec( + `SELECT cid FROM records WHERE uri = ?`, uri + ).toArray() + if (existing.length === 0) { + return { error: 'RecordNotFound', message: 'record not found' } + } + + // Delete from records table + this.sql.exec(`DELETE FROM records WHERE uri = ?`, uri) + + // Rebuild MST + const mst = new MST(this.sql) + const dataRoot = await mst.computeRoot() + + // Get previous commit + const prevCommits = this.sql.exec( + `SELECT cid, rev FROM commits ORDER BY seq DESC LIMIT 1` + ).toArray() + const prevCommit = prevCommits.length > 0 ? prevCommits[0] : null + + // Create commit + const rev = createTid() + const commit = { + did, + version: 3, + data: dataRoot ? new CID(cidToBytes(dataRoot)) : null, + rev, + prev: prevCommit?.cid ? new CID(cidToBytes(prevCommit.cid)) : null + } + + // Sign commit + const commitBytes = cborEncodeDagCbor(commit) + const signingKey = await this.getSigningKey() + const sig = await sign(signingKey, commitBytes) + + const signedCommit = { ...commit, sig } + const signedBytes = cborEncodeDagCbor(signedCommit) + const commitCid = await createCid(signedBytes) + const commitCidStr = cidToString(commitCid) + + // Store commit block + this.sql.exec( + `INSERT OR REPLACE INTO blocks (cid, data) VALUES (?, ?)`, + commitCidStr, signedBytes + ) + + // Store commit reference + this.sql.exec( + `INSERT INTO commits (cid, rev, prev) VALUES (?, ?, ?)`, + commitCidStr, rev, prevCommit?.cid || null + ) + + // Update head and rev + await this.state.storage.put('head', commitCidStr) + await this.state.storage.put('rev', rev) + + // Collect blocks for the event (commit + MST nodes, no record block) + const newBlocks = [] + newBlocks.push({ cid: commitCidStr, data: signedBytes }) + if (dataRoot) { + const mstBlocks = this.collectMstBlocks(dataRoot) + newBlocks.push(...mstBlocks) + } + + // Sequence event with delete action + const eventTime = new Date().toISOString() + const evt = cborEncode({ + ops: [{ action: 'delete', path: `${collection}/${rkey}`, cid: null }], + blocks: buildCarFile(commitCidStr, newBlocks), + rev, + time: eventTime + }) + this.sql.exec( + `INSERT INTO seq_events (did, commit_cid, evt) VALUES (?, ?, ?)`, + did, commitCidStr, evt + ) + + // Broadcast to subscribers + const evtRows = this.sql.exec( + `SELECT * FROM seq_events ORDER BY seq DESC LIMIT 1` + ).toArray() + if (evtRows.length > 0) { + this.broadcastEvent(evtRows[0]) + // Forward to default DO for relay subscribers + if (this.env?.PDS) { + const defaultId = this.env.PDS.idFromName('default') + const defaultPds = this.env.PDS.get(defaultId) + const row = evtRows[0] + const evtArray = Array.from(new Uint8Array(row.evt)) + defaultPds.fetch(new Request('http://internal/forward-event', { + method: 'POST', + body: JSON.stringify({ ...row, evt: evtArray }) + })).catch(e => console.log('forward error:', e)) + } + } + + return { ok: true } + } + formatEvent(evt) { // AT Protocol frame format: header + body // Use DAG-CBOR encoding for body (CIDs need tag 42 + 0x00 prefix) @@ -1337,6 +1447,22 @@ export class PersonalDataServer { } } + async handleDeleteRecord(request) { + const body = await request.json() + if (!body.collection || !body.rkey) { + return Response.json({ error: 'InvalidRequest', message: 'missing collection or rkey' }, { status: 400 }) + } + try { + const result = await this.deleteRecord(body.collection, body.rkey) + if (result.error) { + return Response.json(result, { status: 404 }) + } + return Response.json({}) + } catch (err) { + return Response.json({ error: err.message }, { status: 500 }) + } + } + async handleGetRecord(url) { const collection = url.searchParams.get('collection') const rkey = url.searchParams.get('rkey') @@ -1724,6 +1850,24 @@ async function handleRequest(request, env) { return pds.fetch(request) } + // POST repo endpoints have repo in body + if (url.pathname === '/xrpc/com.atproto.repo.deleteRecord' || + url.pathname === '/xrpc/com.atproto.repo.createRecord') { + const body = await request.json() + const repo = body.repo + if (!repo) { + return Response.json({ error: 'InvalidRequest', message: 'missing repo param' }, { status: 400 }) + } + const id = env.PDS.idFromName(repo) + const pds = env.PDS.get(id) + // Re-create request with body since we consumed it + return pds.fetch(new Request(request.url, { + method: 'POST', + headers: request.headers, + body: JSON.stringify(body) + })) + } + const did = url.searchParams.get('did') if (!did) { return new Response('missing did param', { status: 400 })