// Whom this replica dials, and how long it keeps the connection (design §18, plan phase 3). // // `peers.ts` answers "may this endpoint be served?" and "which endpoints does the directory name?" as // pure functions of a record set. `sync.ts` answers "what do two endpoints say to each other". What // sat between them was a gap the size of a transport: `fixedPeers()` — a list handed in once, never // re-read, never redialled, never dropped — which is exactly right for a test and is not a policy any // replica could run on. // // Filling that gap inside a transport was the tempting shape and the wrong one, for the same reason // §14 moved the wire out of one: none of the questions is an iroh question. // // - **Which endpoints, in what order.** The directory is the answer, and it changes under the // connection: a newly published address would otherwise remain unseen. // - **How often to retry one that will not answer.** A private space is a handful of devices, some // of them laptops that are shut most of the day. Dialling them on every publish is a busy loop // with a network in it. // - **What a link that throws means.** The bus counts a failure and hands the same dead link out // again on the next call, forever, because nothing tells the source that a connection died. // // Answering them here means the daemon and a browser tab hold the same connections for the same // reasons, and a binding implements one function: `dial`. // // Two properties are worth stating because they are the ones that would be easy to lose: // // - **The roster is re-read every round, never cached.** That is what makes retiring a device sever // a connection rather than merely refuse the next one — a device that has left the directory has its // link closed at the top of the next `links()`, before anything is sent to it. It is a // *connection* statement, in §18's vocabulary, and it changes nothing about the index: everything // that peer already sent stays in the store and is judged by the fold exactly as before. // - **This is not a gate, and it is not authorization.** It decides who is dialled; `accept()` // decides who is served, and gates 1 and 2 decide what counts. A peer this file dials is trusted // for nothing — its envelopes go through `admit()` like a stranger's, which is why a stale roster // costs a reconnect and never a wrong record. import type { PeerIdentity } from './peers.js' import type { PeerLink, PeerSource } from './sync.js' /** * What a transport supplies: one authenticated connection to one addressed device. * * Everything else in private mode is above this line. The returned `PeerLink.id` is the endpoint * identity the transport authenticated and must equal the requested `peer.endpointId`. A stale or * misdirected address is not permission to connect to a different endpoint. */ export type PeerDialer = (peer: PeerIdentity) => Promise /** * How long to wait after the nth consecutive failure to reach one endpoint. * * Deterministic, with no jitter, which is a deliberate trade rather than an oversight. Jitter exists * to break up a herd, and the herd here is a private space's device list — single digits, dialling a * peer each, on a schedule that is already staggered by whenever each replica last synced. What * determinism buys instead is that "why is this peer not connected?" has the same answer on every * machine and in every test, which is the question an operator actually asks. */ export const DEFAULT_DIAL_BACKOFF_MS: readonly number[] = [1_000, 5_000, 15_000, 60_000, 300_000] /** Enough connections for a space's devices, and a bound so a large directory cannot become one. */ const DEFAULT_MAX_LINKS = 16 /** * How long one dial may hold up a round. * * A dial to an endpoint that is simply gone — a closed tab whose `deviceAddress` still names it, a * machine that moved — does not fail fast: it waits out a transport handshake that can run tens of * seconds. `links()` fronts every catch-up, and in a browser the first paint of a private space, so * one dead device must not add its handshake to every member's open. The bound decides only when the * round stops waiting: the dial keeps running, and a link that lands late is adopted exactly as an * accepted inbound one would be. */ export const DEFAULT_DIAL_TIMEOUT_MS = 5_000 /** The bound above, spent. Named so the late-landing continuation can tell it from a dial failure. */ class DialTimeout extends Error {} export interface PeerConnectionsOptions { spaceUri: string /** * The devices this replica may dial, as the directory currently describes them — normally * `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 * folds a store to produce it should memoise on whatever version the store already exposes. */ roster: () => readonly PeerIdentity[] dial: PeerDialer /** * The endpoint identities this replica answers to, so it never dials itself. * * A callback rather than a value because it changes under a running bus: rotating a device * publishes a new address, and the endpoint that appears in the directory a moment later is this * one. Dialling yourself is not fatal — `accept()` would authorize you and serve you your own * corpus — it is just a connection spent on learning nothing. */ self?: () => readonly string[] /** * DIDs to dial first: `join.peerHints`, the members expected to be reachable. * * Ordering only, never filtering. A hint is a guess about liveness and confers nothing, so a * replica with hints still dials everyone else the directory names once the hinted peers are * connected — otherwise a stale hint would be a partition. */ prefer?: readonly string[] /** Live connections this replica will hold at once. Extra roster entries wait for a slot. */ maxLinks?: number backoffMs?: readonly number[] /** How long one dial may gate a round before it is counted failed and left to land late. */ dialTimeoutMs?: number now?: () => number } /** One endpoint as this replica currently sees it. Diagnostics — never an input to anything. */ export interface PeerConnectionState { endpointId: string did: string deviceKeyId: string /** Whether a link to this endpoint is held right now. */ connected: boolean /** True when the link was handed in by the transport (`adopt`) rather than dialled from here. */ inbound: boolean /** Consecutive failures — dials that threw, plus links that died in use. Zero on success. */ failures: number /** The clock reading before which no further dial is attempted. */ retryAt?: number lastError?: string } const describe = (error: unknown): string => error instanceof Error ? error.message : String(error) /** * The connections a replica keeps to a private space's other devices. * * A `PeerSource`, so `WirePrivateBus` cannot tell it from `fixedPeers()` — which is the point: * everything above the bus, from `PrivateSpaceIngestor` to dispatch, stays unaware that connections * come and go at all. * * What it is NOT is a transport. It never touches a socket, a relay or a key; it turns the public * directory into an ordered dial list, keeps at most one link per endpoint, drops the ones the * directory has stopped naming, and backs off the ones that will not answer. `dial` supplies the * bytes, and the same instance runs over an iroh endpoint, a WASM endpoint against a public relay, * and `loopbackLink` in a test. */ export class PeerConnections implements PeerSource { readonly spaceUri: string readonly #roster: () => readonly PeerIdentity[] readonly #dial: PeerDialer readonly #self: () => readonly string[] readonly #prefer: ReadonlySet readonly #maxLinks: number readonly #backoff: readonly number[] readonly #dialTimeoutMs: number readonly #now: () => number readonly #links = new Map() readonly #states = new Map() readonly #dialing = new Map>() #closed = false constructor(options: PeerConnectionsOptions) { this.spaceUri = options.spaceUri this.#roster = options.roster this.#dial = options.dial this.#self = options.self ?? (() => []) this.#prefer = new Set(options.prefer ?? []) this.#maxLinks = Math.max(1, options.maxLinks ?? DEFAULT_MAX_LINKS) const backoff = options.backoffMs ?? DEFAULT_DIAL_BACKOFF_MS this.#backoff = backoff.length > 0 ? backoff : DEFAULT_DIAL_BACKOFF_MS this.#dialTimeoutMs = options.dialTimeoutMs ?? DEFAULT_DIAL_TIMEOUT_MS this.#now = options.now ?? (() => Date.now()) } /** * The endpoints worth dialling, in the order they are dialled. * * Hinted DIDs first, then the directory's own `(did, deviceKeyId)` order, minus this replica's own * endpoints and minus duplicates — two device keys of one process legitimately publish one endpoint * id, and holding two links to it would double every gossip. */ #candidates(): PeerIdentity[] { const mine = new Set(this.#self()) const preferred: PeerIdentity[] = [] const rest: PeerIdentity[] = [] const seen = new Set() for (const peer of this.#roster()) { if (mine.has(peer.endpointId) || seen.has(peer.endpointId)) continue seen.add(peer.endpointId) ;(this.#prefer.has(peer.did) ? preferred : rest).push(peer) } return [...preferred, ...rest] } #stateFor(peer: PeerIdentity): PeerConnectionState { const held = this.#states.get(peer.endpointId) if (held) return held const fresh: PeerConnectionState = { endpointId: peer.endpointId, did: peer.did, deviceKeyId: peer.deviceKeyId, connected: false, inbound: false, failures: 0, } this.#states.set(peer.endpointId, fresh) return fresh } #backoffFor(failures: number): number { const index = Math.min(failures, this.#backoff.length) - 1 return this.#backoff[Math.max(0, index)] ?? 0 } #penalise(state: PeerConnectionState, error: unknown): void { state.failures += 1 state.retryAt = this.#now() + this.#backoffFor(state.failures) state.lastError = describe(error) state.connected = false state.inbound = false } #close(endpointId: string): void { const link = this.#links.get(endpointId) if (!link) return this.#links.delete(endpointId) const state = this.#states.get(endpointId) if (state) { state.connected = false state.inbound = false } try { link.close?.() } catch { // A transport that throws on close has already lost the connection this was closing; there is // nothing left to do about it and nothing above this line that would act on knowing. } } /** The dial, raced against the bound. The dial itself keeps running past an expiry. */ #within(dialing: Promise, peer: PeerIdentity): Promise { const bound = this.#dialTimeoutMs if (!Number.isFinite(bound)) return dialing let timer: ReturnType | undefined const expiry = new Promise((_, reject) => { timer = setTimeout( () => reject(new DialTimeout(`${peer.endpointId} did not answer a dial within ${bound}ms`)), bound, ) }) return Promise.race([dialing, expiry]).finally(() => clearTimeout(timer)) } /** * Keep a dialled link, if this replica still wants one. * * The directory can change while a dial is in flight, and a slow dial to a device address changed * meanwhile must not land as a held connection. A dropped state is exactly that event: a round * that ran during the dial found this endpoint gone and forgot it. * ...and an inbound connection from the same device can be adopted while this one is dialling, * in which case the accepted link is the one to keep: dropping it would close a connection the * peer believes it has. */ #land(peer: PeerIdentity, link: PeerLink, state: PeerConnectionState): void { if (link.id !== peer.endpointId) { link.close?.() this.#penalise( state, new Error(`dial authenticated ${link.id}, expected ${peer.endpointId}`), ) return } if (this.#closed || !this.#states.has(peer.endpointId) || this.#links.has(peer.endpointId)) { link.close?.() return } this.#links.set(peer.endpointId, link) state.connected = true state.inbound = false state.failures = 0 delete state.retryAt delete state.lastError } async #connect(peer: PeerIdentity): Promise { const inflight = this.#dialing.get(peer.endpointId) // Two overlapping rounds — a publish and a catch-up in the same tick — must not open two // connections to one endpoint and leave one of them unreferenced. if (inflight) return inflight const state = this.#stateFor(peer) const attempt = (async () => { const dialing = (async () => this.#dial(peer))() let link: PeerLink try { link = await this.#within(dialing, peer) } catch (error) { this.#penalise(state, error) if (error instanceof DialTimeout) { // The dial is still running; only the round has stopped waiting for it. A link that // arrives after that is adopted exactly as an accepted inbound one would be, and a dial // that fails after it is not a second strike — the timeout was this attempt's one. dialing.then( (landed) => this.#land(peer, landed, state), () => undefined, ) } return } this.#land(peer, link, state) })() this.#dialing.set(peer.endpointId, attempt) try { await attempt } finally { this.#dialing.delete(peer.endpointId) } } async links(spaceUri: string): Promise { if (spaceUri !== this.spaceUri) { throw new Error(`these connections serve ${this.spaceUri}, not ${spaceUri}`) } if (this.#closed) return [] const candidates = this.#candidates() const dialable = new Set(candidates.map((peer) => peer.endpointId)) // Revocation, removal and retirement all arrive here as an endpoint the directory has stopped // naming, and all three end the connection before anything else happens this round. for (const endpointId of [...this.#links.keys()]) { if (!dialable.has(endpointId)) this.#close(endpointId) } for (const endpointId of [...this.#states.keys()]) { // A device that left the directory is not a peer in backoff; it is not a peer. Keeping its // state would make a re-published address serve out a stale retry deadline. if (!dialable.has(endpointId)) this.#states.delete(endpointId) } const now = this.#now() const missing = candidates.filter((peer) => { if (this.#links.has(peer.endpointId)) return false const state = this.#states.get(peer.endpointId) return !state || state.retryAt === undefined || state.retryAt <= now }) const slots = Math.max(0, this.#maxLinks - this.#links.size) // Concurrently, because one unreachable peer's dial timeout must not be added to every other // peer's; over a fixed candidate list, so the outcome does not depend on which dial finished // first. await Promise.all(missing.slice(0, slots).map((peer) => this.#connect(peer))) const live: PeerLink[] = [] for (const peer of candidates) { const link = this.#links.get(peer.endpointId) if (link) live.push(link) } return live } /** * A link this bus held threw while in use. * * The connection is closed and the endpoint enters backoff, so the next round dials a fresh one * instead of handing the dead link out again. Only a *transport* failure lands here: a peer that * answers a catch-up with the wrong frame, or refuses one, has answered, and the connection is * fine — `WirePrivateBus` counts that against the peer and leaves the link alone. */ failed(endpointId: string, error: unknown): void { const state = this.#states.get(endpointId) this.#close(endpointId) if (state) this.#penalise(state, error) } /** * Ask every peer again on the next round, whatever its backoff says. * * The ladder answers "this peer is unreachable, stop hammering it", and it is the right answer to * that question. It is the wrong answer to the one this exists for: a peer that is up, and refused * this replica because its device directory is a poll behind the device this replica has just * published (`peers.ts` — that refusal is `unknown`, which means *ask again*). Those look identical * from here, because the peer closes the connection either way, so the first few dials of a newly * published device routinely spend strikes on a peer that was always going to say yes shortly. * * A caller that knows something has changed on the OTHER side — a person who has just come back to * the tab, a replica still waiting for its first peer — can therefore spend a round asking again. * The failure count is deliberately kept: this clears the wait, not the history, so a peer that is * genuinely gone is back on the same rung of the ladder after this one attempt rather than pinned * to the bottom of it. */ retryNow(): void { for (const state of this.#states.values()) delete state.retryAt } /** * Take over a connection the transport accepted, so gossip flows back over it. * * Without this, two replicas behind NATs where only one side can dial would sync in one direction: * the dialer's writes reach the accepter, and the accepter's writes wait for the dialer's next * catch-up. The peer must be one the transport authenticated (`authorizeEndpoint`) — this adopts a * link, it does not authorize one — and it is dropped like any other at the next round if the * directory has stopped naming it. */ adopt(peer: PeerIdentity, link: PeerLink): void { if (this.#closed) { link.close?.() return } const held = this.#links.get(peer.endpointId) if (held === link) return // Replacing the same endpoint consumes no new slot. A different endpoint arriving while the // bound is full is always refused: existing links never depend on inbound arrival order, and a // later round can dial this peer when a slot becomes available. if (!held && this.#links.size >= this.#maxLinks) { link.close?.() return } if (held) this.#close(peer.endpointId) this.#links.set(peer.endpointId, link) const state = this.#stateFor(peer) state.connected = true state.inbound = true state.failures = 0 delete state.retryAt delete state.lastError } /** Every endpoint this replica has an opinion about, sorted. For `radial private status`. */ state(): PeerConnectionState[] { return [...this.#states.values()] .map((state) => ({ ...state })) .sort((left, right) => left.endpointId < right.endpointId ? -1 : left.endpointId > right.endpointId ? 1 : 0, ) } /** How many links are held right now — the one number a status line usually wants. */ get connected(): number { return this.#links.size } close(): void { this.#closed = true for (const endpointId of [...this.#links.keys()]) this.#close(endpointId) this.#states.clear() } }