diff --git a/CLAUDE.md b/CLAUDE.md index 8817015..ff86ff7 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -115,9 +115,28 @@ sends still passes both gates. A refusal is also where a replica **learns whom t `catch-up` and `hello` carry an optional `from`, the dialer's claimed DIDs, which the bus buffers only while refusing and which the answer never reflects — the refusal a stranger gets must not change. `PrivateSpaceIngestor` drains that buffer, reads each claimed DID's *public* repo once, and -keeps it only if that repo binds the endpoint the transport authenticated. A claim is worth a bounded -number of public reads and never a directory entry; do not make it worth more, and do not settle one -anywhere but `authorizeEndpoint`. **Which connections exist** is `private/connections.ts`, decided in +keeps it only if that repo binds the endpoint the transport authenticated — and **only until the +signed genesis projection folds** (ADR §30.1; the boundary is genesis, not the first envelope, or +one stranger's gate-1-valid fragment could block discovery forever). The poll writes the polled repo's directory records into +the store every connection decision is folded from, and the store never deletes, so what an entry is +*worth* is settled at read time (ADR §30.2–30.3): `connectionDirectory` is the directory the endpoint +folds per frame — passed through only while there are no private projections, emptied when a partial +catch-up holds envelopes but no genesis, and kept to EVER-members' devices once genesis folds, so an +entry a stranger earned while the replica was cold stops counting before any corpus can be served. A +partial replica still dials through `catchUpDirectory`, but until genesis arrives its catch-up sends +no summaries or wants and `syncBlobs` asks for nothing — a blob CID names private content. The wire +retains a bounded per-peer cursor across those inventory-free cycles so genesis beyond one cycle's +page bound is still reachable. If genesis arrives before that scan ends, `catchUpDirectory` retains +only that endpoint beside the ever-member roster; the bus sends it empty-inventory catch-up only — +no gossip and no blob requests — until the cursor clears, then fresh summaries resume. This +preserves recovery when the serving peer's grant lies beyond genesis without disclosing the partial +corpus's shape. "Ever" +because the non-member `authorizeEndpoint` tolerates is a *removed* member, who already holds the +corpus. The membership rule is the materializer's own (`everMemberDids` shares `foldMembership`); a +claim-sourced DID is provisional in the ingestor and evicted from the poll list at the first fold +that does not count it. A claim is worth a bounded number of public reads; do not make it worth +more, and do not hand a connection decision a raw `readDeviceDirectory` unless the bootstrap footing +is the point. **Which connections exist** is `private/connections.ts`, decided in `core` for the same reason the wire is: `PeerConnections` re-reads `connectablePeers()` every round (never cached, so newly published devices and address changes are visible), orders it by `join.peerHints` without ever filtering on them, backs off deterministically, and drops diff --git a/docs/adr-private-mode-iroh.md b/docs/adr-private-mode-iroh.md index 994bcb1..4afe6f8 100644 --- a/docs/adr-private-mode-iroh.md +++ b/docs/adr-private-mode-iroh.md @@ -1306,73 +1306,218 @@ trusting a machine. Nor is there repair — an attested cutoff, where a survivin versions that keep counting so a removal need not be all-or-nothing, remains the shape any future revocation should take, and remains unbuilt. -## 30. Decision: a replica learns whom to poll from whoever dials it, and an unnamed connection is held rather than closed +## 30. Decision: connection decisions are settled at read time, and a replica learns whom to poll from whoever dials it *Recorded after "a fresh browser can only open a private space while another browser has it open" -turned out to be a discovery deadlock with a connection lifetime at the bottom of it.* - -The report was exact, and so was the workaround: with a daemon serving a space and no other browser -running, a fresh tab reached `waiting` (§24) and stayed there. Open the space in a second browser and -it filled in seconds. Neither store, fold nor gate was involved. The chain was: - -1. **A fresh replica polls the founder and itself.** Which is all a bookmark, a ticket or a space URI - can name — and a daemon's identities are its *agent profile* DIDs, so they are in none of them. - `connectablePeers()` therefore never listed the daemon, and the tab never dialled it. -2. **The daemon does find the browser.** Its `knownDids` grows from the fold's members, the human is - a member, so its next directory poll saw the browser's freshly published `deviceAddress` and - dialled it. Discovery already worked in exactly one direction. -3. **And that dial went nowhere.** The tab authorized the inbound endpoint against its own directory, - which had never heard of it, and `PrivateEndpoint.#adopt` closed the connection. No frame carried - the dialer's identity, so the refusal taught the tab nothing about whose repo to read next. Both - sides retried forever. A second browser fixed it by publishing an address into a repo the fresh - one already polled, which is precisely the shape of the reported workaround. - -The fix is one sentence in two halves: **a dialer says who it claims to be, and a replica keeps the -connection long enough to hear it.** Everything here is a *connection* statement in §2's vocabulary — -`admission.ts` and `devices.ts` are untouched, `materialize()` counts exactly what it counted, the -digest is unchanged — so, as with §29, there is no lexicon change, no permutation coverage and **no -protocol version bump**: `from` is an optional field an older `readFrame` ignores, and a peer that -ignores it makes a worse connection decision over an identical corpus. - -- **The claim is on the frame, and it is worth nothing on arrival.** `catch-up` and `hello` gain an +turned out to be a discovery deadlock with a connection lifetime at the bottom of it — and revised +five times in review before the design stopped moving. This section states the settled design +first; §30.1–§30.3 keep their original numbers because the source cites them, each rewritten as the +rule it converged on; §30.4 is the history that produced it, kept because every hole had the same +shape and the shape is the lesson.* + +**The problem.** With a daemon serving a space and no other browser running, a fresh tab reached +`waiting` (§24) and stayed there forever. A fresh replica polls the founder and itself — all a +bookmark, a ticket or a space URI can name — and a daemon's identities are its *agent profile* +DIDs, in none of them, so the tab never dialled the daemon. The daemon *did* find the tab (its +`knownDids` grows from the fold's members) and dialled it — and the tab, authorizing the inbound +endpoint against a directory that had never heard of it, closed the connection on the accept. No +frame carried the dialer's identity, so the refusal taught the tab nothing about whose repo to read +next. Both sides retried forever. + +**The fix is one sentence in two halves:** a dialer says who it claims to be, and a replica keeps +the connection long enough to hear it. What that claim is *worth* is the entire rest of this +section, and pricing it correctly converged on one invariant: + +**What any directory entry is worth is settled at read time, by what the corpus can prove at the +moment of the decision — never by how or when the entry got into the store.** The store never +deletes, admission never checks membership, and bootstrap has to run before membership is +answerable, so any rule keyed to an entry's *provenance* eventually trusts an entry that arrived +wrong. Entries stop mattering rather than existing — the same shape as retirement (§29). Concretely +the corpus has three footings, derived from records and never from arrival time (the **signed +genesis projection** — the space projection carrying `deviceKeyIds`, i.e. the one that arrived in +an envelope — is the boundary between bootstrap and membership): + +| Footing | Inbound | Outbound | Claims | +|---|---|---|---| +| **Cold** — no private projections | raw directory | raw directory, full catch-up | spent | +| **Partial** — projections, no signed genesis | serves nobody | raw directory, but empty summaries, empty wants, **no blob requests** | spent | +| **Foldable** — signed genesis folds | ever-member directory | ever-member directory, plus an inventory-free route whose bootstrap cursor has not finished | discarded | + +A cold replica leaks nothing by dialling or serving whoever its directory names; a partial one +holds private bytes it must not describe but still needs a route to genesis; a foldable one can ask +the materializer who was ever a member, and does, per frame. The narrow outbound exception in the +last row prevents genesis itself becoming a page boundary that severs the only route to a later +membership grant; it confers no inbound authorization and carries no private inventory in either +direction. + +The mechanism, in the order a frame meets it: + +- **The claim is on the frame, and it is worth nothing on arrival.** `catch-up` and `hello` carry an optional `from` — the dialer's DIDs, bounded at 8 — which `WirePrivateBus.accept` buffers **only when it is refusing** the frame, before returning the identical refusal a stranger has always got. - The answer must not change: a claim that bought a different reply would be a way to ask an endpoint - which DIDs it already trusts. The buffer is drained by the layer that can settle it - (`PrivateSpaceIngestor`), which polls each claimed DID's **public** repo once and keeps it only if - that repo binds the endpoint the transport authenticated — `authorizeEndpoint`, the same function - every connection goes through. So a claim buys at most a bounded number of public reads, and the - most it can win is a directory entry its owner could have earned by publishing an address anyway. -- **An unnamed inbound connection is held for a bounded grace, and served nothing.** This is the half - the first implementation missed, and it is the half that makes the rest run: both iroh bindings - call `onLink` the moment they accept, so closing there destroys the round trip the claim travels in - — the frame is gone before it was read. `PrivateEndpoint` now keeps such a connection (8 at once, - 30 s each), refuses every frame over it exactly as before, and adopts it if the directory catches - up in the meantime. Nothing about what is *served* changed: the bus authorizes each frame against - the same directory, and a held link is not a lease. -- **A held connection is adopted only into a space with nothing to displace.** `PeerConnections.adopt` - replaces, deliberately — a peer that redialled has said which connection it believes it has — and - that rule is wrong for a link that has been waiting: it is older than anything adopted since, and - promoting it would close a live connection in favour of a socket the peer may already have hung up - on. Fresh connections still replace; held ones only fill a gap. -- **Hints reach the poller, and tickets carry them.** `openPrivateReplica` passed `meta.peerHints` as - `prefer`, which orders a roster and cannot extend one, so the hints ADR §27 put in the bookmark did - nothing in a browser; they are now bootstrap DIDs, as they always were in the CLI and the daemon. - A ticket may now carry `peerHints` too, minted from the fold's active addressed members, so an - invitee's first replica can dial a daemon directly instead of waiting to be dialled. The cost is - stated rather than hidden: those DIDs land in the invitee's **public** `join` bookmark, which names - the operator's agents beside a founder who was already named there. Anyone who considers that too - much can leave hints out and still converge by the paragraph above, one poll slower. + The answer must not change: a claim that bought a different reply would be a way to ask an + endpoint which DIDs it already trusts. The buffer is drained by the layer that can settle it + (`PrivateSpaceIngestor`), which polls each claimed DID's **public** repo — bounded per cycle, per + claim, and by a per-endpoint cooldown — and keeps the DID only if that repo binds the endpoint the + transport authenticated, and only on the cold or partial footing (§30.1). A kept DID is + provisional (`#claimedDids`): the first fold evicts any the fold does not count. +- **An unnamed inbound connection is held for a bounded grace, and served nothing.** Both iroh + bindings call `onLink` the moment they accept, so closing an unrecognized connection there + destroys the round trip the claim travels in — the frame is gone before it was read. + `PrivateEndpoint` keeps such a connection (8 at once, 30 s each), refuses every frame over it + exactly as before, and adopts it if the directory catches up in the meantime — and only into a + space with nothing to displace, because a held link is older than anything adopted since it began + waiting. Fresh connections still replace; held ones only fill a gap. A held link is not a lease. +- **Hints reach the poller, and tickets carry them.** `meta.peerHints` are bootstrap DIDs in every + runtime (a browser once passed them as `prefer`, which orders a roster and cannot extend one), and + a ticket may carry `peerHints` minted from the fold's active addressed members, so an invitee's + first replica can dial a daemon directly instead of waiting to be dialled. The cost is stated + rather than hidden: those DIDs land in the invitee's **public** `join` bookmark. Anyone who + considers that too much can leave hints out and still converge by the claim path, one poll slower. + +Everything here is a *connection* statement in §2's vocabulary — `admission.ts` and `devices.ts` +are untouched, `materialize()` counts exactly what it counted, the digest is unchanged — so, as +with §29, there is no lexicon change, no permutation coverage and **no protocol version bump**: +`from` is an optional field an older `readFrame` ignores, and a peer that ignores it makes a worse +connection decision over an identical corpus. **What this leaves open, deliberately.** A claim is a hint and never a handshake — there is no signature on the wire proving the dialer holds that DID, because the repo read that follows settles -it better and costs the claimer nothing to be honest about. First-open latency is bounded by the -daemon's directory-poll interval, since it must see the browser's device before it dials; that is a -tuning question and not a correctness one. And the polled PDS operator learns that *some* replica was -interested in that DID, which is the one residual disclosure this buys. - -**One note for whoever changes this next.** The bug survived a review because the test doubles were -one-directional: a loopback where the accepter could close its half while the dialer went on being -served reports every accept-side lifecycle decision as working. Both doubles now share one open flag -across a connection's two ends, and the regressions start from the honest state — a replica that has -never heard of the peer about to dial it. +it better and costs the claimer nothing to be honest about. A foldable replica therefore cannot +learn a member whose grant it has not yet folded: if the only reachable peer is that member, it +waits for a peer it does know rather than trusting the claim — the same "wait rather than guess" +the fold makes everywhere else. Fixing that means a signed handshake, which nothing yet needs. +First-open latency is bounded by the daemon's directory-poll interval, a tuning question and not a +correctness one. And the polled PDS operator learns that *some* replica was interested in that DID, +which is the one residual disclosure the claim path buys. + +### 30.1 The claim is spent only before genesis folds + +A claim buys a bounded number of public, unauthenticated repo reads — the same reads +`pollDirectory` makes of a member — and it buys them **only while this replica holds no signed +genesis projection**. Both edges of that boundary are load-bearing, and each was found by asking +what the rule would be worth to somebody it was not written for: + +**A foldable replica must spend nothing**, because the poll is a write. Reading a claimed DID's +repo puts that repo's `device` and `deviceAddress` into the same `RecordStore` every connection +decision is folded from. The attack that follows needs no secrets: a space's topic is derived from +its public `at://` URI (`wire.ts`); anybody may publish a `device` and `deviceAddress` to their own +PDS, since binding a key is independent of membership (§3); so a stranger dials, is refused, claims +*their own* DID — which verifies, being true and irrelevant — and `authorizeEndpoint`, deliberately +not a membership check, says `ok` to the next dial. A private corpus must not be enumerable by +whoever computed the topic, and for the price of two public records it had become exactly that. +"Not a membership check" is safe only while every entry the directory holds got there by being +polled as a member (the non-member it tolerates is a *removed* member, who already holds the +corpus); polling on a stranger's say-so is what broke that premise. Note that "dialable but not +servable" is not a repair — a dialer sends its `summaries`, so a stranger on the roster learns the +corpus's shape without ever being served. + +**A partial replica must still spend**, because the boundary is genesis and not the first envelope. +Admission deliberately does not check membership, so a gate-1-valid fragment from whoever caught +the cold window is exactly the kind of thing a partial replica may hold — and if holding *any* +envelope ended the bootstrap, one junk envelope would end discovery forever: every later claim +drained and discarded, the one legitimate holder never found, the replica stuck the moment its +known peers go offline. The spend stays leak-free on the partial footing because that footing +serves nobody and describes nothing (§30.3). The deadlock this section opened with is a cold tab by +construction, and it still resolves: the cycle that spends the claim is the same one that throws +`Space record not found`, which is §24's `waiting` and not a failure. + +After genesis folds a claim could only re-poll a member (`knownDids` already holds every active +one) or introduce a never-member, which is the case that must not exist — so claims are drained and +discarded, and what a stranger's earlier entry is worth is §30.2's question. + +### 30.2 What an entry is worth is settled at read time + +The claim-spend condition alone is not enough, because the poll's side effect is a **directory +entry**, and the store never deletes anything. A stranger who caught the bootstrap window — polled +in while the replica was cold, in the very cycle that warmed it — would otherwise stay warm +forever: in `knownDids`, re-polled every cycle, dialled by `connectablePeers`, authorized by +`authorizeEndpoint`, served the corpus for the rest of the replica's life. A condition on the +*spend* narrows the violation to a race; it cannot remove it, because the entry outlives the state +that admitted it. + +**The repair is `connectionDirectory`** (`core/src/private/peers.ts`): the directory every +connection decision reads — `PrivateEndpoint.attach` folds it per frame, exactly where +`readDeviceDirectory` was folded before — with the three footings above as its cases. Once the +signed genesis projection folds, only devices of DIDs that have EVER been members remain. "Ever" +and not "active", because the non-member `authorizeEndpoint` was always written for is a *removed* +member: they already hold the corpus, so refusing them would disclose the removal without +protecting anything. A never-member with a directory entry — however the entry got there — is out +of the roster and out of the serve path. + +The membership rule is the materializer's own: `everMemberDids` shares `foldMembership` with +`materialize()` rather than paraphrasing it, folded over only the collections membership can read +so a per-frame caller pays a scan and not a validation of the whole store. Because the filter is +evaluated per frame against the live store, there is no window between "genesis folds" and +"refuses the stranger": an envelope and its projection land in one transaction (§17), so the frame +after genesis folds is refused already. The ingestor's part is hygiene rather than safety — the +first fold evicts a claim-sourced DID the fold does not count, so the stranger's repo also stops +being re-polled. Their records stay in the store — records are never deleted — and stop +*mattering* instead. + +The symmetry is the point: a foldable replica cannot learn a member whose grant it has not folded +(§30.1's giveup), and it will not dial or serve a device whose owner's grant it has not folded +either. The invitee's path is untouched — their grant is an envelope in the corpus every serving +peer already holds, and the bootstrap dance (refusal, claim, poll, dial-back) runs on footings +where nothing is filtered. + +### 30.3 The partial footing: hold the route, disclose nothing + +"No space record folds" and "the replica holds no corpus" are different states, and conflating +them was a hole. The transaction guarantee (§17) makes one envelope and **its own** projection +atomic — it says nothing about order across envelopes. Gossip may deliver a work envelope before +genesis, and bounded catch-up may stop on a page whose ordering has not yet reached genesis. Such +a replica holds private bytes while membership is still unanswerable, and passing its directory +through unfiltered would let a stranger polled during the truly cold window fetch that partial +corpus. + +So the partial footing keeps the route and closes every mouth: inbound authorization serves +nobody; outbound catch-up still dials through the raw directory — that is how genesis arrives — +but sends empty summaries and empty wants; and `syncBlobs` asks for **nothing**, because a +`blob-request` names the CID of bytes a private record carries, which is an identifier of private +content handed to whoever is on the other end of an unfiltered dial. (The blob gate was the last +hole found: summaries and wants were withheld while blob requests still went out, same footing, +same disclosure.) The store still derives what is missing — `missingBlobs` reads records, not the +wire — so member routes can fetch it after genesis, while a route retained only by bootstrap +progress is skipped until its scan ends. + +Empty summaries cannot resume a bounded catch-up by themselves: if every cycle restarted without a +cursor, a genesis envelope beyond the per-cycle page bound would never be reached. While inventory +is withheld, `WirePrivateBus` therefore retains the responder's opaque cursor per peer across +cycles. The cursor describes only progress through that peer's corpus, is bounded like every other +peer-allocated table, and survives genesis when that envelope lands before the bounded scan reaches +its end. `catchUpDirectory` keeps only endpoints with such cursors beside the ever-member roster, +and the bus treats each as outbound-bootstrap-only: its catch-up stays inventory-free, gossip and +blob requests skip it, and inbound authorization remains the ever-member directory. When its cursor +clears, the next connection round removes it unless the envelopes just fetched establish its +membership. Other member peers receive fresh summaries normally throughout. The completed peer gets +fresh summaries on its next round, restarting its scan so a record that landed behind the old cursor +is not skipped. + +The footing is derived from records, never arrival time: a projection carrying `deviceKeyIds` +proves a private envelope is held, and the *space* projection carrying `deviceKeyIds` marks the +transition to foldable. This preserves the bootstrap route without turning a partially filled +replica into a server for whoever caught the claim window. + +### 30.4 How this section got here + +Five findings, in order: the deadlock itself (shipped against `c771e2a` with `c382c96`'s held +connections); a claim spent by a warm replica made the corpus enumerable (§30.1, found the next +day); the entry outliving the cold state that admitted it (§30.2, against `30bf7bf`); genesis +arriving non-atomically with the corpus (§30.3, against `a000ea8`); and, from a final review, the +blob-request disclosure and a claim boundary of "any envelope" that let one stranger's fragment +block discovery permanently — both folded into §30.1 and §30.3 above. Every hole had the same +shape: a rule keyed to how state *arrived* (a claim was trusted for being verified, an entry for +being polled while cold, a corpus for usually arriving genesis-first) where the invariant had to be +keyed to what the records *prove at read time*. That is the test for whatever changes this next: +if a proposed rule mentions when or how something got into the store, it is the next hole. + +The reviews also kept finding the same two fixture bugs, worth naming because each let a real hole +survive a green suite. **One-directional doubles:** a loopback where the accepter could close its +half while the dialer went on being served reports every accept-side lifecycle decision as working +— both doubles now share one open flag across a connection's two ends. **Fixtures that arrange the +broken state:** the §30 regressions offered the corpus before syncing and then asserted a stranger +*was* authorized; UI fixtures served peers whose corpus was one space record to a tab no grant +ever named. The fixtures are now cold by default, `corpus: true` is the warm replica that must +refuse, and the grants are the ones a real peer would hold. When a test has to arrange a state to +make an assertion pass, check whether the state is one the system is meant to be in. diff --git a/packages/core/src/materializer.ts b/packages/core/src/materializer.ts index 33bcd28..254a2c5 100644 --- a/packages/core/src/materializer.ts +++ b/packages/core/src/materializer.ts @@ -31,9 +31,11 @@ import { } from './generated/records.js' import { claimContractViolation, claimDeadline, claimLive } from './claim-lease.js' import { + DEVICE_DIRECTORY_COLLECTIONS, deviceIgnoreReason, deviceViews, readDeviceDirectory, + type DeviceDirectory, type DeviceGate, type DeviceView, } from './devices.js' @@ -374,30 +376,27 @@ function foldMembership( return { memberState, authorizedEvents } } -function trustRecords( - records: StoredRecord[], +/** + * The half of `trustRecords` that answers "whose records count here": locate the space record, run + * the device gate, fold membership. Extracted so `everMemberDids` below asks the identical question + * `materialize` asks — one membership rule, not a connection-side paraphrase of it — and returning + * `undefined` rather than throwing when the space record is absent, because for one caller that + * absence is an answer (a replica still bootstrapping) and not an error. + */ +function membershipOver( + valid: StoredRecord[], spaceUri: string, -): { - space: IndexedRecord - private: boolean - records: AnyIndexed[] - members: MemberView[] - devices: DeviceView[] - ignored: IgnoredRecord[] -} { - const ignored: IgnoredRecord[] = [] - const valid = records.filter((record) => { - const result = validateRecord(record.collection, record.value) - if (!result.success) { - ignored.push({ - uri: record.uri, - reason: `invalid record: ${result.issues.map((issue) => `${issue.path} ${issue.message}`).join(', ')}`, - }) - return false + ignored: IgnoredRecord[], +): + | { + directory: DeviceDirectory + spaceRecord: StoredRecord + isPrivate: boolean + admissible: StoredRecord[] + memberState: Map + authorizedEvents: Set } - return true - }) - + | undefined { // Gate 2, first half: which signing keys each DID has published. Read from the public path only, // and independent of membership — see `devices.ts` for // why the directory cannot depend on the fold it feeds. @@ -408,12 +407,11 @@ function trustRecords( const spaceRecord = valid.find( (record) => record.collection === COLLECTIONS.space && record.uri === spaceUri, ) - if (!spaceRecord) throw new Error(`Space record not found: ${spaceUri}`) - const space = asIndexed(spaceRecord, 'admin') - const rootDid = space.did + if (!spaceRecord) return undefined // Signed OR declared: a space record that arrived in an envelope is a private space's by // construction, and reading the flag alone would let a public copy of it turn the gate off. - const isPrivate = space.value.private === true || spaceRecord.deviceKeyIds !== undefined + const isPrivate = + (spaceRecord.value as SpaceRecord).private === true || spaceRecord.deviceKeyIds !== undefined const gate: DeviceGate = { private: isPrivate, spaceUri, @@ -426,9 +424,76 @@ function trustRecords( const { memberState, authorizedEvents } = foldMembership( admissible, spaceRecord, - rootDid, + spaceRecord.did, ignored, ) + return { directory, spaceRecord, isPrivate, admissible, memberState, authorizedEvents } +} + +/** + * Every DID that has ever held membership in this space — the root implicitly, and every grant + * recipient whether or not a removal followed — or `undefined` when the record set holds no space + * record to fold membership against, which is a replica still bootstrapping and not an empty set. + * + * This exists for the connection layer (`private/peers.ts`), which must not serve a private corpus + * to a DID whose only relationship to the space is a directory entry (ADR §30.2) — and the reason it + * lives HERE is the reason `foldMembership` is a named function: the membership rule exists once. + * "Ever" rather than "active" is deliberate: a removed member already holds the corpus, so refusing + * their connection would disclose the removal without protecting anything (`authorizeEndpoint`). + * + * Folded over only the collections membership can read — the space record, the grants and removals, + * and the device directory the gate verifies them against — so a caller evaluating it per frame + * pays for a scan of those records and not a validation of the whole store. + */ +export function everMemberDids( + records: readonly StoredRecord[], + spaceUri: string, +): Set | undefined { + const relevant: readonly string[] = [ + COLLECTIONS.space, + COLLECTIONS.addMember, + COLLECTIONS.removeMember, + ...DEVICE_DIRECTORY_COLLECTIONS, + ] + const ignored: IgnoredRecord[] = [] + const valid = records.filter( + (record) => + relevant.includes(record.collection) && + validateRecord(record.collection, record.value).success, + ) + const base = membershipOver([...valid], spaceUri, ignored) + return base && new Set(base.memberState.keys()) +} + +function trustRecords( + records: StoredRecord[], + spaceUri: string, +): { + space: IndexedRecord + private: boolean + records: AnyIndexed[] + members: MemberView[] + devices: DeviceView[] + ignored: IgnoredRecord[] +} { + const ignored: IgnoredRecord[] = [] + const valid = records.filter((record) => { + const result = validateRecord(record.collection, record.value) + if (!result.success) { + ignored.push({ + uri: record.uri, + reason: `invalid record: ${result.issues.map((issue) => `${issue.path} ${issue.message}`).join(', ')}`, + }) + return false + } + return true + }) + + const base = membershipOver(valid, spaceUri, ignored) + if (!base) throw new Error(`Space record not found: ${spaceUri}`) + const { directory, spaceRecord, isPrivate, admissible, memberState, authorizedEvents } = base + const space = asIndexed(spaceRecord, 'admin') + const rootDid = space.did const trusted: AnyIndexed[] = [] for (const record of admissible) { diff --git a/packages/core/src/private/bus.ts b/packages/core/src/private/bus.ts index a0883ac..d27e61e 100644 --- a/packages/core/src/private/bus.ts +++ b/packages/core/src/private/bus.ts @@ -32,6 +32,14 @@ export interface CatchUpRequest { export interface PrivateBus { /** Untrusted identities claimed by refused inbound peers; callers must verify them publicly. */ drainPeerCandidates?(): Array<{ endpointId: string; dids: string[] }> + /** + * Endpoints whose inventory-free bootstrap scan stopped at a cursor. + * + * A transport-backed endpoint uses this only to retain those specific outbound routes after + * genesis folds, until the bounded scan reaches its end. They are not authorized, gossiped to or + * asked for blobs merely by appearing here. + */ + bootstrapPeers?(): readonly string[] /** Join a space's topic. Idempotent. */ join(spaceUri: string): Promise /** Gossip one envelope to every connected peer in the space. */ diff --git a/packages/core/src/private/connections.ts b/packages/core/src/private/connections.ts index 79f1b26..fd27233 100644 --- a/packages/core/src/private/connections.ts +++ b/packages/core/src/private/connections.ts @@ -77,7 +77,9 @@ export interface PeerConnectionsOptions { spaceUri: string /** * The devices this replica may dial, as the directory currently describes them — normally - * `connectablePeers(readDeviceDirectory(...))`. + * `connectablePeers(catchUpDirectory(...))`, so a never-member is not dialled once a corpus has + * folded except while its bounded inventory-free bootstrap cursor is unfinished, and a partial + * replica retains a route to fetch genesis (ADR §30.2–30.3). * * Called on **every** round rather than cached, because that is what makes a retirement close a * connection rather than merely refuse the next one. It should therefore be cheap: a caller that diff --git a/packages/core/src/private/endpoint.ts b/packages/core/src/private/endpoint.ts index 4e19c88..7380520 100644 --- a/packages/core/src/private/endpoint.ts +++ b/packages/core/src/private/endpoint.ts @@ -40,11 +40,17 @@ // bounded grace instead, during which it is refused frame by frame exactly as before and is // adopted the moment the directory does name it. -import { readDeviceDirectory, type DeviceDirectory } from '../devices.js' +import type { DeviceDirectory } from '../devices.js' import { compareCodePoints } from '../order.js' import { publishDeviceAddress, type AddressWriter } from './directory-writes.js' import type { EnvelopeStore } from './envelope-store.js' -import { authorizeEndpoint, connectablePeers, type PeerIdentity } from './peers.js' +import { + authorizeEndpoint, + catchUpDirectory, + connectablePeers, + connectionDirectory, + type PeerIdentity, +} from './peers.js' import { PeerConnections } from './connections.js' import { WirePrivateBus, type FrameEndpoint, type PeerLink } from './sync.js' import { errorFrame, spaceTopic, type PrivateFrame } from './wire.js' @@ -277,18 +283,34 @@ export class PrivateEndpoint { // Folded on every call rather than cached, which costs a scan of the replica per round and is // the trade §16 asks for: a cached roster leaves a device the directory has stopped naming // connected for as long as the cache lives. The store has no version to memoise on, and - // inventing one here would be inventing a way for it to be wrong. - const directory = (): DeviceDirectory => readDeviceDirectory(options.records.records(), []) + // inventing one here would be inventing a way for it to be wrong. The CONNECTION directory, + // not the raw one: once this store folds a space record, a device whose owner has never been a + // member is never served and is dialled only to finish an already-started inventory-free scan, + // however its records got here (ADR §30.2–30.3). Because this closure is re-read per frame, the + // serve filter takes effect the moment the corpus does. + const directory = (): DeviceDirectory => + connectionDirectory(options.records.records(), options.spaceUri, []) + // A partial replica serves nobody until genesis makes membership foldable, but it must retain + // outbound routes to fetch that genesis. If a bounded scan stops beyond genesis, retain only + // that peer until its cursor clears; the bus still sends it no inventory, gossip or blob CIDs. + let bus: WirePrivateBus | undefined + const catchUpPeers = (): DeviceDirectory => + catchUpDirectory( + options.records.records(), + options.spaceUri, + [], + new Set(bus?.bootstrapPeers() ?? []), + ) const connections = new PeerConnections({ spaceUri: options.spaceUri, // Re-read every round — a device published since the last poll is dialable on the next one, // and a device that stopped publishing an address has its link closed (§16). - roster: () => connectablePeers(directory()), + roster: () => connectablePeers(catchUpPeers()), dial: (peer) => this.#transport.dial(peer), self: () => [this.#transport.endpointId], ...(options.prefer ? { prefer: options.prefer } : {}), }) - const bus = new WirePrivateBus({ + bus = new WirePrivateBus({ spaceUri: options.spaceUri, envelopes: options.envelopes, peers: connections, diff --git a/packages/core/src/private/peers.ts b/packages/core/src/private/peers.ts index ec6fa32..0d5665d 100644 --- a/packages/core/src/private/peers.ts +++ b/packages/core/src/private/peers.ts @@ -29,9 +29,11 @@ // is deliberately observer-local, because it is what makes "why did that peer refuse me?" answerable // from a directory dump rather than from a log on the other machine. -import type { DeviceDirectory } from '../devices.js' -import { deviceKey } from '../devices.js' +import type { DeviceBinding, DeviceDirectory } from '../devices.js' +import { deviceKey, readDeviceDirectory } from '../devices.js' +import { everMemberDids } from '../materializer.js' import { compareCodePoints } from '../order.js' +import type { StoredRecord } from '../store.js' /** A device a replica can dial, as the public directory currently describes it. */ export interface PeerIdentity { @@ -64,6 +66,100 @@ export type PeerAuthorization = | { ok: true; peer: PeerIdentity } | { ok: false; status: PeerRefusal; reason: string } +/** + * The directory as CONNECTION decisions must read it (ADR §30.2). + * + * `readDeviceDirectory`, then — once this record set holds a space record to fold membership + * against — kept to the devices of DIDs that have ever been members here. `authorizeEndpoint` reads + * this result for every inbound frame. Outbound catch-up uses `catchUpDirectory` below so a partial + * replica can recover genesis without serving or describing what it already holds. + * + * The three states are the three footings a replica can be on: + * + * - **No private projections — bootstrapping.** The directory passes through unfiltered, because + * a replica holding no corpus has nothing to leak by dialling or serving whoever its directory + * names (ADR §30). + * - **Private projections but no space record — partial catch-up.** Nobody is served. Envelope + * storage is atomic with EACH projection, not with the genesis projection, so gossip or a + * bounded catch-up may put private records here before the space record that lets membership + * fold. `catchUpDirectory` below still lets the replica dial for the missing genesis, while the + * ingestor withholds summaries and wants until it arrives (ADR §30.3). + * - **Space record present — a corpus.** Only ever-members remain. "Ever" and not "active", + * because the non-member `authorizeEndpoint` is written for is a REMOVED member: they already + * hold the corpus, so refusing their connection would disclose the removal without protecting + * anything. A DID with a directory entry and no grant — a stranger polled off a claim while + * this replica was cold — is exactly what this exists to exclude. + * + * A pure function of the record set, like everything else in this file, so two replicas holding the + * same records refuse the same peers. + */ +export function connectionDirectory( + records: readonly StoredRecord[], + spaceUri: string, + ignored: Array<{ uri: string; reason: string }>, +): DeviceDirectory { + const directory = readDeviceDirectory(records, ignored) + const members = everMemberDids(records, spaceUri) + if (!members) { + // A projection carrying signers proves this replica holds at least one private envelope. Without + // genesis there is no membership answer, so the only safe serve roster is nobody. Do not confuse + // this with a truly cold replica: per-envelope atomicity does not make genesis arrive first. + if (!records.some((record) => record.deviceKeyIds !== undefined)) return directory + for (const binding of directory.bindings.values()) { + ignored.push({ + uri: binding.uri, + reason: 'space membership is unavailable while private envelopes are present', + }) + } + return { bindings: new Map(), addresses: new Map(), retired: new Map() } + } + return memberDirectory(directory, members, ignored) +} + +/** + * The directory used only to dial OUT for catch-up. + * + * A partial replica must refuse inbound frames (`connectionDirectory`) but must keep a route to the + * peer that can supply its missing genesis. Until genesis arrives this therefore returns the raw + * public directory. The caller must reveal no local inventory on those bootstrap requests; + * `PrivateSpaceIngestor.sync()` sends empty summaries and wants on exactly this footing (ADR §30.3). + * Once membership can fold, dialing uses the same ever-member filter as serving, plus only the + * endpoints whose bounded inventory-free scan still has a cursor. `WirePrivateBus` exposes those + * routes and withholds gossip, blobs and inventory from them until the cursor clears. + */ +export function catchUpDirectory( + records: readonly StoredRecord[], + spaceUri: string, + ignored: Array<{ uri: string; reason: string }>, + bootstrapEndpoints: ReadonlySet = new Set(), +): DeviceDirectory { + const directory = readDeviceDirectory(records, ignored) + const members = everMemberDids(records, spaceUri) + return members ? memberDirectory(directory, members, ignored, bootstrapEndpoints) : directory +} + +const memberDirectory = ( + directory: DeviceDirectory, + members: ReadonlySet, + ignored: Array<{ uri: string; reason: string }>, + bootstrapEndpoints: ReadonlySet = new Set(), +): DeviceDirectory => { + const bindings = new Map() + for (const [key, binding] of directory.bindings) { + const address = directory.addresses.get(key) + if (members.has(binding.did) || (address && bootstrapEndpoints.has(address.endpointId))) { + bindings.set(key, binding) + continue + } + ignored.push({ uri: binding.uri, reason: 'device owner has never been a member of this space' }) + } + return { + bindings, + addresses: new Map([...directory.addresses].filter(([key]) => bindings.has(key))), + retired: new Map([...directory.retired].filter(([key]) => bindings.has(key))), + } +} + /** * Every device this replica could dial: bound and currently addressed. * @@ -103,6 +199,15 @@ export function connectablePeers( * removed member's records are on the public path. And it is not sufficient: whatever an authorized * peer sends still goes through `admit()`, which is why an out-of-date directory here costs a * reconnect and never a wrong record. + * + * **What makes the first of those safe is a premise about the directory this is handed, not about + * this function.** The non-member this is written for is a *removed* member: serving somebody who + * already holds the corpus discloses nothing. A DID whose only relationship to the space is a + * directory entry — a stranger polled off a claim while the replica was cold (ADR §30) — must + * therefore never reach this scan, and `connectionDirectory` above is where that is enforced: it + * serves nobody on a partial corpus and, once a space record folds, keeps only ever-members + * (ADR §30.2–30.3). Hand this function a raw `readDeviceDirectory` only where the bootstrap footing + * is the point. */ export function authorizeEndpoint( endpointId: string, diff --git a/packages/core/src/private/sync.ts b/packages/core/src/private/sync.ts index 40ce629..f4bd504 100644 --- a/packages/core/src/private/sync.ts +++ b/packages/core/src/private/sync.ts @@ -36,6 +36,7 @@ import type { EnvelopeStore, VersionSummary } from './envelope-store.js' import { envelopeKey, type PrivateEnvelope } from './envelope.js' import type { CatchUpRequest, PrivateBus } from './bus.js' import type { PeerAuthorization } from './peers.js' +import { compareCodePoints } from '../order.js' import { WIRE_VERSION, decodeFrame, @@ -285,12 +286,17 @@ export interface WirePrivateBusOptions { * * A bound rather than "until the cursor is absent", because the cursor is the *peer's* claim about * whether there is more, and a peer that always says yes would otherwise own this replica's sync - * loop. Stopping early costs a round; the next `sync()` resumes with fresh summaries. + * loop. Stopping early costs a round. A normal peer starts the next one from fresh summaries; a + * peer whose bootstrap scan is deliberately withholding inventory resumes from its last cursor + * instead, even if genesis landed meanwhile, because empty summaries cannot otherwise move past + * the same bounded prefix. */ maxRounds?: number } const DEFAULT_MAX_ROUNDS = 32 +/** A bootstrap cursor is tiny, but endpoint ids are still a peer-controlled namespace. */ +const MAX_BOOTSTRAP_CURSORS = 64 /** * A `PrivateBus` over frames. @@ -315,6 +321,8 @@ export class WirePrivateBus implements PrivateBus { readonly #blobLimits: BlobLimits readonly #subscribers = new Map void>>() readonly #status = new Map() + /** Per-peer progress while an inventory-free bootstrap scan is unfinished (ADR §30.3). */ + readonly #bootstrapCursors = new Map() #topic: string | undefined #closed = false readonly #candidates = new Map() @@ -340,6 +348,17 @@ export class WirePrivateBus implements PrivateBus { return this.#topic } + /** + * Routes whose inventory-free scan has not reached the end yet. + * + * Connection management may keep these endpoints dialable after genesis folds, but the bus keeps + * treating them as bootstrap-only: no gossip, no blobs, and no local inventory until their cursor + * clears. Sorted so the answer is diagnostic as well as actionable. + */ + bootstrapPeers(): readonly string[] { + return [...this.#bootstrapCursors.keys()].sort(compareCodePoints) + } + #assertSpace(spaceUri: string): void { if (spaceUri !== this.spaceUri) { // One bus, one space — the same shape a replica directory has (design §18). A caller passing @@ -389,6 +408,9 @@ export class WirePrivateBus implements PrivateBus { const topic = await this.#topicFor(spaceUri) const frame: PrivateFrame = { kind: 'gossip', version: WIRE_VERSION, topic, envelope } for (const link of await this.#links(spaceUri)) { + // A retained bootstrap route is not authorization. Until its inventory-free scan ends it gets + // only the catch-up requests that advance that scan, never newly written private content. + if (this.#bootstrapCursors.has(link.id)) continue const status = this.#peerStatus(link.id) try { await link.send(frame) @@ -412,13 +434,23 @@ export class WirePrivateBus implements PrivateBus { async catchUp(spaceUri: string, request: CatchUpRequest): Promise { this.#assertSpace(spaceUri) const topic = await this.#topicFor(spaceUri) + // Before signed genesis lands, the ingestor sends neither summaries nor wants: disclosing either + // would describe a partial private corpus to the raw bootstrap directory (ADR §30.3). A cursor + // is progress through the PEER's corpus and says nothing about ours, so retain it across bounded + // calls. If genesis lands before that scan ends, only that peer continues inventory-free; other, + // proven members receive the caller's fresh summaries as normal. // Keyed rather than appended: several peers legitimately hold the same envelope, and returning // one copy per peer would make the cost of joining a space grow with the number of members. const collected = new Map() for (const link of await this.#links(spaceUri)) { const status = this.#peerStatus(link.id) status.received = 0 - let cursor: string | undefined + const retained = this.#bootstrapCursors.get(link.id) + const inventoryWithheld = + retained !== undefined || (request.summaries.length === 0 && request.wants.length === 0) + const summaries = inventoryWithheld ? [] : request.summaries + const wants = inventoryWithheld ? [] : request.wants + let cursor = retained for (let round = 0; round < this.#maxRounds; round += 1) { let answer: PrivateFrame try { @@ -426,8 +458,8 @@ export class WirePrivateBus implements PrivateBus { kind: 'catch-up', version: WIRE_VERSION, topic, - summaries: request.summaries, - wants: request.wants, + summaries, + wants, ...(this.#selfDids().length > 0 ? { from: [...new Set(this.#selfDids())].slice(0, 8) } : {}), ...(cursor !== undefined ? { cursor } : {}), }) @@ -441,19 +473,25 @@ export class WirePrivateBus implements PrivateBus { } if (answer.kind !== 'envelopes' || answer.topic !== topic) { this.#fail(status, 'peer answered a catch-up with something else') + this.#bootstrapCursors.delete(link.id) break } for (const envelope of answer.envelopes) { collected.set(envelopeKey(envelope), envelope) status.received += 1 } - if (answer.cursor === undefined) break + if (answer.cursor === undefined) { + this.#bootstrapCursors.delete(link.id) + break + } // A peer that repeats a cursor is not making progress; stopping is cheaper than trusting it. if (answer.cursor === cursor) { this.#fail(status, 'peer repeated a catch-up cursor') + this.#bootstrapCursors.delete(link.id) break } cursor = answer.cursor + if (inventoryWithheld) this.#rememberBootstrapCursor(link.id, cursor) } } // Sorted so two replicas offered the same set process it in the same order — the admission @@ -464,6 +502,15 @@ export class WirePrivateBus implements PrivateBus { .map(([, envelope]) => envelope) } + /** Keep bootstrap progress bounded even when authorized endpoints churn before genesis arrives. */ + #rememberBootstrapCursor(endpointId: string, cursor: string): void { + this.#bootstrapCursors.delete(endpointId) + this.#bootstrapCursors.set(endpointId, cursor) + while (this.#bootstrapCursors.size > MAX_BOOTSTRAP_CURSORS) { + this.#bootstrapCursors.delete(this.#bootstrapCursors.keys().next().value as string) + } + } + /** * The bytes a record named, from the first peer that has them. * @@ -477,6 +524,9 @@ export class WirePrivateBus implements PrivateBus { this.#assertSpace(spaceUri) const topic = await this.#topicFor(spaceUri) for (const link of await this.#links(spaceUri)) { + // A blob CID identifies private content. A route retained only to finish an inventory-free + // scan must not learn one; once its cursor clears, the member roster decides whether it stays. + if (this.#bootstrapCursors.has(link.id)) continue const status = this.#peerStatus(link.id) const bytes = await this.#fetchFrom(link, status, topic, need) if (bytes) return bytes @@ -670,6 +720,7 @@ export class WirePrivateBus implements PrivateBus { async close(): Promise { this.#closed = true this.#subscribers.clear() + this.#bootstrapCursors.clear() } } diff --git a/packages/core/test/private-endpoint.test.mjs b/packages/core/test/private-endpoint.test.mjs index 1861fc8..209ad6b 100644 --- a/packages/core/test/private-endpoint.test.mjs +++ b/packages/core/test/private-endpoint.test.mjs @@ -92,6 +92,39 @@ function publishDevice(records, did, endpointId) { } } +/** The moment membership becomes foldable, before every later membership envelope has arrived. */ +function putGenesis(records) { + const value = { + $type: COLLECTIONS.space, + name: 'Private', + description: 'No PDS holds a copy.', + private: true, + createdAt: '2026-03-01T00:00:02Z', + } + records.put({ + did: 'did:plc:privatefounder', + collection: COLLECTIONS.space, + rkey: 'space1', + uri: SPACE, + cid: 'cid-space', + rev: '0000000000003', + deviceKeyIds: ['founder-node-1'], + value, + }) + return { + version: 1, + space: SPACE, + did: 'did:plc:privatefounder', + collection: COLLECTIONS.space, + rkey: 'space1', + record: value, + rev: '0000000000003', + recordCid: 'cid-space', + deviceKeyId: 'founder-node-1', + signature: 'test-signature', + } +} + const claimFrame = async (from) => ({ kind: 'catch-up', version: WIRE_VERSION, @@ -232,3 +265,99 @@ describe('an inbound connection no attached space names', () => { assert.equal(second.closed, true, 'a held connection outlived the endpoint holding it') }) }) + +describe('an outbound bootstrap connection at the genesis boundary', () => { + it('stays until its inventory-free cursor finishes, then returns to the member roster', async () => { + const records = new MemoryRecordStore() + publishDevice(records, AGENT, AGENT_ENDPOINT) + let requests = 0 + let sends = 0 + let keepPaging = true + const link = { + id: AGENT_ENDPOINT, + closed: false, + async request(frame) { + requests += 1 + return { + kind: 'envelopes', + version: WIRE_VERSION, + topic: frame.topic, + envelopes: [], + ...(keepPaging ? { cursor: `cursor-${requests}` } : {}), + } + }, + async send() { + sends += 1 + }, + close() { + link.closed = true + }, + } + const endpoint = await PrivateEndpoint.bind({ + transport: async () => ({ + endpointId: 'e'.repeat(64), + relays: [], + secretKey: new Uint8Array(32), + async dial(peer) { + assert.equal(peer.endpointId, AGENT_ENDPOINT) + return link + }, + async close() {}, + }), + }) + const attached = endpoint.attach({ + spaceUri: SPACE, + records, + envelopes: new MemoryEnvelopeStore(), + }) + try { + // The bounded cold scan stops with a cursor owned by this endpoint. + await attached.bus.catchUp(SPACE, { summaries: [], wants: [] }) + assert.deepEqual(attached.bus.bootstrapPeers(), [AGENT_ENDPOINT]) + const coldRequests = requests + + // Genesis lands at the page boundary, but this peer's grant may be on the next page. The + // production roster must not close the one route capable of supplying it before asking again. + const genesis = putGenesis(records) + await attached.bus.publish(SPACE, genesis) + assert.equal(sends, 0, 'a retained bootstrap route received gossip') + assert.equal( + await attached.bus.fetchBlob(SPACE, { cid: 'b'.repeat(59), mimeType: 'image/png', size: 1 }), + undefined, + ) + assert.equal(requests, coldRequests, 'a retained bootstrap route learned a blob CID') + keepPaging = false + await attached.bus.catchUp(SPACE, { + summaries: [ + { + did: 'did:plc:privatefounder', + collection: COLLECTIONS.space, + rkey: 'space1', + cids: ['cid-space'], + }, + ], + wants: [], + }) + assert.equal(requests, coldRequests + 1, 'the route was dropped when genesis folded') + assert.deepEqual(attached.bus.bootstrapPeers(), []) + + // This fixture supplied no grant, so after the scan ends the ordinary ever-member roster + // removes it. Retention is progress state, never membership or authorization. + await attached.bus.catchUp(SPACE, { + summaries: [ + { + did: 'did:plc:privatefounder', + collection: COLLECTIONS.space, + rkey: 'space1', + cids: ['cid-space'], + }, + ], + wants: [], + }) + assert.equal(requests, coldRequests + 1) + assert.equal(link.closed, true) + } finally { + await endpoint.close() + } + }) +}) diff --git a/packages/core/test/private-peers.test.mjs b/packages/core/test/private-peers.test.mjs index 19e6384..6fa0a7b 100644 --- a/packages/core/test/private-peers.test.mjs +++ b/packages/core/test/private-peers.test.mjs @@ -2,7 +2,14 @@ import assert from 'node:assert/strict' import { describe, it } from 'node:test' import { COLLECTIONS } from '../dist/generated/records.js' import { deviceKey, deviceViews, readDeviceDirectory } from '../dist/devices.js' -import { authorizeEndpoint, connectablePeers, peerEndpoints } from '../dist/private/peers.js' +import { + authorizeEndpoint, + catchUpDirectory, + connectablePeers, + connectionDirectory, + peerEndpoints, +} from '../dist/private/peers.js' +import { privateScenario } from '../dist/test/scenario.js' const ALICE = 'did:plc:alice' const BOB = 'did:plc:bob' @@ -184,3 +191,149 @@ describe('a retired device', () => { } }) }) + +// The connection directory (ADR §30.2): what `authorizeEndpoint`'s "not a membership check" rests +// on. A directory entry can be earned by a stranger — anybody may publish a `device`, and a cold +// replica polls a claimed repo (ADR §30) — so what the entry is WORTH is settled at read time: once +// a space record folds, a never-member's devices vanish from every connection decision. +describe('the connection directory', () => { + it('passes the directory through while no private envelope has landed — the bootstrap footing', () => { + // A replica holding no corpus has nothing to leak by dialling or serving whoever its directory + // names, and membership is a question this record set cannot answer yet. + const records = [ + deviceRecord(BOB, 'bob-node-1'), + addressRecord(BOB, 'bob-node-1', 'endpoint-bob'), + ] + const directory = connectionDirectory(records, `at://${ALICE}/${COLLECTIONS.space}/space`, []) + assert.deepEqual(peerEndpoints(connectablePeers(directory)), ['endpoint-bob']) + assert.equal(authorizeEndpoint('endpoint-bob', directory).ok, true) + }) + + it('serves nobody when private envelopes arrive before the space record', () => { + const scenario = privateScenario() + const endpoint = 'endpoint-stranger' + const records = [ + ...scenario.records.filter((record) => record.uri !== scenario.spaceUri), + addressRecord(scenario.strangerDid, 'stranger-1', endpoint), + ] + assert.equal( + records.some((record) => record.deviceKeyIds !== undefined), + true, + 'fixture is not a partial private corpus', + ) + + const ignored = [] + const serving = connectionDirectory(records, scenario.spaceUri, ignored) + assert.deepEqual(connectablePeers(serving), []) + assert.equal(authorizeEndpoint(endpoint, serving).ok, false) + assert.equal(authorizeEndpoint(scenario.addressEndpointId, serving).ok, false) + assert.equal( + ignored.some( + (entry) => + entry.reason === 'space membership is unavailable while private envelopes are present', + ), + true, + ) + + // Recovery is still possible: outbound catch-up retains the public routes that may supply the + // missing genesis. The ingestor regression asserts those requests disclose no inventory. + assert.equal( + peerEndpoints(connectablePeers(catchUpDirectory(records, scenario.spaceUri, []))).includes( + scenario.addressEndpointId, + ), + true, + ) + }) + + it('keeps a never-member out of every connection decision once a corpus has folded', () => { + // The stranger did everything ADR §30.1's attack does: published their own device and address + // to their own repo, and got both into this replica's store. The entry is real; it buys nothing. + const scenario = privateScenario() + const strangerAddress = addressRecord(scenario.strangerDid, 'stranger-1', 'endpoint-stranger') + const records = [...scenario.records, strangerAddress] + + // Order-independent, like the fold it borrows membership from. + for (const ordering of [records, [...records].reverse()]) { + const ignored = [] + const directory = connectionDirectory(ordering, scenario.spaceUri, ignored) + assert.equal( + peerEndpoints(connectablePeers(directory)).includes('endpoint-stranger'), + false, + 'a never-member was dialable', + ) + const refusal = authorizeEndpoint('endpoint-stranger', directory) + assert.equal(refusal.ok, false, 'a never-member was servable') + // The same `unknown` a typo earns: a refusal must not say what has been noticed. + assert.equal(refusal.status, 'unknown') + // The founder's device is untouched beside it. + assert.equal(authorizeEndpoint(scenario.addressEndpointId, directory).ok, true) + assert.equal( + ignored.some( + (entry) => entry.reason === 'device owner has never been a member of this space', + ), + true, + ) + + // A route whose inventory-free bootstrap scan stopped at a cursor survives just long enough + // to finish that scan. This is outbound-only; `connectionDirectory` above still refuses it. + const retained = catchUpDirectory( + ordering, + scenario.spaceUri, + [], + new Set(['endpoint-stranger']), + ) + assert.equal( + peerEndpoints(connectablePeers(retained)).includes('endpoint-stranger'), + true, + 'an unfinished bootstrap route was dropped at genesis', + ) + } + }) + + it('still dials and serves a removed member — they already hold the corpus', () => { + // "Ever a member" and not "active": refusing a removed member's connection would disclose the + // removal without protecting anything, and their history is still counting (design §18). + const spaceUri = `at://${ALICE}/${COLLECTIONS.space}/space` + const sealed = (did, rkey, value, cid) => ({ + did, + collection: value.$type, + rkey, + uri: `at://${did}/${value.$type}/${rkey}`, + cid, + rev: '0000000000009', + value, + deviceKeyIds: ['alice-node-1'], + }) + const grant = sealed(ALICE, 'grant-bob', { + $type: COLLECTIONS.addMember, + space: { uri: spaceUri, cid: 'cid-space' }, + did: BOB, + kind: 'human', + role: 'member', + createdAt: '2026-01-01T00:00:01Z', + }, 'cid-grant-bob') + const records = [ + deviceRecord(ALICE, 'alice-node-1'), + deviceRecord(BOB, 'bob-node-1'), + addressRecord(BOB, 'bob-node-1', 'endpoint-bob'), + sealed(ALICE, 'space', { + $type: COLLECTIONS.space, + name: 'Private', + description: 'No PDS holds a copy.', + private: true, + createdAt: '2026-01-01T00:00:00Z', + }, 'cid-space'), + grant, + sealed(ALICE, 'remove-bob', { + $type: COLLECTIONS.removeMember, + space: { uri: spaceUri, cid: 'cid-space' }, + did: BOB, + membership: { uri: grant.uri, cid: grant.cid }, + createdAt: '2026-01-01T00:00:02Z', + }, 'cid-remove-bob'), + ] + const directory = connectionDirectory(records, spaceUri, []) + assert.deepEqual(peerEndpoints(connectablePeers(directory)), ['endpoint-bob']) + assert.equal(authorizeEndpoint('endpoint-bob', directory).ok, true) + }) +}) diff --git a/packages/core/test/private-sync.test.mjs b/packages/core/test/private-sync.test.mjs index a2b0c26..4862a09 100644 --- a/packages/core/test/private-sync.test.mjs +++ b/packages/core/test/private-sync.test.mjs @@ -221,14 +221,62 @@ describe('the wire bus', () => { assert.deepEqual(names(offered), names(envelopes)) }) - it('stops paging at maxRounds rather than letting a peer own the loop', async () => { + it('bounds one call at maxRounds and resumes an inventory-free catch-up on the next', async () => { const envelopes = await corpus(5) const { left } = await pair([], envelopes, { limits: { maxEnvelopes: 1, maxBytes: 1 << 20 }, maxRounds: 2, }) - const offered = await left.bus.catchUp(SPACE, { summaries: [], wants: [] }) - assert.equal(offered.length, 2) + const first = await left.bus.catchUp(SPACE, { summaries: [], wants: [] }) + assert.deepEqual(left.bus.bootstrapPeers(), ['endpoint-right']) + const second = await left.bus.catchUp(SPACE, { summaries: [], wants: [] }) + const third = await left.bus.catchUp(SPACE, { summaries: [], wants: [] }) + assert.equal(first.length, 2, 'one peer owned more than the per-call bound') + assert.equal(second.length, 2) + assert.equal(third.length, 1) + assert.deepEqual(names([...first, ...second, ...third]), names(envelopes)) + assert.deepEqual(left.bus.bootstrapPeers(), []) + }) + + it('finishes bootstrap progress before fresh summaries find records behind the cursor', async () => { + const envelopes = await corpus(5) + const { left, right } = await pair([], envelopes, { + limits: { maxEnvelopes: 1, maxBytes: 1 << 20 }, + maxRounds: 2, + }) + const first = await left.bus.catchUp(SPACE, { summaries: [], wants: [] }) + assert.equal(first.length, 2) + + // This lands lexically behind the retained cursor. Once the inventory-free scan finishes, fresh + // summaries replace that cursor so a later write in an earlier key range is still found. + const behind = await seal({ + space: SPACE, + did: DID, + collection: COLLECTIONS.message, + rkey: 'message--behind-cursor', + record: messageRecord('landed behind the cursor'), + rev: REV[0], + deviceKeyId: 'human-node-1', + privateKey: await signer(), + }) + right.store.put(behind) + // A caller that has folded genesis may now supply summaries, but this peer still receives empty + // inventory until its retained cursor reaches the end. Otherwise the production roster would + // have to drop the route before a grant beyond the page boundary could arrive. + const second = await left.bus.catchUp(SPACE, { summaries: summarize(first), wants: [] }) + const third = await left.bus.catchUp(SPACE, { + summaries: summarize([...first, ...second]), + wants: [], + }) + assert.deepEqual(left.bus.bootstrapPeers(), []) + + // The next ordinary round starts from fresh summaries and finds the write that landed behind + // the old cursor while the bounded bootstrap scan was in progress. + const next = await left.bus.catchUp(SPACE, { + summaries: summarize([...first, ...second, ...third]), + wants: [], + }) + assert.equal(next.some((envelope) => envelopeKey(envelope) === envelopeKey(behind)), true) }) it('returns one copy of an envelope two peers both hold', async () => { diff --git a/packages/ingest/src/private.ts b/packages/ingest/src/private.ts index 4ac2445..6a2e306 100644 --- a/packages/ingest/src/private.ts +++ b/packages/ingest/src/private.ts @@ -1,5 +1,6 @@ import { parseAtUri } from '@radial/atproto' import { + COLLECTIONS, DEFAULT_ADMISSION_LIMITS, DEVICE_DIRECTORY_COLLECTIONS, EnvelopeWriter, @@ -175,6 +176,7 @@ const MAX_PEER_CANDIDATE_COOLDOWNS = 64 export class PrivateSpaceIngestor { readonly spaceDid: string + readonly #spaceRkey: string readonly knownDids = new Set() readonly quarantine: Quarantine readonly wants: WantList @@ -192,6 +194,8 @@ export class PrivateSpaceIngestor { readonly #onArrival: (() => void) | undefined readonly #onPeerCandidate: (() => void) | undefined readonly #candidateCooldown = new Map() + /** DIDs `knownDids` holds on a claim's word alone, pending the fold's answer — see `sync()`. */ + readonly #claimedDids = new Set() constructor( readonly spaceUri: string, @@ -200,7 +204,9 @@ export class PrivateSpaceIngestor { readonly records: RecordStore, options: PrivateSpaceIngestorOptions = {}, ) { - this.spaceDid = parseAtUri(spaceUri).did + const space = parseAtUri(spaceUri) + this.spaceDid = space.did + this.#spaceRkey = space.rkey this.knownDids.add(this.spaceDid) for (const did of options.bootstrapDids ?? []) this.knownDids.add(did) this.wants = options.wants ?? new MemoryWantList() @@ -367,10 +373,15 @@ export class PrivateSpaceIngestor { // Catch-up carries the want list explicitly. A summary can say what this replica HAS; only a // want can say what it knows it lost, and on a bus with no durable middle the difference is the - // difference between a re-fetch and a silent permanent fork (`quarantine.ts`). + // difference between a re-fetch and a silent permanent fork (`quarantine.ts`). A private + // envelope may arrive before genesis by gossip or bounded catch-up, though. On that partial + // footing this replica still dials so it can recover, but describes none of what it holds to a + // directory whose membership it cannot yet filter. The signed genesis projection is the point + // at which summaries, wants and blob requests become safe to disclose (ADR §30.3). + const hasGenesis = this.#hasGenesis() const offered = await this.bus.catchUp(this.spaceUri, { - summaries: summarize(this.envelopes.all()), - wants: this.wants.all(), + summaries: hasGenesis ? summarize(this.envelopes.all()) : [], + wants: hasGenesis ? this.wants.all() : [], }) for (const envelope of offered) await this.offer(envelope, report) @@ -399,12 +410,39 @@ export class PrivateSpaceIngestor { for (const member of index.members) { if (member.active) this.knownDids.add(member.did) } + // A claim-sourced DID was provisional, and this fold is what settles it (ADR §30.2). One that + // the fold now counts as an active member is on the ordinary footing and stays; one it does not + // stops being polled, so a stranger who talked a cold replica into one bounded read is not + // re-read every cycle for the rest of this replica's life. Their directory records stay in the + // store — records are never deleted — and stop MATTERING instead: `connectionDirectory` keeps a + // never-member out of every connection decision from the moment the space record folds. + for (const did of this.#claimedDids) { + if (!index.members.some((member) => member.did === did && member.active)) { + this.knownDids.delete(did) + } + } + this.#claimedDids.clear() return index } /** What the last `sync()` admitted, quarantined and refused. For diagnostics, not for trust. */ lastReport: AdmissionReport = emptyReport() + /** + * Whether the SIGNED genesis projection has landed — the space record carrying `deviceKeyIds`, + * which is the one that arrived in an envelope rather than off a public repo. + * + * This is the boundary every disclosure and the claim spend are gated on (ADR §30.1, §30.3), and + * it is a statement about records, never arrival time: gossip and bounded catch-up may land work + * envelopes before genesis, so "holds envelopes" and "can answer membership" are different facts, + * and everything gated here cares about the second. + */ + #hasGenesis(): boolean { + return ( + this.records.get(this.spaceDid, COLLECTIONS.space, this.#spaceRkey)?.deviceKeyIds !== undefined + ) + } + /** * The bytes every record here names and this replica does not hold. * @@ -426,6 +464,14 @@ export class PrivateSpaceIngestor { async syncBlobs(report: AdmissionReport = emptyReport()): Promise { const blobs = this.#blobs if (!blobs) return report + // A blob-request names a CID a private record carries, and until genesis folds this replica + // dials through an unfiltered directory (ADR §30.3) — so asking would hand whoever answered + // during bootstrap an identifier of private content. Same footing rule as summaries and wants: + // the store still says what is missing, and the first post-genesis cycle asks for it. + if (!this.#hasGenesis()) { + report.blobs.missing = this.missingBlobs().length + return report + } const needed = this.missingBlobs() // Round-robin rather than "the first N by CID". A blob nobody can serve — a record naming bytes // that were never written, or that every peer has lost — is retried forever by construction, and @@ -480,25 +526,59 @@ export class PrivateSpaceIngestor { } /** - * Whom a refused peer said it was, checked against the only thing that can settle it: its repo. + * Whom a refused peer said it was, checked against the only thing that can settle it: its repo — + * and only while this replica is still bootstrapping. * * A claim is not a credential and is never kept as one. What it buys is a bounded number of * *public, unauthenticated* repo reads — the same reads `pollDirectory` makes of a member — and a * DID survives them only if its own repo publishes a `deviceAddress` naming the endpoint the * transport authenticated, which is the question `authorizeEndpoint` already answers for every - * connection. So the worst a stranger can do by claiming somebody else's DID is make this replica - * read a public repo, and the most it can win by claiming its own is a directory entry it would - * have earned anyway by publishing an address it controls. + * connection. + * + * **A poll is a write, and the write outlives the state that allowed it.** It puts the polled + * repo's `device` and `deviceAddress` into the SAME store every connection decision is folded + * from, so a claim that is polled at all is a claim that lands in this replica's directory — and + * a directory entry, unlike the cold state that admitted it, is forever: records are never + * deleted. That is why what a stranger's entry is WORTH is settled at read time and not here — + * `connectionDirectory` (ADR §30.2) keeps a never-member out of every dial and every serve from + * the moment a space record folds, so the worst a claim can buy is a bounded poll into a replica + * that can serve and disclose nothing. + * + * The bootstrap condition below is still load-bearing on top of that, and the boundary is the + * signed GENESIS projection, not the first envelope (ADR §30.1). After genesis folds a claim has + * no use — `pollDirectory` already polls every DID in `knownDids`, and `knownDids` already holds + * every active member — so spending one there could only re-poll a member or write a stranger + * into the store for `connectionDirectory` to spend every subsequent frame filtering out. Before + * genesis folds a claim must still spend, however many envelopes have landed: admission + * deliberately does not check membership, so a gate-1-valid fragment from whoever caught the + * cold window is exactly what a partial replica may hold — and a boundary of "any envelope" + * would let that one fragment end discovery forever, claims drained and discarded while the one + * legitimate holder goes unfound. The spend stays safe on the partial footing because a partial + * replica serves nobody and describes nothing — empty summaries and wants, no blob requests + * (ADR §30.3). What is genuinely given up is narrow: a replica whose genesis has folded cannot + * learn a member it has not yet folded a grant for, so if the only reachable peer is that + * member, it waits for a peer it does know rather than trusting the claim. That is the same + * "wait rather than guess" the fold makes everywhere else (ADR §24). A DID kept here is + * provisional either way — `sync()` re-settles every claim-sourced DID against the fold the + * moment there is one. * - * Bounded three ways, because the claim arrives from anybody who can reach this endpoint: - * candidates per cycle, DIDs per candidate, and a per-endpoint cooldown so a peer redialling in a - * loop cannot buy a poll per dial. The cooldown table is bounded too — its keys are a stranger's - * to invent. + * Bounded three ways beyond that, because the claim arrives from anybody who can reach this + * endpoint: candidates per cycle, DIDs per candidate, and a per-endpoint cooldown so a peer + * redialling in a loop cannot buy a poll per dial. The cooldown table is bounded too — its keys + * are a stranger's to invent. */ async #pollPeerCandidates(signal?: AbortSignal): Promise { if (!this.#poller || !this.bus.drainPeerCandidates) return + // Drained either way: the buffer is a stranger's to fill, so it is emptied whether or not this + // replica is in a state to spend anything on it. + const candidates = this.bus.drainPeerCandidates().slice(0, MAX_PEER_CANDIDATES_PER_SYNC) + if (candidates.length === 0) return + // Genesis, not the first envelope — see the doc comment above: a partial replica serves and + // discloses nothing, so the spend is still free of leaks, and "any envelope" would let one + // stranger's fragment block discovery permanently. + if (this.#hasGenesis()) return const now = Date.now() - for (const candidate of this.bus.drainPeerCandidates().slice(0, MAX_PEER_CANDIDATES_PER_SYNC)) { + for (const candidate of candidates) { const last = this.#candidateCooldown.get(candidate.endpointId) ?? 0 if (now - last < PEER_CANDIDATE_COOLDOWN_MS) continue this.#rememberCandidate(candidate.endpointId, now) @@ -516,9 +596,12 @@ export class PrivateSpaceIngestor { (peer) => peer.did === did && peer.endpointId === candidate.endpointId, ) if (!bound) continue - // Membership is what keeps it here afterwards, if it belongs here at all: this decides whose - // repo is read and never what counts, so a DID that turns out to be nobody costs one poll. + // Provisional, and marked as such: `sync()` keeps this DID only if the fold counts it as an + // active member once there is a fold to ask. This decides whose repo is read and never what + // counts, so a DID that turns out to be nobody costs one poll — one poll into an empty + // replica, which is the only state this runs in. this.knownDids.add(did) + this.#claimedDids.add(did) verified = true } // The peer is dialable from this side now, and the backoff rung it is on was earned by diff --git a/packages/ingest/test/private.test.mjs b/packages/ingest/test/private.test.mjs index 0b1af9e..43e1c74 100644 --- a/packages/ingest/test/private.test.mjs +++ b/packages/ingest/test/private.test.mjs @@ -13,7 +13,9 @@ import { WirePrivateBus, authorizeEndpoint, connectablePeers, + connectionDirectory, createTidSource, + envelopeKey, exportPublicKey, generateDeviceKey, indexDigest, @@ -24,6 +26,7 @@ import { recordCid, seal, spaceTopic, + wantFor, } from '../../core/dist/index.js' import { MemorySyncStateStore, @@ -204,6 +207,62 @@ const shuffle = (items, seed) => { } 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 = [ @@ -648,6 +707,7 @@ function wireReplica(id, spaceUri, directory, options = {}) { 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, @@ -738,6 +798,43 @@ describe('private ingestion over the wire', () => { 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')] @@ -983,10 +1080,22 @@ describe('private ingestion over the wire', () => { describe('a claim from a peer this replica refused', () => { const AGENT = 'did:plc:privateagent' const IMPOSTOR = 'did:plc:privateimpostor' + const STRANGER = 'did:plc:totalstranger' const AGENT_ENDPOINT = 'endpoint-agent' + const STRANGER_ENDPOINT = 'endpoint-stranger' - /** A replica whose directory names nobody, and a fake set of public repos behind its poller. */ - async function refusingReplica(repos = new Map()) { + /** + * A replica whose directory names nobody, and a fake set of public repos behind its poller. + * + * Cold by default — no envelopes — which is one of the two states a claim is spent in; the + * boundary is the signed genesis projection, not the first envelope (ADR §30.1). A poll writes + * the polled repo's directory records into the store the roster and the authorization are both + * folded from, so a replica that spent a claim after genesis folded would be one a stranger + * could talk itself into being dialled and served by. `corpus: true` is that warm replica, which + * must refuse (ADR §30); `partial: true` holds a work envelope but no genesis, and must still + * spend — a stranger's fragment must not be able to end discovery. + */ + async function refusingReplica(repos = new Map(), options = {}) { const { spaceUri, rootDevice, memberDevice, envelopes } = await privateCorpus() const records = new MemoryRecordStore() for (const record of [ @@ -1001,9 +1110,14 @@ describe('a claim from a peer this replica refused', () => { spaceUri, envelopes: envelopeStore, peers: { async links() { return [] } }, - // No address in this directory names anybody: the honest state of a replica whose peers it - // has never read. Every inbound frame is refused, which is where a claim is captured. - authorizePeer: () => ({ ok: false, status: 'unknown', reason: 'no published device address' }), + // Either the hardcoded refusal — the honest state of a replica whose peers it has never read, + // and where a claim is captured — or the production wiring (`endpoint.ts`), which is what a + // test asking "would this stranger be served?" has to ask: the CONNECTION directory, filtered + // to ever-members from the moment a space record folds (ADR §30.2). + authorizePeer: options.liveDirectory + ? (endpointId) => + authorizeEndpoint(endpointId, connectionDirectory(records.records(), spaceUri, [])) + : () => ({ ok: false, status: 'unknown', reason: 'no published device address' }), }) const polled = [] const poller = { @@ -1023,14 +1137,35 @@ describe('a claim from a peer this replica refused', () => { woken += 1 }, }) - // Folded already, so `sync()` is exercising the candidate poll and not a cold start. - for (const envelope of envelopes) await ingestor.offer(envelope) + if (options.corpus) for (const envelope of envelopes) await ingestor.offer(envelope) + if (options.partial) { + const fragment = envelopes.find((envelope) => envelope.collection !== COLLECTIONS.space) + assert.ok(fragment) + assert.equal((await ingestor.offer(fragment)).admitted, 1) + } const topic = await spaceTopic(spaceUri) return { ingestor, records, + spaceUri, polled, woken: () => woken, + /** The corpus arriving over catch-up, after construction — how a cold replica warms up. */ + offerCorpus: async () => { + for (const envelope of envelopes) await ingestor.offer(envelope) + }, + /** + * One cycle. A replica holding no space record cannot fold one, and `sync()` says so by + * throwing — the candidate poll has already run by then, which is what makes ADR §24's + * `waiting` a state the loop passes through rather than a failure. + */ + sync: async () => { + try { + await ingestor.sync() + } catch (error) { + if (!String(error).includes('Space record not found')) throw error + } + }, /** One refused catch-up carrying a claim — a dialer's first frame, exactly. */ claim: (endpointId, from) => bus.accept(endpointId, { @@ -1058,7 +1193,7 @@ describe('a claim from a peer this replica refused', () => { assert.equal(answer.kind, 'error') assert.equal(answer.status, 'unavailable') - await replica.ingestor.sync() + await replica.sync() assert.deepEqual(replica.polled.filter((did) => did === AGENT), [AGENT]) assert.equal(replica.ingestor.knownDids.has(AGENT), true) // ...and connection management is told, because the peer's backoff was earned by refusals from @@ -1078,7 +1213,7 @@ describe('a claim from a peer this replica refused', () => { for (const record of impostor) replica.records.put(record) await replica.claim(AGENT_ENDPOINT, [AGENT]) - await replica.ingestor.sync() + await replica.sync() assert.deepEqual(replica.polled.filter((did) => did === AGENT), [AGENT], 'one poll, no more') assert.equal(replica.ingestor.knownDids.has(AGENT), false) assert.equal(replica.woken(), 0) @@ -1091,7 +1226,7 @@ describe('a claim from a peer this replica refused', () => { const replica = await refusingReplica() for (let dial = 0; dial < 3; dial += 1) { await replica.claim(AGENT_ENDPOINT, [AGENT]) - await replica.ingestor.sync() + await replica.sync() } assert.deepEqual(replica.polled.filter((did) => did === AGENT), [AGENT]) assert.equal(replica.ingestor.knownDids.has(AGENT), false) @@ -1102,7 +1237,7 @@ describe('a claim from a peer this replica refused', () => { for (let index = 0; index < 6; index += 1) { await replica.claim(`endpoint-stranger-${index}`, [`did:plc:stranger${index}`]) } - await replica.ingestor.sync() + await replica.sync() assert.equal( replica.polled.filter((did) => did.startsWith('did:plc:stranger')).length, 4, @@ -1111,4 +1246,129 @@ describe('a claim from a peer this replica refused', () => { // Nothing was kept: none of those repos says anything at all. assert.equal(replica.woken(), 0) }) + + // The regression the bootstrap condition exists for. + // + // A claim is settled by reading a repo, and reading a repo WRITES its `device` and `deviceAddress` + // into the same store the roster and `authorizeEndpoint` are folded from. So a replica that spent + // a claim while holding a corpus would hand a stranger the one thing connection authorization is + // there to withhold: a private corpus must not be enumerable by whoever computed the topic, and a + // topic is derived from a public URI. Everything the stranger needs here is public — they publish + // their own device and address to their own PDS, and claim nothing but their own DID. + it('will not let a stranger into the directory of a replica that holds a corpus', async () => { + const strangerDevice = await device(STRANGER, 'stranger-node') + const replica = await refusingReplica( + new Map([[ + STRANGER, + [ + directoryRecord(strangerDevice, '2026-02-01T00:00:03Z'), + addressRecord(strangerDevice, STRANGER_ENDPOINT), + ], + ]]), + // Warm, and authorized the way production authorizes (`endpoint.ts`). + { corpus: true, liveDirectory: true }, + ) + + // The premise: the stranger is nobody here. Membership is the fold's answer, and it is no. + const index = materialize(replica.records, { spaceUri: `at://${ROOT}/${COLLECTIONS.space}/space` }) + assert.equal(index.members.some((member) => member.did === STRANGER), false) + + const first = await replica.claim(STRANGER_ENDPOINT, [STRANGER]) + assert.equal(first.kind, 'error', 'a stranger is refused on the first dial') + + await replica.sync() + + // Not polled at all: the repo read is the step that would put them in the directory, so the + // refusal has to happen before it and not after. + assert.deepEqual(replica.polled.filter((did) => did === STRANGER), []) + assert.equal(replica.ingestor.knownDids.has(STRANGER), false) + assert.equal(replica.woken(), 0) + assert.equal( + authorizeEndpoint(STRANGER_ENDPOINT, readDeviceDirectory(replica.records.records(), [])).ok, + false, + 'a stranger talked their way into the directory', + ) + + // And the frame that matters: dialling again after a whole cycle still yields no corpus. + const second = await replica.claim(STRANGER_ENDPOINT, [STRANGER]) + assert.equal(second.kind, 'error', 'a stranger was served the private corpus') + assert.equal(second.status, 'unavailable') + }) + + // The boundary is genesis, not the first envelope (ADR §30.1). + // + // Admission deliberately does not check membership, so a gate-1-valid fragment from whoever + // caught the cold window is exactly the kind of thing a partial replica may hold — and if holding + // ANY envelope ended the bootstrap, that one fragment would end discovery forever: every later + // claim drained and discarded, the one legitimate holder never found, the replica stuck the + // moment its known peers go offline. Spending here stays leak-free because a partial replica + // serves nobody and describes nothing — empty summaries and wants, no blob requests (ADR §30.3). + it('still spends a claim while envelopes have arrived but genesis has not', async () => { + const replica = await refusingReplica( + new Map([[AGENT, await agentRepo(AGENT, AGENT_ENDPOINT)]]), + { partial: true, liveDirectory: true }, + ) + // The refusal itself is §30.3's: a partial replica authorizes no inbound peer. + const answer = await replica.claim(AGENT_ENDPOINT, [AGENT]) + assert.equal(answer.kind, 'error') + + await replica.sync() + assert.deepEqual(replica.polled.filter((did) => did === AGENT), [AGENT]) + assert.equal(replica.ingestor.knownDids.has(AGENT), true) + assert.equal(replica.woken(), 1) + }) + + // The residual the bootstrap condition alone leaves open (ADR §30.2). + // + // The claim above is spent only while this replica is cold — but the poll it buys writes the + // stranger's directory records into a store that never deletes anything, and the replica does not + // stay cold. Whether that entry outlives the state that admitted it is the whole question, and it + // is answered at read time: `connectionDirectory` keeps a never-member out of every dial and every + // serve from the moment the space record folds, and the fold evicts the claim-sourced DID from the + // poll list. A stranger who catches the bootstrap window gets one poll and an empty catch-up, and + // nothing after. + it('strips a stranger who claimed during bootstrap the moment the corpus folds', async () => { + const strangerDevice = await device(STRANGER, 'stranger-node') + const replica = await refusingReplica( + new Map([[ + STRANGER, + [ + directoryRecord(strangerDevice, '2026-02-01T00:00:03Z'), + addressRecord(strangerDevice, STRANGER_ENDPOINT), + ], + ]]), + // Cold, and authorized the way production authorizes (`endpoint.ts`). + { liveDirectory: true }, + ) + + // The §30 bootstrap, working as designed: the claim is spent, the stranger's repo genuinely + // binds the endpoint, and a replica with nothing to lose keeps the DID provisionally. + await replica.claim(STRANGER_ENDPOINT, [STRANGER]) + await replica.sync() + assert.deepEqual(replica.polled.filter((did) => did === STRANGER), [STRANGER]) + assert.equal(replica.ingestor.knownDids.has(STRANGER), true) + + // The corpus arrives. No sync has run yet — the filter is evaluated per frame against the live + // store, so there must be no window between "holds a corpus" and "refuses the stranger". + await replica.offerCorpus() + assert.equal( + authorizeEndpoint( + STRANGER_ENDPOINT, + connectionDirectory(replica.records.records(), replica.spaceUri, []), + ).ok, + false, + 'a stranger admitted while cold was still authorized against the corpus', + ) + const refused = await replica.claim(STRANGER_ENDPOINT, [STRANGER]) + assert.equal(refused.kind, 'error', 'a stranger was served the private corpus') + assert.equal(refused.status, 'unavailable') + + // The fold settles the provisional DID: evicted, so the poller stops reading their repo. This + // cycle still polls them once — `pollDirectory` runs before the fold — and none after it does. + await replica.sync() + assert.equal(replica.ingestor.knownDids.has(STRANGER), false) + const polls = replica.polled.filter((did) => did === STRANGER).length + await replica.sync() + assert.equal(replica.polled.filter((did) => did === STRANGER).length, polls) + }) }) diff --git a/packages/ui/src/lib/private-space.test.ts b/packages/ui/src/lib/private-space.test.ts index 9f2b06a..a39474b 100644 --- a/packages/ui/src/lib/private-space.test.ts +++ b/packages/ui/src/lib/private-space.test.ts @@ -159,6 +159,38 @@ async function spaceEnvelope( }) } +/** + * A membership grant sealed by the founder — the other half of a corpus a peer would serve. + * + * A peer neither dials nor serves a device whose owner has never been a member of the space + * (`connectionDirectory`, ADR §30.2), so a corpus that is only the space record is a space nobody + * but the founder could ever sync: a realistic serving peer holds the grant of everyone it serves. + */ +async function grantEnvelope( + keys: BrowserDeviceKeyStore, + key: DeviceKeyRecord, + space: PrivateEnvelope, + member: { did: string; kind: 'human' | 'agent'; rev: string }, +): Promise { + return seal({ + space: SPACE, + did: FOUNDER, + collection: COLLECTIONS.addMember, + rkey: `grant${member.did.replace(/[^a-z0-9]/gi, '')}`, + record: { + $type: COLLECTIONS.addMember, + space: { uri: SPACE, cid: space.recordCid }, + did: member.did, + kind: member.kind, + role: member.kind === 'agent' ? 'agent' : 'member', + createdAt: NOW, + }, + rev: member.rev, + deviceKeyId: key.deviceKeyId, + privateKey: (await keys.signer(key)).privateKey, + }) +} + /** The fold on screen — only ever asked for under a `ready` status, as a route would. */ function index(): MaterializedIndex { const current = session.space @@ -1000,7 +1032,7 @@ describe('a tab with a transport endpoint', () => { * fake PDS the tab publishes into, which is what makes "the founder will serve this tab" a claim * about a `deviceAddress` the tab really wrote rather than about a list handed in by the test. */ - async function peerHolding(h: Harness, envelope: PrivateEnvelope) { + async function peerHolding(h: Harness, corpus: PrivateEnvelope[]) { const records = new MemoryRecordStore() const envelopes = new MemoryEnvelopeStore(records) let bound: { endpoint: FrameEndpoint; onLink?: (link: PeerLink) => void } | undefined @@ -1020,7 +1052,7 @@ describe('a tab with a transport endpoint', () => { }) await ingestor.start() await ingestor.pollDirectory() - expect((await ingestor.offer(envelope)).admitted).toBe(1) + for (const envelope of corpus) expect((await ingestor.offer(envelope)).admitted).toBe(1) return { endpoint, ingestor, @@ -1037,7 +1069,7 @@ describe('a tab with a transport endpoint', () => { /** The same peer, but able to dial the tab — for the direction gossip travels. */ async function peerHoldingThatDials( h: Harness, - envelope: PrivateEnvelope, + corpus: PrivateEnvelope[], tabFrames: () => FrameEndpoint | undefined, ) { const records = new MemoryRecordStore() @@ -1052,7 +1084,7 @@ describe('a tab with a transport endpoint', () => { }) await ingestor.start() await ingestor.pollDirectory() - expect((await ingestor.offer(envelope)).admitted).toBe(1) + for (const envelope of corpus) expect((await ingestor.offer(envelope)).admitted).toBe(1) return { poll: () => ingestor.pollDirectory(), gossip: (offered: PrivateEnvelope) => attached.bus.publish(SPACE, offered), @@ -1094,7 +1126,12 @@ describe('a tab with a transport endpoint', () => { const founder = await founderDevice(h, h.keys) await publishFounderAddress(h, founder.key) const envelope = await spaceEnvelope(h.keys, founder.key) - const peer = await peerHolding(h, envelope) + const grant = await grantEnvelope(h.keys, founder.key, envelope, { + did: MEMBER, + kind: 'human', + rev: '3kaaaaaaaaaa3', + }) + const peer = await peerHolding(h, [envelope, grant]) // The tab: signed in as the member, holding a bootstrap and nothing else. signInAs(h, MEMBER) @@ -1140,7 +1177,7 @@ describe('a tab with a transport endpoint', () => { // And now the thing a tab could not do. const index = await replica.sync() expect(index.space.value.name).toBe('Skunkworks') - expect(replica.stores.envelopes.all()).toHaveLength(1) + expect(replica.stores.envelopes.all()).toHaveLength(2) expect(replica.connections?.connected).toBe(1) } finally { replica.close() @@ -1158,7 +1195,7 @@ describe('a tab with a transport endpoint', () => { */ async function daemonPeer( h: Harness, - envelope: PrivateEnvelope, + corpus: PrivateEnvelope[], tab: () => { endpoint: FrameEndpoint; onLink?: (link: PeerLink) => void } | undefined, ) { const keys = await BrowserDeviceKeyStore.open(memoryDatabase()) @@ -1196,7 +1233,7 @@ describe('a tab with a transport endpoint', () => { }) await ingestor.start() await ingestor.pollDirectory() - expect((await ingestor.offer(envelope)).admitted).toBe(1) + for (const envelope of corpus) expect((await ingestor.offer(envelope)).admitted).toBe(1) return { endpoint, frames: () => bound?.endpoint, @@ -1210,6 +1247,21 @@ describe('a tab with a transport endpoint', () => { const h = await harness() const founder = await founderDevice(h, h.keys) const envelope = await spaceEnvelope(h.keys, founder.key) + const grants = [ + await grantEnvelope(h.keys, founder.key, envelope, { + did: MEMBER, + kind: 'human', + rev: '3kaaaaaaaaaa3', + }), + // The daemon is a member too — an agent's membership is a grant like anybody's (§11), and a + // peer that were not one would be evicted from this tab's directory the moment the corpus + // folded (`connectionDirectory`, ADR §30.2). + await grantEnvelope(h.keys, founder.key, envelope, { + did: AGENT, + kind: 'agent', + rev: '3kaaaaaaaaaa4', + }), + ] // The reported setup: the corpus is on an agent daemon, the founder's own machine is off, and // this browser has never held the space. Its bootstrap is the founder and itself — the ticket, // the bookmark and the hints can carry nothing else, because none of them existed when the @@ -1217,7 +1269,7 @@ describe('a tab with a transport endpoint', () => { signInAs(h, MEMBER) await memberDevice(h) let tabBound: { endpoint: FrameEndpoint; onLink?: (link: PeerLink) => void } | undefined - const daemon = await daemonPeer(h, envelope, () => tabBound) + const daemon = await daemonPeer(h, [envelope, ...grants], () => tabBound) const endpoint = await browserEndpoint({ transport: loopbackTransport({ endpointId: TAB_ENDPOINT, @@ -1277,7 +1329,14 @@ describe('a tab with a transport endpoint', () => { const founder = await founderDevice(h, h.keys) await publishFounderAddress(h, founder.key) const envelope = await spaceEnvelope(h.keys, founder.key) - const peer = await peerHolding(h, envelope) + const grant = await grantEnvelope(h.keys, founder.key, envelope, { + did: MEMBER, + kind: 'human', + rev: '3kaaaaaaaaaa3', + }) + // A member's grant is in the corpus and it changes nothing: authorization is about the + // ENDPOINT, and this tab never says where it is. + const peer = await peerHolding(h, [envelope, grant]) signInAs(h, MEMBER) await memberDevice(h) @@ -1321,7 +1380,19 @@ describe('a tab with a transport endpoint', () => { signInAs(h, MEMBER) await memberDevice(h) let tabFrames: FrameEndpoint | undefined - const peer = await peerHoldingThatDials(h, envelope, () => tabFrames) + // The grant is what makes this tab a device the peer will dial at all (ADR §30.2). + const peer = await peerHoldingThatDials( + h, + [ + envelope, + await grantEnvelope(h.keys, founder.key, envelope, { + did: MEMBER, + kind: 'human', + rev: '3kaaaaaaaaaa3', + }), + ], + () => tabFrames, + ) const endpoint = await browserEndpoint({ transport: loopbackTransport({ endpointId: TAB_ENDPOINT, @@ -1384,7 +1455,12 @@ describe('a tab with a transport endpoint', () => { const founder = await founderDevice(h, h.keys) await publishFounderAddress(h, founder.key) const envelope = await spaceEnvelope(h.keys, founder.key) - const peer = await peerHolding(h, envelope) + const grant = await grantEnvelope(h.keys, founder.key, envelope, { + did: MEMBER, + kind: 'human', + rev: '3kaaaaaaaaaa3', + }) + const peer = await peerHolding(h, [envelope, grant]) signInAs(h, MEMBER) const mine = await memberDevice(h) // This tab is already in the peer's directory, so the catch-up below is served. Getting there is