import assert from 'node:assert/strict' import { describe, it } from 'node:test' import { COLLECTIONS, DEVICE_DIRECTORY_COLLECTIONS, MemoryEnvelopeStore, MemoryPrivateBlobStore, MemoryPrivateBusNetwork, MemoryRecordStore, MemoryWantList, Quarantine, WIRE_VERSION, WirePrivateBus, authorizeEndpoint, connectablePeers, connectionDirectory, createTidSource, envelopeKey, exportPublicKey, generateDeviceKey, indexDigest, loopbackLink, materialize, privateBlobRef, readDeviceDirectory, recordCid, seal, spaceTopic, wantFor, } from '../../core/dist/index.js' import { MemorySyncStateStore, PrivateSpaceIngestor, RepoPoller, directoryPoller, privateWriter, } from '../dist/index.js' const ROOT = 'did:plc:privateroot' const MEMBER = 'did:plc:privatemember' async function device(did, keyId) { const pair = await generateDeviceKey(true) return { did, keyId, publicKey: await exportPublicKey(pair.publicKey), privateKey: pair.privateKey, } } /** The public device directory record, exactly as a poller would have written it. */ const directoryRecord = (signer, createdAt) => ({ did: signer.did, collection: COLLECTIONS.device, rkey: signer.keyId, uri: `at://${signer.did}/${COLLECTIONS.device}/${signer.keyId}`, cid: `cid-device-${signer.keyId}`, rev: '0000000000001', value: { $type: COLLECTIONS.device, deviceKeyId: signer.keyId, publicKey: signer.publicKey, algorithm: 'ed25519', kind: 'node', createdAt, }, }) /** The public address hint that makes a device dialable — discovery metadata, never trust. */ const addressRecord = (signer, endpointId) => ({ did: signer.did, collection: COLLECTIONS.deviceAddress, rkey: `addr-${signer.keyId}`, uri: `at://${signer.did}/${COLLECTIONS.deviceAddress}/addr-${signer.keyId}`, cid: `cid-addr-${signer.keyId}`, rev: '0000000000002', value: { $type: COLLECTIONS.deviceAddress, deviceKeyId: signer.keyId, endpointId, createdAt: '2026-02-01T00:00:02Z', }, }) /** A whole private space, sealed: space, membership, project, goal, and two messages. */ async function privateCorpus() { const rootDevice = await device(ROOT, 'root-node-1') const memberDevice = await device(MEMBER, 'member-browser-1') const tid = createTidSource(() => 1_700_000_000_000, 1) const spaceUri = `at://${ROOT}/${COLLECTIONS.space}/space` const values = [] const spaceValue = { $type: COLLECTIONS.space, name: 'Private', description: 'No PDS holds a copy.', private: true, createdAt: '2026-03-01T00:00:00Z', } const spaceRef = { uri: spaceUri, cid: await recordCid(spaceValue) } values.push([rootDevice, 'space', spaceValue]) const grantValue = { $type: COLLECTIONS.addMember, space: spaceRef, did: MEMBER, kind: 'human', role: 'member', createdAt: '2026-03-01T00:00:01Z', } values.push([rootDevice, 'grant', grantValue]) const projectValue = { $type: COLLECTIONS.project, space: spaceRef, name: 'radial', gitUrl: 'https://tangled.org/disnetdev.com/radial', defaultBranch: 'main', checks: [], autoReview: {}, createdAt: '2026-03-01T00:00:02Z', } const projectRef = { uri: `at://${ROOT}/${COLLECTIONS.project}/project`, cid: await recordCid(projectValue), } values.push([rootDevice, 'project', projectValue]) const goalValue = { $type: COLLECTIONS.goal, space: spaceRef, project: projectRef, title: 'Ship private mode', body: 'Everything here stays off every PDS.', createdAt: '2026-03-01T00:00:03Z', } const goalRef = { uri: `at://${MEMBER}/${COLLECTIONS.goal}/goal`, cid: await recordCid(goalValue), } values.push([memberDevice, 'goal', goalValue]) for (const [index, body] of ['First note.', 'Second note.'].entries()) { values.push([ memberDevice, `note-${index}`, { $type: COLLECTIONS.message, goal: goalRef, body, mentions: [], createdAt: `2026-03-01T00:01:0${index}Z`, }, ]) } const envelopes = [] for (const [signer, rkey, record] of values) { envelopes.push( await seal({ space: spaceUri, did: signer.did, collection: record.$type, rkey, record, rev: tid(), deviceKeyId: signer.keyId, privateKey: signer.privateKey, }), ) } return { spaceUri, rootDevice, memberDevice, envelopes } } /** One replica: its own stores, its own ingestor, its own endpoint on the shared bus. */ function replica(network, id, spaceUri, directory, options = {}) { const records = new MemoryRecordStore() for (const record of directory) records.put(record) const wants = new MemoryWantList() const envelopes = new MemoryEnvelopeStore(records) // A peer serves catch-up out of its own envelope store, not out of what it happened to publish. const bus = network.endpoint(id, { corpus: () => envelopes.all() }) const ingestor = new PrivateSpaceIngestor(spaceUri, bus, envelopes, records, { wants, quarantine: new Quarantine(wants, options.quarantineLimits), bootstrapDids: [ROOT], ...(options.spacePin ? { spacePin: options.spacePin } : {}), }) return { records, bus, wants, envelopes, ingestor } } const shuffle = (items, seed) => { const result = [...items] let state = seed >>> 0 || 1 const random = () => { state ^= state << 13 state ^= state >>> 17 state ^= state << 5 return state >>> 0 } for (let index = result.length - 1; index > 0; index -= 1) { const target = random() % (index + 1) ;[result[index], result[target]] = [result[target], result[index]] } return result } describe('private space ingestion', () => { it('withholds a partial corpus inventory until genesis arrives', async () => { const { spaceUri, rootDevice, memberDevice, envelopes } = await privateCorpus() const records = new MemoryRecordStore() for (const record of [ directoryRecord(rootDevice, '2026-02-01T00:00:00Z'), directoryRecord(memberDevice, '2026-02-01T00:00:01Z'), ]) { records.put(record) } const store = new MemoryEnvelopeStore(records) const wants = new MemoryWantList() const blobs = new MemoryPrivateBlobStore() const requests = [] const blobAsks = [] const bus = { async join() {}, async publish() {}, subscribe() { return () => {} }, async catchUp(_space, request) { requests.push(request) return [] }, async fetchBlob(_space, need) { blobAsks.push(need.cid) return undefined }, async close() {}, } const ingestor = new PrivateSpaceIngestor(spaceUri, bus, store, records, { wants, blobs }) const fragment = envelopes.find((envelope) => envelope.collection !== COLLECTIONS.space) assert.ok(fragment) assert.equal((await ingestor.offer(fragment)).admitted, 1) wants.add(wantFor(fragment)) // A record naming bytes, arrived before genesis: a blob-request would put its CID on the wire. const ref = await privateBlobRef(new TextEncoder().encode('private bytes'), 'image/png') const image = await seal({ space: spaceUri, did: memberDevice.did, collection: COLLECTIONS.image, rkey: 'image-1', record: { $type: COLLECTIONS.image, blob: ref, createdAt: '2026-03-01T00:03:00Z' }, rev: '3ms26zx46yg2z', deviceKeyId: memberDevice.keyId, privateKey: memberDevice.privateKey, }) assert.equal((await ingestor.offer(image)).admitted, 1) await assert.rejects(ingestor.sync(), /Space record not found/) assert.deepEqual(requests, [{ summaries: [], wants: [] }]) // The store still knows what is missing; the wire was never told (ADR §30.3) — a blob CID is // an identifier of private content, and the partial footing dials an unfiltered directory. assert.deepEqual(blobAsks, []) assert.equal(ingestor.missingBlobs().length, 1) assert.equal(ingestor.lastReport.blobs.missing, 1) }) it('folds a corpus delivered in any order into one index and one digest', async () => { const { spaceUri, rootDevice, memberDevice, envelopes } = await privateCorpus() const directory = [ directoryRecord(rootDevice, '2026-02-01T00:00:00Z'), directoryRecord(memberDevice, '2026-02-01T00:00:01Z'), ] let expected for (let seed = 1; seed <= 12; seed += 1) { const network = new MemoryPrivateBusNetwork() const one = replica(network, `a${seed}`, spaceUri, directory) await one.ingestor.start() // Delivered in a different order every time, and every envelope twice. const ordering = shuffle(envelopes, seed) for (const envelope of [...ordering, ...ordering]) await one.ingestor.offer(envelope) const index = await one.ingestor.sync() const digest = indexDigest(index) expected ??= digest assert.deepEqual(digest, expected, `seed ${seed}`) assert.equal(index.private, true) assert.equal(index.goals.length, 1) assert.equal(index.goals[0].messages.length, 2) // Duplicates are idempotent: one version per record, whatever arrived. assert.equal(one.envelopes.all().length, envelopes.length) } }) it('converges two replicas that were partitioned while writes happened', async () => { const { spaceUri, rootDevice, memberDevice, envelopes } = await privateCorpus() const directory = [ directoryRecord(rootDevice, '2026-02-01T00:00:00Z'), directoryRecord(memberDevice, '2026-02-01T00:00:01Z'), ] const network = new MemoryPrivateBusNetwork() const alice = replica(network, 'alice', spaceUri, directory) const bob = replica(network, 'bob', spaceUri, directory) await alice.ingestor.start() await bob.ingestor.start() network.partition('alice', 'bob') for (const envelope of envelopes) { await alice.ingestor.offer(envelope) await alice.bus.publish(spaceUri, envelope) } const aliceIndex = await alice.ingestor.sync() // Bob has the space record from nowhere yet, so his fold has nothing to build. await assert.rejects(bob.ingestor.sync(), /Space record not found/) network.heal('alice', 'bob') const bobIndex = await bob.ingestor.sync() assert.deepEqual(indexDigest(bobIndex), indexDigest(aliceIndex)) }) it('quarantines what it cannot verify yet and folds it once the key is published', async () => { const { spaceUri, rootDevice, memberDevice, envelopes } = await privateCorpus() const network = new MemoryPrivateBusNetwork() // The member's device record has not been polled yet: their envelopes are sound, unverifiable, // and must not be thrown away. const late = replica(network, 'late', spaceUri, [ directoryRecord(rootDevice, '2026-02-01T00:00:00Z'), ]) await late.ingestor.start() for (const envelope of envelopes) await late.ingestor.offer(envelope) const before = await late.ingestor.sync() assert.equal(before.goals.length, 0, 'the member’s goal is not folded yet') assert.ok(late.ingestor.outstanding().quarantined > 0) late.records.put(directoryRecord(memberDevice, '2026-02-01T00:00:01Z')) const after = await late.ingestor.sync() assert.equal(after.goals.length, 1) assert.equal(after.goals[0].messages.length, 2) assert.equal(late.ingestor.outstanding().quarantined, 0) }) it('recovers an evicted envelope through the want list rather than losing it', async () => { const { spaceUri, rootDevice, memberDevice, envelopes } = await privateCorpus() const directory = [directoryRecord(rootDevice, '2026-02-01T00:00:00Z')] const network = new MemoryPrivateBusNetwork() // A peer that holds everything, so there is somewhere for the want to be answered from. const keeper = replica(network, 'keeper', spaceUri, [ ...directory, directoryRecord(memberDevice, '2026-02-01T00:00:01Z'), ]) await keeper.ingestor.start() for (const envelope of envelopes) await keeper.ingestor.offer(envelope) await keeper.ingestor.sync() // A replica whose quarantine holds one envelope at a time: the member's records arrive while // their key is unknown, so all but the newest are evicted. const squeezed = replica(network, 'squeezed', spaceUri, directory, { quarantineLimits: { maxEnvelopes: 1, maxBytes: 1024 * 1024 }, }) await squeezed.ingestor.start() for (const envelope of envelopes) await squeezed.ingestor.offer(envelope) assert.ok(squeezed.wants.all().length > 0, 'every eviction leaves a want') // The key is published, and catch-up re-requests exactly what was evicted. squeezed.records.put(directoryRecord(memberDevice, '2026-02-01T00:00:01Z')) const recovered = await squeezed.ingestor.sync() assert.equal(recovered.goals.length, 1) assert.equal(recovered.goals[0].messages.length, 2) assert.deepEqual(squeezed.wants.all(), [], 'a satisfied want is cleared') assert.deepEqual(indexDigest(recovered), indexDigest(await keeper.ingestor.sync())) }) it('polls members’ repos for the directory and nothing else', async () => { const { spaceUri, rootDevice, memberDevice, envelopes } = await privateCorpus() // A member's repo holding what a real repo holds: their device record, and — because whoever // holds the repo can put anything in it — a grant making a stranger an admin of this space. const forgedGrant = { uri: `at://${MEMBER}/${COLLECTIONS.addMember}/forged`, cid: 'cid-forged-grant', value: { $type: COLLECTIONS.addMember, space: { uri: spaceUri, cid: envelopes.find((envelope) => envelope.collection === COLLECTIONS.space).recordCid, }, did: 'did:plc:attacker', kind: 'human', role: 'admin', createdAt: '2026-03-01T09:00:00Z', }, } const deviceEntry = (signer) => ({ uri: `at://${signer.did}/${COLLECTIONS.device}/${signer.keyId}`, cid: `cid-device-${signer.keyId}`, value: { $type: COLLECTIONS.device, deviceKeyId: signer.keyId, publicKey: signer.publicKey, algorithm: 'ed25519', kind: 'node', createdAt: '2026-02-01T00:00:00Z', }, }) const listed = new Set() const transport = { async resolvePds() { return new URL('https://pds.test') }, async getLatestCommit() { return { rev: '0001', commitCid: 'head-1' } }, async listRecords({ did, collection }) { listed.add(collection) if (collection === COLLECTIONS.device) { const signer = did === ROOT ? rootDevice : memberDevice return { records: [deviceEntry(signer)] } } if (collection === COLLECTIONS.addMember && did === MEMBER) { return { records: [forgedGrant] } } return { records: [] } }, async getRecord() { throw new Error('unused') }, } const network = new MemoryPrivateBusNetwork() const records = new MemoryRecordStore() const envelopeStore = new MemoryEnvelopeStore(records) const ingestor = new PrivateSpaceIngestor( spaceUri, network.endpoint('polled', { corpus: () => envelopeStore.all() }), envelopeStore, records, { poller: directoryPoller(transport, records, new MemorySyncStateStore(), { now: () => 'now' }), bootstrapDids: [ROOT, MEMBER], }, ) await ingestor.start() for (const envelope of envelopes) await ingestor.offer(envelope) const index = await ingestor.sync() // The directory arrived, so the corpus verifies and folds. assert.deepEqual([...listed].sort(), [...DEVICE_DIRECTORY_COLLECTIONS].sort()) assert.equal(index.goals.length, 1) assert.equal(index.devices.length, 2) // And the forged grant is not even in the store, let alone in the fold. assert.equal(records.records().some((record) => record.uri === forgedGrant.uri), false) assert.equal(index.members.some((member) => member.did === 'did:plc:attacker'), false) }) it('refuses a poller that would list anything beyond the directory', async () => { const { spaceUri } = await privateCorpus() const records = new MemoryRecordStore() const network = new MemoryPrivateBusNetwork() assert.throws( () => new PrivateSpaceIngestor( spaceUri, network.endpoint('wide'), new MemoryEnvelopeStore(records), records, { poller: new RepoPoller({}, records, new MemorySyncStateStore()) }, ), /device directory only/, ) }) it('refuses an envelope addressed to another space', async () => { const { spaceUri, rootDevice, envelopes } = await privateCorpus() const network = new MemoryPrivateBusNetwork() const one = replica(network, 'one', spaceUri, [directoryRecord(rootDevice, '2026-02-01T00:00:00Z')]) const report = await one.ingestor.offer({ ...envelopes[0], space: 'at://did:plc:other/x/y' }) assert.equal(report.admitted, 0) assert.equal(report.quarantined, 0) assert.match(report.rejected[0].reason, /different space/) }) }) describe('the private write path', () => { const spaceValue = { $type: COLLECTIONS.space, name: 'Private', description: 'No PDS holds a copy.', private: true, createdAt: '2026-03-01T00:00:00Z', } /** A founder with a device and a replica that has already read their public device record. */ async function founderReplica(network, id = 'writer') { const rootDevice = await device(ROOT, 'root-node-1') const spaceUri = `at://${ROOT}/${COLLECTIONS.space}/space` const site = replica(network, id, spaceUri, [directoryRecord(rootDevice, '2026-02-01T00:00:00Z')]) const signer = { did: rootDevice.did, deviceKeyId: rootDevice.keyId, privateKey: rootDevice.privateKey, } await site.ingestor.start() return { ...site, spaceUri, rootDevice, signer, writer: privateWriter(site.ingestor, signer) } } it('lands a write in the writer’s own store and gossips it to a peer', async () => { const network = new MemoryPrivateBusNetwork() const author = await founderReplica(network, 'author') const peer = replica(network, 'peer', author.spaceUri, [ directoryRecord(author.rootDevice, '2026-02-01T00:00:00Z'), ]) await peer.ingestor.start() const ref = await author.writer.create(COLLECTIONS.space, spaceValue, { rkey: 'space' }) assert.equal(ref.uri, author.spaceUri) assert.equal(author.records.get(ROOT, COLLECTIONS.space, 'space').cid, ref.cid) // Gossip reaches the peer's inbox; a sync drains it through the same admission gate. await peer.ingestor.sync() assert.equal(peer.records.get(ROOT, COLLECTIONS.space, 'space').cid, ref.cid) assert.equal(indexDigest(materialize(peer.records, { spaceUri: author.spaceUri })).overall, indexDigest(materialize(author.records, { spaceUri: author.spaceUri })).overall) }) it('refuses to write with a key whose device record this replica has not read', async () => { const network = new MemoryPrivateBusNetwork() const rootDevice = await device(ROOT, 'root-node-1') const spaceUri = `at://${ROOT}/${COLLECTIONS.space}/space` // No directory: the writer's own key is unknown to its own replica, which is the cold-start // ordering mistake (publish the device record, sync the directory, then write). const site = replica(network, 'cold', spaceUri, []) const writer = privateWriter(site.ingestor, { did: rootDevice.did, deviceKeyId: rootDevice.keyId, privateKey: rootDevice.privateKey, }) await assert.rejects( () => writer.create(COLLECTIONS.space, spaceValue, { rkey: 'space' }), /cannot verify its own envelope/, ) assert.equal(site.records.get(ROOT, COLLECTIONS.space, 'space'), undefined) }) it('is the same gate a stranger’s envelope passes: an invalid record never lands', async () => { const network = new MemoryPrivateBusNetwork() const author = await founderReplica(network) await assert.rejects( () => author.writer.create(COLLECTIONS.space, { ...spaceValue, name: 42 }), /name/, ) assert.equal(author.envelopes.all().length, 0) }) }) describe('the ticket’s space pin', () => { /** * A second space record at the same name, correctly signed by the founder's own key. * * Every gate-1 check passes on this envelope, and every gate-2 check would too: the founder is the * space's author, the signature verifies against a key their public repo publishes, and the record * is lexicon-valid. The only thing that makes it the wrong space is the ticket the invitee was * handed out of band — which is exactly why the pin cannot live in either gate. */ async function equivocation(spaceUri, rootDevice, name) { const rev = createTidSource(() => 1_800_000_000_000, 9)() return seal({ space: spaceUri, did: ROOT, collection: COLLECTIONS.space, rkey: 'space', record: { $type: COLLECTIONS.space, name, description: 'A second genesis at the same name.', private: true, createdAt: '2026-03-01T00:00:00Z', }, rev, deviceKeyId: rootDevice.keyId, privateKey: rootDevice.privateKey, }) } it('stops a catch-up that would fold a space record the ticket did not pin', async () => { const { spaceUri, rootDevice, envelopes } = await privateCorpus() const genuine = envelopes.find((envelope) => envelope.collection === COLLECTIONS.space) const forged = await equivocation(spaceUri, rootDevice, 'Not the space you were invited to') const directory = [directoryRecord(rootDevice, '2026-02-01T00:00:00Z')] const network = new MemoryPrivateBusNetwork() // The invitee holds the ticket, and the only space record they are offered is the wrong one. const invitee = replica(network, 'invitee', spaceUri, directory, { spacePin: { uri: spaceUri, cid: genuine.recordCid }, }) await invitee.ingestor.start() await invitee.ingestor.offer(forged) await assert.rejects(invitee.ingestor.sync(), /pins/) // Gate 1 is untouched: the envelope is stored, because what a replica STORES may not depend on // a ticket only it holds. What is refused is folding it as this space. assert.equal(invitee.envelopes.all().length, 1) assert.match(invitee.ingestor.spaceMismatch().reason, /out of band/) // And a replica with no ticket — the founder's own, or one bootstrapped some other way — folds // it without complaint. The refusal is local by construction, so no fold rule diverged. const unpinned = replica(network, 'unpinned', spaceUri, directory) await unpinned.ingestor.start() await unpinned.ingestor.offer(forged) const index = await unpinned.ingestor.sync() assert.equal(index.space.cid, forged.recordCid) assert.equal(unpinned.ingestor.spaceMismatch(), undefined) }) it('folds the pinned space, and is unmoved when the second genesis arrives beside it', async () => { const { spaceUri, rootDevice, memberDevice, envelopes } = await privateCorpus() const genuine = envelopes.find((envelope) => envelope.collection === COLLECTIONS.space) const forged = await equivocation(spaceUri, rootDevice, 'Not the space you were invited to') const network = new MemoryPrivateBusNetwork() const invitee = replica( network, 'invitee', spaceUri, [ directoryRecord(rootDevice, '2026-02-01T00:00:00Z'), directoryRecord(memberDevice, '2026-02-01T00:00:01Z'), ], { spacePin: { uri: spaceUri, cid: genuine.recordCid } }, ) await invitee.ingestor.start() for (const envelope of [...envelopes, forged]) await invitee.ingestor.offer(envelope) // The later version loses the selection (a space record is not a sanctioned rewrite), so the // pin agrees with the fold and the space works normally. const index = await invitee.ingestor.sync() assert.equal(index.space.cid, genuine.recordCid) assert.equal(index.goals.length, 1) assert.equal(invitee.ingestor.spaceMismatch(), undefined) }) it('refuses a local write into a space this replica cannot identify', async () => { const { spaceUri, rootDevice, envelopes } = await privateCorpus() const genuine = envelopes.find((envelope) => envelope.collection === COLLECTIONS.space) const forged = await equivocation(spaceUri, rootDevice, 'Not the space you were invited to') const network = new MemoryPrivateBusNetwork() const invitee = replica(network, 'invitee', spaceUri, [directoryRecord(rootDevice, '2026-02-01T00:00:00Z')], { spacePin: { uri: spaceUri, cid: genuine.recordCid }, }) await invitee.ingestor.start() await invitee.ingestor.offer(forged) const writer = privateWriter(invitee.ingestor, { did: rootDevice.did, deviceKeyId: rootDevice.keyId, privateKey: rootDevice.privateKey, }) await assert.rejects( () => writer.create(COLLECTIONS.message, { $type: COLLECTIONS.message, goal: { uri: `at://${ROOT}/${COLLECTIONS.goal}/goal`, cid: genuine.recordCid }, body: 'Written into the wrong space.', mentions: [], createdAt: '2026-03-05T00:00:00Z', }), /pins/, ) // Not written, not published, not even sealed into the store: the refusal is before the write. assert.equal(invitee.envelopes.all().length, 1) }) it('refuses a pin for a different space than the replica holds', async () => { const { spaceUri, rootDevice } = await privateCorpus() const network = new MemoryPrivateBusNetwork() assert.throws( () => replica(network, 'confused', spaceUri, [directoryRecord(rootDevice, '2026-02-01T00:00:00Z')], { spacePin: { uri: `at://${ROOT}/${COLLECTIONS.space}/other`, cid: 'bafyreiother' }, }), /a replica holds one/, ) }) }) /** * A replica on the WIRE bus: the same ingestor, over encoded frames instead of an object graph. * * The point of every test below is that nothing above the bus changes. `PrivateSpaceIngestor`, * `EnvelopeWriter` and the fold were built against `PrivateBus` rather than against * `MemoryPrivateBus`, so a real protocol underneath them is a substitution and not a port. */ function wireReplica(id, spaceUri, directory, options = {}) { const records = new MemoryRecordStore() for (const record of directory) records.put(record) const wants = new MemoryWantList() const envelopes = new MemoryEnvelopeStore(records) const blobs = new MemoryPrivateBlobStore() const peers = [] const bus = new WirePrivateBus({ spaceUri, envelopes, blobs, endpointId: id, peers: { async links() { return peers } }, authorizePeer: options.authorizePeer ?? ((endpointId) => ({ ok: true, peer: { did: ROOT, deviceKeyId: 'test-device', endpointId }, })), ...(options.limits ? { limits: options.limits } : {}), ...(options.maxRounds !== undefined ? { maxRounds: options.maxRounds } : {}), }) const ingestor = new PrivateSpaceIngestor(spaceUri, bus, envelopes, records, { wants, quarantine: new Quarantine(wants, options.quarantineLimits), bootstrapDids: [ROOT], blobs, ...(options.maxBlobsPerSync ? { maxBlobsPerSync: options.maxBlobsPerSync } : {}), }) return { id, records, bus, wants, envelopes, blobs, ingestor, peers } } /** Connect two wire replicas, with a switch a test can flip to partition them. */ function connect(left, right, options = {}) { const offline = { value: false } const link = (peer) => loopbackLink(peer.id, peer.bus, { localId: left.id, offline: () => offline.value, ...(options.duplicates ? { duplicates: options.duplicates } : {}), }) left.peers.push(link(right)) const reverse = (peer) => loopbackLink(peer.id, peer.bus, { localId: right.id, offline: () => offline.value, ...(options.duplicates ? { duplicates: options.duplicates } : {}), }) right.peers.push(reverse(left)) return offline } describe('private ingestion over the wire', () => { it('converges two replicas over encoded frames, duplicates and a partition included', async () => { const { spaceUri, rootDevice, memberDevice, envelopes } = await privateCorpus() const directory = [ directoryRecord(rootDevice, '2026-02-01T00:00:00Z'), directoryRecord(memberDevice, '2026-02-01T00:00:01Z'), ] const founder = wireReplica('endpoint-founder', spaceUri, directory) const invitee = wireReplica('endpoint-invitee', spaceUri, directory) const offline = connect(founder, invitee, { duplicates: 2 }) await founder.ingestor.start() await invitee.ingestor.start() // Half the corpus is written while the two cannot reach each other, so gossip is lost and only // catch-up can make them whole. offline.value = true for (const envelope of envelopes.slice(0, 3)) { await founder.ingestor.offer(envelope) await founder.bus.publish(spaceUri, envelope) } // Nothing reaches them — not even the space record, so there is nothing here to fold yet. assert.deepEqual(await invitee.bus.catchUp(spaceUri, { summaries: [], wants: [] }), []) assert.equal(invitee.envelopes.all().length, 0) offline.value = false // The rest gossips through, twice over, and the earlier writes come back at catch-up. for (const envelope of envelopes.slice(3)) { await founder.ingestor.offer(envelope) await founder.bus.publish(spaceUri, envelope) } const theirs = await invitee.ingestor.sync() const ours = await founder.ingestor.sync() assert.deepEqual(indexDigest(theirs), indexDigest(ours)) assert.equal(theirs.goals.length, 1) assert.equal(theirs.goals[0].messages.length, 2) // Idempotent under duplication: one envelope per version on both sides. assert.equal(invitee.envelopes.all().length, envelopes.length) assert.equal(founder.envelopes.all().length, envelopes.length) }) it('pages a cold replica through a corpus larger than one batch', async () => { const { spaceUri, rootDevice, memberDevice, envelopes } = await privateCorpus() const directory = [ directoryRecord(rootDevice, '2026-02-01T00:00:00Z'), directoryRecord(memberDevice, '2026-02-01T00:00:01Z'), ] const keeper = wireReplica('endpoint-keeper', spaceUri, directory) for (const envelope of envelopes) await keeper.ingestor.offer(envelope) const cold = wireReplica('endpoint-cold', spaceUri, directory, { limits: { maxEnvelopes: 2, maxBytes: 1024 * 1024 }, }) connect(cold, keeper) const caught = await cold.ingestor.sync() assert.equal(cold.envelopes.all().length, envelopes.length) assert.deepEqual(indexDigest(caught), indexDigest(await keeper.ingestor.sync())) }) it('resumes bounded catch-up while inventory is withheld until genesis arrives', async () => { const { spaceUri, rootDevice, memberDevice, envelopes } = await privateCorpus() const directory = [ directoryRecord(rootDevice, '2026-02-01T00:00:00Z'), directoryRecord(memberDevice, '2026-02-01T00:00:01Z'), ] const keeper = wireReplica('endpoint-keeper', spaceUri, directory, { limits: { maxEnvelopes: 1, maxBytes: 1024 * 1024 }, }) for (const envelope of envelopes) await keeper.ingestor.offer(envelope) // The responder orders by envelope identity, not arrival. Put genesis beyond several one-page // sync cycles: before ADR §30.3's empty summaries learned to retain their cursor, every cycle // restarted at the first envelope and this replica waited forever. const ordered = [...envelopes].sort((left, right) => envelopeKey(left) < envelopeKey(right) ? -1 : envelopeKey(left) > envelopeKey(right) ? 1 : 0, ) const genesisOffset = ordered.findIndex((envelope) => envelope.collection === COLLECTIONS.space) assert.ok(genesisOffset > 1, 'fixture does not put genesis beyond the per-cycle bound') const cold = wireReplica('endpoint-cold', spaceUri, directory, { maxRounds: 1, }) connect(cold, keeper) for (let cycle = 0; cycle < genesisOffset; cycle += 1) { await assert.rejects( cold.ingestor.sync(), /Space record not found/, `genesis arrived before its sorted offset on cycle ${cycle}`, ) } const opened = await cold.ingestor.sync() assert.equal(opened.space.uri, spaceUri) assert.equal(cold.envelopes.all().length, genesisOffset + 1) }) it('re-requests an evicted envelope explicitly, over the wire', async () => { const { spaceUri, rootDevice, memberDevice, envelopes } = await privateCorpus() const founderOnly = [directoryRecord(rootDevice, '2026-02-01T00:00:00Z')] const keeper = wireReplica('endpoint-keeper', spaceUri, [ ...founderOnly, directoryRecord(memberDevice, '2026-02-01T00:00:01Z'), ]) for (const envelope of envelopes) await keeper.ingestor.offer(envelope) // A replica whose quarantine holds one envelope at a time: the member's records arrive before // their key does, so all but the newest are evicted — and every eviction leaves a want. const squeezed = wireReplica('endpoint-squeezed', spaceUri, founderOnly, { quarantineLimits: { maxEnvelopes: 1, maxBytes: 1024 * 1024 }, }) connect(squeezed, keeper) for (const envelope of envelopes) await squeezed.ingestor.offer(envelope) assert.ok(squeezed.wants.all().length > 0, 'every eviction leaves a want') // Nothing in the summaries would surface those: the projections were never stored either, so // this is the case where only the explicit want-list can ask. The key is published, and the // catch-up carries the wants across the wire. squeezed.records.put(directoryRecord(memberDevice, '2026-02-01T00:00:01Z')) const recovered = await squeezed.ingestor.sync() assert.equal(recovered.goals.length, 1) assert.equal(recovered.goals[0].messages.length, 2) assert.deepEqual(squeezed.wants.all(), [], 'a satisfied want is cleared') assert.deepEqual(indexDigest(recovered), indexDigest(await keeper.ingestor.sync())) }) it('runs the write path unchanged: a command writes, and a peer folds it', async () => { const { spaceUri, rootDevice, memberDevice, envelopes } = await privateCorpus() const directory = [ directoryRecord(rootDevice, '2026-02-01T00:00:00Z'), directoryRecord(memberDevice, '2026-02-01T00:00:01Z'), ] const author = wireReplica('endpoint-author', spaceUri, directory) const peer = wireReplica('endpoint-peer', spaceUri, directory) connect(author, peer) await author.ingestor.start() await peer.ingestor.start() for (const envelope of envelopes) { await author.ingestor.offer(envelope) await peer.ingestor.offer(envelope) } await author.ingestor.sync() assert.equal((await peer.ingestor.sync()).goals.length, 1) const goalEnvelope = envelopes.find((envelope) => envelope.collection === COLLECTIONS.goal) const goalRef = { uri: `at://${goalEnvelope.did}/${goalEnvelope.collection}/${goalEnvelope.rkey}`, cid: goalEnvelope.recordCid, } const writer = privateWriter(author.ingestor, { did: memberDevice.did, deviceKeyId: memberDevice.keyId, privateKey: memberDevice.privateKey, }) const written = await writer.create(COLLECTIONS.message, { $type: COLLECTIONS.message, goal: goalRef, body: 'Written through the seam, delivered over the wire.', mentions: [], createdAt: '2026-03-01T00:02:00Z', }) assert.ok(written.uri.startsWith(`at://${MEMBER}/`)) const folded = await peer.ingestor.sync() assert.equal(folded.goals[0].messages.length, 3) assert.deepEqual(indexDigest(folded), indexDigest(await author.ingestor.sync())) }) it('serves catch-up only to a published, addressed endpoint', async () => { const { spaceUri, rootDevice, memberDevice, envelopes } = await privateCorpus() const directory = [ directoryRecord(rootDevice, '2026-02-01T00:00:00Z'), directoryRecord(memberDevice, '2026-02-01T00:00:01Z'), addressRecord(rootDevice, 'endpoint-founder'), ] let keeper const authorizePeer = (endpointId) => authorizeEndpoint(endpointId, readDeviceDirectory(keeper.records.records(), [])) keeper = wireReplica('endpoint-keeper', spaceUri, directory, { authorizePeer }) for (const envelope of envelopes) await keeper.ingestor.offer(envelope) const folded = readDeviceDirectory(keeper.records.records(), []) const retrieveAs = async (endpointId) => { const asker = wireReplica(endpointId, spaceUri, []) asker.peers.push(loopbackLink(keeper.id, keeper.bus, { localId: endpointId })) return asker.bus.catchUp(spaceUri, { summaries: [], wants: [] }) } // Topic knowledge is public. The authenticated endpoint identity and current directory decide // whether the responder reaches its corpus; the loopback transport cannot omit that decision. assert.equal((await retrieveAs('endpoint-founder')).length, envelopes.length) assert.deepEqual(await retrieveAs('endpoint-stranger'), []) assert.equal(authorizeEndpoint('endpoint-founder', folded).ok, true) assert.equal(authorizeEndpoint('endpoint-stranger', folded).status, 'unknown') assert.deepEqual( connectablePeers(folded).map((peer) => peer.endpointId), ['endpoint-founder'], ) }) it('pulls the bytes a record names, in chunks, and stops wanting them', async () => { const { spaceUri, rootDevice, memberDevice, envelopes } = await privateCorpus() const directory = [ directoryRecord(rootDevice, '2026-02-01T00:00:00Z'), directoryRecord(memberDevice, '2026-02-01T00:00:01Z'), ] const holder = wireReplica('endpoint-holder', spaceUri, directory) const asker = wireReplica('endpoint-asker', spaceUri, directory) connect(holder, asker) // A picture large enough that one frame cannot carry it: paging is the normal case, not an edge. const picture = new Uint8Array(700 * 1024) for (let index = 0; index < picture.length; index += 1) picture[index] = index % 251 const ref = await privateBlobRef(picture, 'image/png') await holder.blobs.put(picture) const imageValue = { $type: COLLECTIONS.image, blob: ref, alt: 'a private picture', createdAt: '2026-03-01T00:03:00Z', } const image = await seal({ space: spaceUri, did: memberDevice.did, collection: COLLECTIONS.image, rkey: 'image-1', record: imageValue, rev: '3ms26zx46yg2z', deviceKeyId: memberDevice.keyId, privateKey: memberDevice.privateKey, }) for (const envelope of [...envelopes, image]) await holder.ingestor.offer(envelope) await asker.ingestor.start() // The record arrives first and the bytes with it — what to ask for is derived from what landed. const report = await asker.ingestor.sync() assert.ok(report) assert.equal(asker.ingestor.lastReport.blobs.fetched, 1) assert.deepEqual(asker.blobs.get(ref.ref.$link), picture) assert.deepEqual(asker.ingestor.missingBlobs(), []) assert.equal(asker.ingestor.outstanding().blobs, 0) }) it('leaves a blob nobody holds visible and retryable, never silently absent', async () => { const { spaceUri, rootDevice, memberDevice, envelopes } = await privateCorpus() const directory = [ directoryRecord(rootDevice, '2026-02-01T00:00:00Z'), directoryRecord(memberDevice, '2026-02-01T00:00:01Z'), ] const holder = wireReplica('endpoint-holder', spaceUri, directory) const asker = wireReplica('endpoint-asker', spaceUri, directory) connect(holder, asker) const picture = new TextEncoder().encode('the bytes nobody kept') const ref = await privateBlobRef(picture, 'image/png') const image = await seal({ space: spaceUri, did: memberDevice.did, collection: COLLECTIONS.image, rkey: 'image-1', record: { $type: COLLECTIONS.image, blob: ref, createdAt: '2026-03-01T00:03:00Z', }, rev: '3ms26zx46yg2z', deviceKeyId: memberDevice.keyId, privateKey: memberDevice.privateKey, }) for (const envelope of [...envelopes, image]) await holder.ingestor.offer(envelope) await asker.ingestor.start() await asker.ingestor.sync() // Derived from the record, so it survives a restart with no want list to lose. assert.deepEqual(asker.ingestor.missingBlobs(), [ { cid: ref.ref.$link, mimeType: 'image/png', size: picture.length }, ]) assert.equal(asker.ingestor.lastReport.blobs.missing, 1) // The picture turns up on the peer later; the next ordinary cycle collects it, unasked. await holder.blobs.put(picture) await asker.ingestor.sync() assert.deepEqual(asker.blobs.get(ref.ref.$link), picture) assert.deepEqual(asker.ingestor.missingBlobs(), []) }) it('never lets an unfetchable blob starve the ones behind it', async () => { const { spaceUri, rootDevice, memberDevice, envelopes } = await privateCorpus() const directory = [ directoryRecord(rootDevice, '2026-02-01T00:00:00Z'), directoryRecord(memberDevice, '2026-02-01T00:00:01Z'), ] // One blob per cycle, so a fixed prefix of the sorted list would only ever ask for one of them. const holder = wireReplica('endpoint-holder', spaceUri, directory) const asker = wireReplica('endpoint-asker', spaceUri, directory, { maxBlobsPerSync: 1 }) connect(holder, asker) const images = [] // The unfetchable one is first BY CID, which is the order the fetcher walks: a fixed prefix of // the sorted list would ask for it, and only it, every cycle. for (const [index, body] of ['lost forever 0', 'still around'].entries()) { const bytes = new TextEncoder().encode(body) const ref = await privateBlobRef(bytes, 'image/png') images.push({ bytes, ref }) await holder.ingestor.offer( await seal({ space: spaceUri, did: memberDevice.did, collection: COLLECTIONS.image, rkey: `image-${index}`, record: { $type: COLLECTIONS.image, blob: ref, createdAt: '2026-03-01T00:03:00Z' }, rev: ['3ms26zx46yg2z', '3ms26zx46yg3z'][index], deviceKeyId: memberDevice.keyId, privateKey: memberDevice.privateKey, }), ) } for (const envelope of envelopes) await holder.ingestor.offer(envelope) // Only the second one is holdable; the first names bytes nobody kept, and always will. await holder.blobs.put(images[1].bytes) await asker.ingestor.start() for (let cycle = 0; cycle < 3; cycle += 1) await asker.ingestor.sync() assert.deepEqual(asker.blobs.get(images[1].ref.ref.$link), images[1].bytes) assert.deepEqual( asker.ingestor.missingBlobs().map((need) => need.cid), [images[0].ref.ref.$link], ) }) })