// The migration handlers, driven through the factory: the status report, // blob accounting and CAR import for an account arriving here, and the // activate/deactivate switch-over. getServiceAuth rides along because a // migration's service tokens come from it. import { describe, expect, it } from 'vitest'; import { generateKeyPair, importPrivateKey, } from '../packages/core/src/crypto.js'; import { createMigrationHandlers } from '../packages/core/src/handlers/migration.js'; import { createXrpcServerHandlers } from '../packages/core/src/handlers/xrpc-server.js'; import { buildMst } from '../packages/core/src/mst.js'; import { buildCarFile, CID, cborEncodeDagCbor, cidToBytes, cidToString, createCid, createTid, } from '../packages/core/src/repo.js'; const DID = 'did:plc:migrationtestaccount2345'; const FULL_ACCESS = { did: DID, scope: 'com.atproto.access' }; const NARROW = { did: DID, scope: 'repo:app.bsky.feed.post' }; /** * @param {Object} [options] * @param {string|null} [options.did] * @param {string} [options.status] * @param {string[]} [options.storedBlobs] - CIDs the blob store holds * @param {Array<{cid: string, recordUri: string}>} [options.referenced] */ function createMigration({ did = DID, status = 'active', storedBlobs = [], referenced = [], } = {}) { /** @type {{status: string, handle: string}} */ const state = { status, handle: 'alice.example.com' }; /** @type {Array<{type: string, body: Object}>} */ const events = []; /** @type {string[]} */ const crawled = []; const calls = { relinked: 0, syncEvents: 0 }; const migration = createMigrationHandlers( /** @type {any} */ ({ actorStorage: { getAccountStatus: async () => state.status, setAccountStatus: async (/** @type {string} */ s) => { state.status = s; }, getHandle: async () => state.handle, getLatestCommit: async () => ({ cid: 'bafycommit', rev: '3lrev' }), listAllRecords: async () => [{ key: 'app.bsky.feed.post/3aaa' }], getBlob: async (/** @type {string} */ cid) => storedBlobs.includes(cid) ? new Uint8Array([1]) : null, putBlock: async () => {}, putRecord: async () => {}, putCommit: async () => {}, setDid: async () => {}, }, readOnly: false, getDid: async () => did, getSigningKey: async () => ({}), readOnlyError: () => Response.json({ error: 'ReadOnly' }, { status: 401 }), emitEvent: async ( /** @type {string} */ type, /** @type {Object} */ body, ) => { events.push({ type, body }); }, emitSyncEvent: async () => { calls.syncEvents++; }, notifyRelay: async (/** @type {string} */ d) => { crawled.push(d); }, collectRepoBlocks: async () => [{ cid: 'bafyblock' }], collectReferencedBlobs: async () => referenced, relinkBlobRecords: async () => { calls.relinked++; return referenced.length; }, }), ); return { migration, state, events, crawled, calls }; } /** * @param {ReturnType} migration * @param {string} path * @param {{method?: string, body?: Uint8Array, auth?: Object|null, search?: string}} [options] */ function call( migration, path, { method, body, auth = null, search = '' } = {}, ) { const url = new URL(`https://pds.example.com${path}${search}`); const request = new Request(url, { method: method || 'GET', body: /** @type {any} */ (body), }); const handler = /** @type {any} */ (migration.routes[path]).handler; return handler(request, url, auth); } /** A minimal but valid repo CAR: records, MST, signed commit. */ /** * @param {string} did * @param {Array<{collection: string, rkey: string, record: object}>} recordsToInclude */ async function buildTestCar(did, recordsToInclude) { /** @type {Array<{cid: string, data: Uint8Array}>} */ const blocks = []; /** @type {Array<{key: string, cid: string}>} */ const mstEntries = []; for (const { collection, rkey, record } of recordsToInclude) { const recordBytes = cborEncodeDagCbor(record); const recordCid = cidToString(await createCid(recordBytes)); blocks.push({ cid: recordCid, data: recordBytes }); mstEntries.push({ key: `${collection}/${rkey}`, cid: recordCid }); } mstEntries.sort((a, b) => a.key.localeCompare(b.key)); const mstRoot = await buildMst(mstEntries, async (cid, data) => { blocks.push({ cid, data }); }); const commit = { did, version: 3, rev: createTid(), prev: null, data: mstRoot ? new CID(cidToBytes(mstRoot)) : null, sig: new Uint8Array(64), }; const commitBytes = cborEncodeDagCbor(commit); const commitCid = cidToString(await createCid(commitBytes)); blocks.push({ cid: commitCid, data: commitBytes }); return buildCarFile(commitCid, blocks); } describe('checkAccountStatus', () => { const PATH = '/xrpc/com.atproto.server.checkAccountStatus'; it('requires authentication and an initialized server', async () => { const { migration } = createMigration(); expect((await call(migration, PATH)).status).toBe(401); const { migration: uninit } = createMigration({ did: null }); expect((await call(uninit, PATH, { auth: FULL_ACCESS })).status).toBe(400); }); it('reports the repo shape and how many expected blobs arrived', async () => { const { migration } = createMigration({ referenced: [ { cid: 'bafyblob1', recordUri: 'at://x/1' }, { cid: 'bafyblob2', recordUri: 'at://x/2' }, ], storedBlobs: ['bafyblob1'], }); const res = await call(migration, PATH, { auth: FULL_ACCESS }); expect(await res.json()).toEqual({ activated: true, validDid: true, repoCommit: 'bafycommit', repoRev: '3lrev', repoBlocks: 1, indexedRecords: 1, privateStateValues: 0, expectedBlobs: 2, importedBlobs: 1, }); }); }); describe('listMissingBlobs', () => { const PATH = '/xrpc/com.atproto.repo.listMissingBlobs'; const referenced = [ { cid: 'bafyb', recordUri: 'at://x/2' }, { cid: 'bafya', recordUri: 'at://x/1' }, { cid: 'bafya', recordUri: 'at://x/1-again' }, { cid: 'bafyc', recordUri: 'at://x/3' }, ]; it('lists each missing blob once, sorted, honoring limit and cursor', async () => { const { migration } = createMigration({ referenced, storedBlobs: ['bafyc'], }); const first = await call(migration, PATH, { auth: FULL_ACCESS, search: '?limit=1', }); const firstBody = await first.json(); expect(firstBody.blobs).toEqual([{ cid: 'bafya', recordUri: 'at://x/1' }]); expect(firstBody.cursor).toBe('bafya'); const rest = await call(migration, PATH, { auth: FULL_ACCESS, search: '?cursor=bafya', }); const restBody = await rest.json(); expect(restBody.blobs).toEqual([{ cid: 'bafyb', recordUri: 'at://x/2' }]); expect(restBody.cursor).toBeUndefined(); }); it('requires authentication', async () => { const { migration } = createMigration(); expect((await call(migration, PATH)).status).toBe(401); }); }); describe('importRepo', () => { const PATH = '/xrpc/com.atproto.repo.importRepo'; it('gates on auth, scope, initialization, and a non-empty body', async () => { const { migration } = createMigration(); expect((await call(migration, PATH, { method: 'POST' })).status).toBe(401); expect( (await call(migration, PATH, { method: 'POST', auth: NARROW })).status, ).toBe(403); const { migration: uninit } = createMigration({ did: null }); expect( ( await call(uninit, PATH, { method: 'POST', auth: FULL_ACCESS, body: new Uint8Array([1]), }) ).status, ).toBe(400); const empty = await call(migration, PATH, { method: 'POST', auth: FULL_ACCESS, body: new Uint8Array(0), }); expect((await empty.json()).message).toBe('Empty CAR body'); }); it('rejects a CAR that does not parse', async () => { const { migration } = createMigration(); const res = await call(migration, PATH, { method: 'POST', auth: FULL_ACCESS, body: new Uint8Array([0xff, 0xff, 0xff]), }); expect(res.status).toBe(400); expect((await res.json()).message).toMatch(/^Invalid CAR:/); }); it("rejects another account's repository", async () => { const { migration } = createMigration(); const car = await buildTestCar('did:plc:someoneelseentirely2345a', [ { collection: 'app.bsky.feed.post', rkey: '3aaa', record: { text: 'x' } }, ]); const res = await call(migration, PATH, { method: 'POST', auth: FULL_ACCESS, body: car, }); expect(res.status).toBe(400); expect((await res.json()).message).toContain('belongs to'); }); it('imports this account, relinks blobs, and announces the state', async () => { const { migration, calls } = createMigration(); const car = await buildTestCar(DID, [ { collection: 'app.bsky.feed.post', rkey: '3aaa', record: { text: 'x' } }, ]); const res = await call(migration, PATH, { method: 'POST', auth: FULL_ACCESS, body: car, }); expect(res.status).toBe(200); expect(calls.relinked).toBe(1); expect(calls.syncEvents).toBe(1); }); }); describe('activateAccount and deactivateAccount', () => { const ACTIVATE = '/xrpc/com.atproto.server.activateAccount'; const DEACTIVATE = '/xrpc/com.atproto.server.deactivateAccount'; it('requires a full-access session on both', async () => { const { migration } = createMigration(); for (const path of [ACTIVATE, DEACTIVATE]) { expect((await call(migration, path, { method: 'POST' })).status).toBe( 401, ); expect( (await call(migration, path, { method: 'POST', auth: NARROW })).status, ).toBe(403); } }); it('deactivation stores the state and announces #account inactive', async () => { const { migration, state, events } = createMigration(); const res = await call(migration, DEACTIVATE, { method: 'POST', auth: FULL_ACCESS, }); expect(res.status).toBe(200); expect(state.status).toBe('deactivated'); expect(events).toEqual([ { type: 'account', body: { active: false, status: 'deactivated' } }, ]); expect(await migration.isDeactivated()).toBe(true); expect((await migration.checkAccountActive())?.status).toBe(409); }); it('activation re-announces identity and asks the relay to crawl', async () => { const { migration, state, events, crawled } = createMigration({ status: 'deactivated', }); const res = await call(migration, ACTIVATE, { method: 'POST', auth: FULL_ACCESS, }); expect(res.status).toBe(200); expect(state.status).toBe('active'); expect(events).toEqual([ { type: 'account', body: { active: true } }, { type: 'identity', body: { handle: 'alice.example.com' } }, ]); expect(crawled).toEqual([DID]); expect(await migration.checkAccountActive()).toBeNull(); }); }); describe('getServiceAuth', () => { const PATH = '/xrpc/com.atproto.server.getServiceAuth'; /** * @param {Object} [options] * @param {string|null} [options.did] * @param {boolean} [options.withKey] * @param {Response|null} [options.scopeError] */ async function createServer({ did = DID, withKey = true, scopeError = null, } = {}) { const { privateKey } = await generateKeyPair(); const signingKey = withKey ? await importPrivateKey(privateKey) : null; return createXrpcServerHandlers( /** @type {any} */ ({ getDid: async () => did, getSigningKey: async () => signingKey, checkRpcScope: () => scopeError, }), ); } /** * @param {ReturnType} server * @param {string} search * @param {Object|null} auth */ function callServiceAuth(server, search, auth) { const url = new URL(`https://pds.example.com${PATH}${search}`); const handler = /** @type {any} */ (server.routes[PATH]).handler; return handler(new Request(url), url, auth); } it('mints a service JWT for the bare DID audience', async () => { const server = await createServer(); const res = await callServiceAuth( server, '?aud=did:web:api.bsky.app%23bsky_appview&lxm=app.bsky.feed.getFeedSkeleton', FULL_ACCESS, ); expect(res.status).toBe(200); const { token } = await res.json(); const payload = JSON.parse(atob(token.split('.')[1])); expect(payload.iss).toBe(DID); expect(payload.aud).toBe('did:web:api.bsky.app'); expect(payload.lxm).toBe('app.bsky.feed.getFeedSkeleton'); }); it('gates on auth, parameters, scope, and initialization', async () => { const server = await createServer(); expect((await callServiceAuth(server, '?aud=did:web:x', null)).status).toBe( 401, ); expect( (await callServiceAuth(server, '?aud=did:web:x', FULL_ACCESS)).status, ).toBe(400); const denied = await createServer({ scopeError: Response.json({ error: 'InvalidScope' }, { status: 403 }), }); expect( (await callServiceAuth(denied, '?aud=did:web:x&lxm=m', FULL_ACCESS)) .status, ).toBe(403); const uninit = await createServer({ withKey: false }); expect( (await callServiceAuth(uninit, '?aud=did:web:x&lxm=m', FULL_ACCESS)) .status, ).toBe(400); }); });