iroh_db/
network.rs

1use std::collections::BTreeMap;
2use std::sync::{
3    Arc, RwLock,
4    atomic::{AtomicU64, Ordering},
5};
6use std::time::Duration;
7
8use iroh::{
9    Endpoint, EndpointAddr, RelayMode,
10    address_lookup::memory::MemoryLookup,
11    endpoint::{Connection, RecvStream, SendStream, presets},
12    protocol::{AcceptError, ProtocolHandler, Router},
13};
14use iroh_blobs::provider::events::{
15    AbortReason, ConnectMode, EventMask, EventSender, ObserveMode, ProviderMessage, RequestMode,
16};
17use iroh_db_blobs::{BlobEngine, BlobRef, BlobStream, FetchMetrics, FetchOptions};
18use iroh_db_core::{AuthorId, BlobHash, CommitEnvelope, CommitId, DomainId};
19use iroh_db_security::{
20    CapabilityCertificate, ControlTransition, DomainSeed, Invitation, Permission,
21    verify_gossip_signature,
22};
23use iroh_db_sync::{
24    GossipHint, MAX_CONTROL_FRAME, MAX_PROTOCOL_VERSION, MAX_REFERENCES, MAX_SESSION_REFERENCES,
25    MIN_PROTOCOL_VERSION, SyncFrame, missing_commits, negotiate_protocol,
26};
27use iroh_gossip::{
28    api::{Event as GossipEvent, GossipReceiver, GossipSender},
29    net::Gossip,
30    proto::TopicId,
31};
32use n0_future::StreamExt as _;
33
34use crate::{DbError, IrohDb, Snapshot};
35
36/// The sole iroh-db control-plane ALPN.
37pub const SYNC_ALPN: &[u8] = b"/iroh-db/sync/4";
38
39/// Result counters for one complete bidirectional reconciliation round.
40#[derive(Debug, Clone, Copy, PartialEq, Eq)]
41pub struct SyncReport {
42    received: usize,
43    sent: usize,
44}
45
46impl SyncReport {
47    /// Number of remote commits fetched and accepted locally.
48    pub const fn received(self) -> usize {
49        self.received
50    }
51
52    /// Number of local commits pushed and accepted remotely.
53    pub const fn sent(self) -> usize {
54        self.sent
55    }
56}
57
58/// A running iroh endpoint routing sync control and immutable blob protocols.
59pub struct SyncNode {
60    core: SyncCore,
61    gossip: Gossip,
62    lookup: MemoryLookup,
63    router: Router,
64}
65
66impl SyncNode {
67    /// Returns this endpoint's currently advertised direct/relay address.
68    pub fn addr(&self) -> EndpointAddr {
69        self.router.endpoint().addr()
70    }
71
72    /// Reconciles both peers over authenticated QUIC and iroh-blobs.
73    pub async fn sync_with(&self, remote: EndpointAddr) -> Result<SyncReport, DbError> {
74        self.core.sync_with(remote).await
75    }
76
77    /// Makes one signed invitation available to its endpoint-bound subject.
78    pub fn offer_invitation(&self, invitation: &Invitation) -> Result<(), DbError> {
79        let _control_reader = self
80            .core
81            .db
82            .inner
83            .control_gate
84            .try_read()
85            .map_err(|_| DbError::Operation("control plane is busy".into()))?;
86        self.core.require_local_permission(Permission::Invite)?;
87        if invitation.descriptor().domain_id() != self.core.domain_id() {
88            return Err(protocol_error("invitation targets another domain"));
89        }
90        if invitation.domain_seed().epoch() != self.core.db.epoch() {
91            return Err(DbError::UnauthorizedDevice);
92        }
93        let payload = invitation.encode_canonical()?;
94        self.core
95            .pending_invitations
96            .write()
97            .map_err(|_| protocol_error("invitation registry is unavailable"))?
98            .insert(invitation.subject(), payload);
99        Ok(())
100    }
101
102    /// Requests this authenticated endpoint's offered invitation and imports it locally.
103    pub async fn request_invitation(&self, remote: EndpointAddr) -> Result<IrohDb, DbError> {
104        let expected_issuer = author_id(remote.id);
105        let connection = self
106            .core
107            .endpoint
108            .connect(remote, SYNC_ALPN)
109            .await
110            .map_err(network_error)?;
111        let (mut send, mut recv) = connection.open_bi().await.map_err(network_error)?;
112        send_frame(&mut send, &hello_frame(random_nonce()?)).await?;
113        let SyncFrame::Hello {
114            min_protocol,
115            max_protocol,
116            ..
117        } = recv_frame(&mut recv).await?
118        else {
119            return Err(protocol_error("expected hello"));
120        };
121        negotiate_protocol(
122            MIN_PROTOCOL_VERSION,
123            MAX_PROTOCOL_VERSION,
124            min_protocol,
125            max_protocol,
126        )
127        .map_err(network_error)?;
128        send_frame(&mut send, &SyncFrame::InvitationRequest).await?;
129        let SyncFrame::Invitation { payload } = recv_frame(&mut recv).await? else {
130            return Err(DbError::UnauthorizedDevice);
131        };
132        let invitation = Invitation::decode_canonical(&payload)?;
133        if invitation.issuer() != expected_issuer {
134            return Err(DbError::UnauthorizedDevice);
135        }
136        let domain = self.core.db.import_invitation(&invitation)?;
137        send.finish().map_err(network_error)?;
138        Ok(domain)
139    }
140
141    /// Fetches and installs a trusted encrypted snapshot before ordinary tail reconciliation.
142    pub async fn bootstrap_from(
143        &self,
144        remote: EndpointAddr,
145        snapshot: Snapshot,
146    ) -> Result<(), DbError> {
147        self.core.require_local_blob_fetch()?;
148        self.core.require_trusted(author_id(remote.id))?;
149        let connection = self
150            .core
151            .endpoint
152            .connect(remote, iroh_blobs::ALPN)
153            .await
154            .map_err(network_error)?;
155        let engine = self.core.db.blob_engine();
156        engine
157            .fetch_from(&self.core.db.inner.blobs, connection, &snapshot.blob_ref())
158            .await?;
159        self.core.db.install_snapshot(snapshot).await
160    }
161
162    /// Downloads a complete encrypted blob from one authorized provider into pinned storage.
163    pub async fn fetch_blob_from(
164        &self,
165        remote: EndpointAddr,
166        reference: &BlobRef,
167    ) -> Result<(), DbError> {
168        self.core.require_local_blob_fetch()?;
169        self.core.require_trusted(author_id(remote.id))?;
170        let connection = self
171            .core
172            .endpoint
173            .connect(remote, iroh_blobs::ALPN)
174            .await
175            .map_err(network_error)?;
176        self.core
177            .blob_engine()
178            .fetch_from(&self.core.db.inner.blobs, connection, reference)
179            .await
180            .map_err(DbError::from)
181    }
182
183    /// Opens a lazy seekable stream whose missing chunks rotate across authorized providers.
184    pub async fn stream_blob_from(
185        &self,
186        providers: Vec<EndpointAddr>,
187        reference: &BlobRef,
188        options: FetchOptions,
189    ) -> Result<(BlobStream, FetchMetrics), DbError> {
190        self.core.require_local_blob_fetch()?;
191        let mut connections = Vec::with_capacity(providers.len());
192        for provider in providers {
193            self.core.require_trusted(author_id(provider.id))?;
194            connections.push(
195                self.core
196                    .endpoint
197                    .connect(provider, iroh_blobs::ALPN)
198                    .await
199                    .map_err(network_error)?,
200            );
201        }
202        self.core
203            .blob_engine()
204            .open_from_providers(
205                self.core.db.inner.blobs.as_ref().clone(),
206                connections,
207                reference,
208                options,
209            )
210            .await
211            .map_err(DbError::from)
212    }
213
214    /// Joins this domain's unguessable gossip topic using optional bootstrap endpoints.
215    pub async fn join_gossip(
216        &self,
217        bootstrap: Vec<EndpointAddr>,
218    ) -> Result<GossipSession, DbError> {
219        self.core.require_local_permission(Permission::Read)?;
220        for peer in &bootstrap {
221            self.lookup.add_endpoint_info(peer.clone());
222        }
223        let peers: Vec<_> = bootstrap.into_iter().map(|peer| peer.id).collect();
224        let topic = self.core.gossip_topic()?;
225        let topic = if peers.is_empty() {
226            self.gossip.subscribe(topic, peers).await
227        } else {
228            self.gossip.subscribe_and_join(topic, peers).await
229        }
230        .map_err(network_error)?;
231        let (sender, receiver) = topic.split();
232        Ok(GossipSession {
233            core: self.core.clone(),
234            sender,
235            receiver,
236            counter: Arc::new(AtomicU64::new(initial_gossip_counter())),
237            seen: BTreeMap::new(),
238        })
239    }
240
241    /// Stops network protocols without closing the embedded database.
242    pub async fn shutdown(self) -> Result<(), DbError> {
243        self.router
244            .shutdown()
245            .await
246            .map_err(|error| DbError::Network(error.to_string()))
247    }
248}
249
250/// An active, lossy gossip subscription that can announce and act on sync hints.
251pub struct GossipSession {
252    core: SyncCore,
253    sender: GossipSender,
254    receiver: GossipReceiver,
255    counter: Arc<AtomicU64>,
256    seen: BTreeMap<AuthorId, u64>,
257}
258
259impl GossipSession {
260    /// Broadcasts a signed hash-only hint for the current commit frontier.
261    pub async fn announce(&self) -> Result<GossipHint, DbError> {
262        self.core.require_local_permission(Permission::Read)?;
263        let commits = self.core.commit_ids()?;
264        let commit = commits
265            .last()
266            .copied()
267            .ok_or_else(|| protocol_error("cannot announce an empty commit set"))?;
268        let frontier_digest = digest_commits(&commits);
269        let counter = self.counter.fetch_add(1, Ordering::Relaxed);
270        let signer = self.core.db.inner.credentials.signer();
271        let unsigned = GossipHint::new(signer.author(), commit, frontier_digest, counter, [0; 64]);
272        let signature = signer.sign_gossip(&unsigned.signing_bytes().map_err(sync_error)?);
273        let hint = GossipHint::new(signer.author(), commit, frontier_digest, counter, signature);
274        self.sender
275            .broadcast(hint.encode_canonical().map_err(sync_error)?.into())
276            .await
277            .map_err(network_error)?;
278        Ok(hint)
279    }
280
281    /// Waits for a valid hint and runs ordinary authenticated reconciliation with its provider.
282    pub async fn sync_next(&mut self) -> Result<SyncReport, DbError> {
283        self.core.require_local_permission(Permission::Read)?;
284        loop {
285            let event = self
286                .receiver
287                .next()
288                .await
289                .ok_or_else(|| protocol_error("gossip subscription closed"))?
290                .map_err(network_error)?;
291            let GossipEvent::Received(message) = event else {
292                continue;
293            };
294            let provider = author_id(message.delivered_from);
295            if !self.core.peer_has_permission(provider, Permission::Read) {
296                continue;
297            }
298            let Ok(hint) = GossipHint::decode_canonical(&message.content) else {
299                continue;
300            };
301            if !self
302                .core
303                .peer_has_permission(hint.author(), Permission::Read)
304                || verify_gossip_signature(
305                    hint.author(),
306                    &hint.signing_bytes().map_err(sync_error)?,
307                    hint.signature(),
308                )
309                .is_err()
310                || self
311                    .seen
312                    .get(&hint.author())
313                    .is_some_and(|counter| *counter >= hint.counter())
314            {
315                continue;
316            }
317            self.seen.insert(hint.author(), hint.counter());
318            return self
319                .core
320                .sync_with(EndpointAddr::new(message.delivered_from))
321                .await;
322        }
323    }
324}
325
326#[derive(Clone)]
327struct SyncCore {
328    db: IrohDb,
329    endpoint: Endpoint,
330    pending_invitations: Arc<RwLock<BTreeMap<AuthorId, Vec<u8>>>>,
331}
332
333#[derive(Clone, Copy)]
334struct AuthorityPosition {
335    epoch: u64,
336    sequence: u64,
337    head: Option<[u8; 32]>,
338}
339
340struct HandshakeState {
341    remote_authority: AuthorityPosition,
342    remote_capability_chain: Vec<Vec<u8>>,
343}
344
345#[derive(Clone)]
346struct SyncProtocol {
347    core: SyncCore,
348}
349
350impl std::fmt::Debug for SyncProtocol {
351    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
352        formatter.write_str("SyncProtocol")
353    }
354}
355
356impl IrohDb {
357    /// Starts a local iroh endpoint with sync and iroh-blobs protocol handlers.
358    pub async fn start_sync(&self) -> Result<SyncNode, DbError> {
359        let lookup = MemoryLookup::new();
360        let endpoint = Endpoint::builder(presets::Minimal)
361            .secret_key(self.inner.credentials.endpoint_secret_key())
362            .address_lookup(lookup.clone())
363            .relay_mode(RelayMode::Disabled)
364            .bind()
365            .await
366            .map_err(|error| DbError::Network(error.to_string()))?;
367        let core = SyncCore {
368            db: self.clone(),
369            endpoint: endpoint.clone(),
370            pending_invitations: Arc::new(RwLock::new(BTreeMap::new())),
371        };
372        let gossip = Gossip::builder().spawn(endpoint.clone());
373        let blob_events = core.blob_provider_events();
374        let router = Router::builder(endpoint)
375            .accept(SYNC_ALPN, SyncProtocol { core: core.clone() })
376            .accept(
377                iroh_blobs::ALPN,
378                self.inner.blobs.protocol_with_events(blob_events),
379            )
380            .accept(iroh_gossip::ALPN, gossip.clone())
381            .spawn();
382        Ok(SyncNode {
383            core,
384            gossip,
385            lookup,
386            router,
387        })
388    }
389}
390
391impl ProtocolHandler for SyncProtocol {
392    async fn accept(&self, connection: Connection) -> Result<(), AcceptError> {
393        let peer = connection.remote_id();
394        let (mut send, mut recv) = connection.accept_bi().await?;
395        self.core
396            .serve(peer, &mut send, &mut recv)
397            .await
398            .map_err(AcceptError::from_err)?;
399        send.finish()?;
400        connection.closed().await;
401        Ok(())
402    }
403}
404
405impl SyncCore {
406    #[tracing::instrument(name = "iroh_db.sync", skip_all)]
407    async fn sync_with(&self, remote: EndpointAddr) -> Result<SyncReport, DbError> {
408        let result = self.sync_with_inner(remote).await;
409        let (received, sent) = result
410            .as_ref()
411            .map_or((0, 0), |report| (report.received(), report.sent()));
412        self.db
413            .inner
414            .telemetry
415            .record_sync(received, sent, result.is_err());
416        result
417    }
418
419    async fn sync_with_inner(&self, remote: EndpointAddr) -> Result<SyncReport, DbError> {
420        let peer = author_id(remote.id);
421        self.require_local_blob_fetch()?;
422        self.require_trusted(peer)?;
423        let connection = self
424            .endpoint
425            .connect(remote.clone(), SYNC_ALPN)
426            .await
427            .map_err(network_error)?;
428        let blobs_connection = self
429            .endpoint
430            .connect(remote, iroh_blobs::ALPN)
431            .await
432            .map_err(network_error)?;
433        let (mut send, mut recv) = connection.open_bi().await.map_err(network_error)?;
434
435        let handshake = self.outbound_handshake(peer, &mut send, &mut recv).await?;
436        let mut remote_authority = handshake.remote_authority;
437        let mut report = SyncReport {
438            received: 0,
439            sent: 0,
440        };
441        while self.db.epoch() != remote_authority.epoch {
442            let boundary = self.db.epoch().min(remote_authority.epoch);
443            let round = self
444                .outbound_commit_round(
445                    peer,
446                    &mut send,
447                    &mut recv,
448                    blobs_connection.clone(),
449                    boundary,
450                )
451                .await?;
452            report.received += round.received;
453            report.sent += round.sent;
454            if self.db.epoch() > remote_authority.epoch {
455                let next_sequence = remote_authority
456                    .sequence
457                    .checked_add(1)
458                    .ok_or_else(|| protocol_error("remote control sequence exhausted"))?;
459                let disclosure = self.db.inner.control_gate.read().await;
460                let grant = self.control_grant(next_sequence, peer)?;
461                send_frame(&mut send, &grant).await?;
462                drop(disclosure);
463                expect_done(&mut recv).await?;
464                remote_authority.epoch = remote_authority
465                    .epoch
466                    .checked_add(1)
467                    .ok_or_else(|| protocol_error("remote control epoch exhausted"))?;
468                remote_authority.sequence = next_sequence;
469                remote_authority.head = Some(self.control_id(next_sequence)?);
470            } else {
471                let grant = recv_frame(&mut recv).await?;
472                self.apply_control_grant(grant)?;
473                send_frame(&mut send, &SyncFrame::Done).await?;
474            }
475        }
476        self.finalize_peer_authority(
477            peer,
478            remote_authority,
479            &handshake.remote_capability_chain,
480            &mut send,
481            &mut recv,
482        )
483        .await?;
484        let round = self
485            .outbound_commit_round(
486                peer,
487                &mut send,
488                &mut recv,
489                blobs_connection,
490                self.db.epoch(),
491            )
492            .await?;
493        report.received += round.received;
494        report.sent += round.sent;
495        send.finish().map_err(network_error)?;
496        Ok(report)
497    }
498
499    async fn outbound_commit_round(
500        &self,
501        peer: AuthorId,
502        send: &mut SendStream,
503        recv: &mut RecvStream,
504        blobs_connection: Connection,
505        max_epoch: u64,
506    ) -> Result<SyncReport, DbError> {
507        let local = self.commit_ids_through(max_epoch)?;
508        send_commit_pages(send, FrameKind::Frontier, &local).await?;
509        let remote = recv_commit_pages(recv, FrameKind::Frontier).await?;
510        let receive = missing_commits(&local, &remote).map_err(sync_error)?;
511        send_commit_pages(send, FrameKind::Need, &receive).await?;
512        let to_send = recv_commit_pages(recv, FrameKind::Need).await?;
513        let expected_send = missing_commits(&remote, &local).map_err(sync_error)?;
514        if to_send != expected_send {
515            return Err(protocol_error("peer requested an invalid reference set"));
516        }
517
518        let received = self
519            .fetch_and_apply_closure(blobs_connection, peer, &receive)
520            .await?;
521        send_commit_pages(send, FrameKind::Refs, &to_send).await?;
522        send_frame(send, &SyncFrame::Done).await?;
523        expect_done(recv).await?;
524        Ok(SyncReport {
525            received,
526            sent: to_send.len(),
527        })
528    }
529
530    async fn outbound_handshake(
531        &self,
532        peer: AuthorId,
533        send: &mut SendStream,
534        recv: &mut RecvStream,
535    ) -> Result<HandshakeState, DbError> {
536        self.require_local_permission(Permission::Read)?;
537        let local_nonce = random_nonce()?;
538        send_frame(send, &hello_frame(local_nonce)).await?;
539        let SyncFrame::Hello {
540            nonce: remote_nonce,
541            min_protocol,
542            max_protocol,
543            ..
544        } = recv_frame(recv).await?
545        else {
546            return Err(protocol_error("expected hello"));
547        };
548        negotiate_protocol(
549            MIN_PROTOCOL_VERSION,
550            MAX_PROTOCOL_VERSION,
551            min_protocol,
552            max_protocol,
553        )
554        .map_err(network_error)?;
555        send_frame(send, &self.authority_state()?).await?;
556        let SyncFrame::AuthorityState {
557            epoch: remote_epoch,
558            sequence,
559            head,
560        } = recv_frame(recv).await?
561        else {
562            return Err(protocol_error("expected authority state"));
563        };
564        let proof_epoch = self.db.epoch().min(remote_epoch);
565        send_frame(
566            send,
567            &SyncFrame::DomainRequest {
568                domain_id: self.domain_id(),
569                epoch: proof_epoch,
570                proof: self.proof(local_nonce, remote_nonce, self.db.author_id(), proof_epoch)?,
571                capability_chain: self.local_capability_chain()?,
572            },
573        )
574        .await?;
575        match recv_frame(recv).await? {
576            SyncFrame::DomainRequest {
577                domain_id,
578                epoch,
579                proof,
580                capability_chain,
581            } if domain_id == self.domain_id()
582                && epoch == proof_epoch
583                && proof == self.proof(local_nonce, remote_nonce, peer, proof_epoch)? =>
584            {
585                self.authenticate_peer_capability(peer, &capability_chain)
586                    .await?;
587                Ok(HandshakeState {
588                    remote_authority: AuthorityPosition {
589                        epoch: remote_epoch,
590                        sequence,
591                        head,
592                    },
593                    remote_capability_chain: capability_chain,
594                })
595            }
596            _ => Err(DbError::UnauthorizedDevice),
597        }
598    }
599
600    async fn fetch_and_apply_closure(
601        &self,
602        connection: Connection,
603        peer: AuthorId,
604        roots: &[CommitId],
605    ) -> Result<usize, DbError> {
606        self.require_local_blob_fetch()?;
607        let mut queued: std::collections::BTreeSet<_> = roots.iter().copied().collect();
608        let mut pending = BTreeMap::<CommitId, CommitEnvelope>::new();
609        while let Some(commit_id) = queued.pop_first() {
610            if self
611                .db
612                .inner
613                .store
614                .contains_commit(self.domain_id(), commit_id)?
615                || pending.contains_key(&commit_id)
616            {
617                continue;
618            }
619            if pending.len() >= MAX_SESSION_REFERENCES {
620                return Err(protocol_error("dependency closure exceeds session budget"));
621            }
622            self.db
623                .inner
624                .blobs
625                .fetch(
626                    connection.clone(),
627                    BlobHash::from_bytes(commit_id.to_bytes()),
628                )
629                .await
630                .map_err(network_error)?;
631            let bytes = self
632                .read_transferred(
633                    commit_id,
634                    tokio::time::Instant::now() + Duration::from_secs(10),
635                )
636                .await?;
637            let envelope = CommitEnvelope::decode_canonical(&bytes)?;
638            if envelope.commit_id() != commit_id
639                || envelope.header().domain_id() != self.domain_id()
640            {
641                return Err(protocol_error("transferred commit identity mismatch"));
642            }
643            for dependency in envelope.header().dependencies() {
644                if !self
645                    .db
646                    .inner
647                    .store
648                    .contains_commit(self.domain_id(), *dependency)?
649                    && !pending.contains_key(dependency)
650                {
651                    queued.insert(*dependency);
652                }
653            }
654            pending.insert(commit_id, envelope);
655        }
656        let mut applied = 0_usize;
657        while !pending.is_empty() {
658            let ready = pending.iter().find_map(|(&commit_id, envelope)| {
659                envelope
660                    .header()
661                    .dependencies()
662                    .iter()
663                    .all(|dependency| {
664                        self.db
665                            .inner
666                            .store
667                            .contains_commit(self.domain_id(), *dependency)
668                            .unwrap_or(false)
669                    })
670                    .then_some(commit_id)
671            });
672            let commit_id = ready.ok_or_else(|| protocol_error("dependency graph is cyclic"))?;
673            let envelope = pending
674                .remove(&commit_id)
675                .expect("selected pending commit exists");
676            if matches!(
677                self.db.accept_remote_commit(peer, envelope).await?,
678                crate::RemoteApply::Applied(_)
679            ) {
680                applied = applied.saturating_add(1);
681            }
682        }
683        Ok(applied)
684    }
685
686    #[allow(clippy::too_many_lines)]
687    async fn serve(
688        &self,
689        peer_endpoint: iroh::EndpointId,
690        send: &mut SendStream,
691        recv: &mut RecvStream,
692    ) -> Result<(), DbError> {
693        let peer = author_id(peer_endpoint);
694        let SyncFrame::Hello {
695            nonce: remote_nonce,
696            min_protocol,
697            max_protocol,
698            ..
699        } = recv_frame(recv).await?
700        else {
701            return Err(protocol_error("expected hello"));
702        };
703        negotiate_protocol(
704            MIN_PROTOCOL_VERSION,
705            MAX_PROTOCOL_VERSION,
706            min_protocol,
707            max_protocol,
708        )
709        .map_err(network_error)?;
710        let local_nonce = random_nonce()?;
711        send_frame(send, &hello_frame(local_nonce)).await?;
712        let next = recv_frame(recv).await?;
713        if next == SyncFrame::InvitationRequest {
714            return self.serve_invitation(peer, send).await;
715        }
716        let SyncFrame::AuthorityState {
717            epoch: remote_epoch,
718            sequence: remote_sequence,
719            head: remote_head,
720        } = next
721        else {
722            return Err(protocol_error("expected authority state"));
723        };
724        send_frame(send, &self.authority_state()?).await?;
725        let proof_epoch = self.db.epoch().min(remote_epoch);
726        let SyncFrame::DomainRequest {
727            domain_id,
728            epoch,
729            proof,
730            capability_chain,
731        } = recv_frame(recv).await?
732        else {
733            return Err(DbError::UnauthorizedDevice);
734        };
735        if domain_id != self.domain_id()
736            || epoch != proof_epoch
737            || proof != self.proof(remote_nonce, local_nonce, peer, proof_epoch)?
738        {
739            return Err(DbError::UnauthorizedDevice);
740        }
741        self.authenticate_peer_capability(peer, &capability_chain)
742            .await?;
743        send_frame(
744            send,
745            &SyncFrame::DomainRequest {
746                domain_id: self.domain_id(),
747                epoch: proof_epoch,
748                proof: self.proof(remote_nonce, local_nonce, self.db.author_id(), proof_epoch)?,
749                capability_chain: self.local_capability_chain()?,
750            },
751        )
752        .await?;
753        let blobs_connection = self
754            .endpoint
755            .connect(EndpointAddr::new(peer_endpoint), iroh_blobs::ALPN)
756            .await
757            .map_err(network_error)?;
758        let mut remote_authority = AuthorityPosition {
759            epoch: remote_epoch,
760            sequence: remote_sequence,
761            head: remote_head,
762        };
763        while self.db.epoch() != remote_authority.epoch {
764            let boundary = self.db.epoch().min(remote_authority.epoch);
765            self.serve_commit_round(peer, send, recv, blobs_connection.clone(), boundary)
766                .await?;
767            if self.db.epoch() > remote_authority.epoch {
768                let next_sequence = remote_authority
769                    .sequence
770                    .checked_add(1)
771                    .ok_or_else(|| protocol_error("remote control sequence exhausted"))?;
772                let disclosure = self.db.inner.control_gate.read().await;
773                let grant = self.control_grant(next_sequence, peer)?;
774                send_frame(send, &grant).await?;
775                drop(disclosure);
776                expect_done(recv).await?;
777                remote_authority.epoch = remote_authority
778                    .epoch
779                    .checked_add(1)
780                    .ok_or_else(|| protocol_error("remote control epoch exhausted"))?;
781                remote_authority.sequence = next_sequence;
782                remote_authority.head = Some(self.control_id(next_sequence)?);
783            } else {
784                let grant = recv_frame(recv).await?;
785                self.apply_control_grant(grant)?;
786                send_frame(send, &SyncFrame::Done).await?;
787            }
788        }
789        self.finalize_peer_authority(peer, remote_authority, &capability_chain, send, recv)
790            .await?;
791        self.serve_commit_round(peer, send, recv, blobs_connection, self.db.epoch())
792            .await
793    }
794
795    async fn finalize_peer_authority(
796        &self,
797        peer: AuthorId,
798        remote_authority: AuthorityPosition,
799        capability_chain: &[Vec<u8>],
800        send: &mut SendStream,
801        recv: &mut RecvStream,
802    ) -> Result<(), DbError> {
803        self.reconcile_authority_head(remote_authority, send, recv)
804            .await?;
805        self.authenticate_peer_capability(peer, capability_chain)
806            .await
807    }
808
809    async fn serve_invitation(&self, peer: AuthorId, send: &mut SendStream) -> Result<(), DbError> {
810        let control_reader = self.db.inner.control_gate.read().await;
811        self.require_local_permission(Permission::Invite)?;
812        let payload = self
813            .pending_invitations
814            .read()
815            .map_err(|_| protocol_error("invitation registry is unavailable"))?
816            .get(&peer)
817            .cloned()
818            .ok_or(DbError::UnauthorizedDevice)?;
819        let invitation = Invitation::decode_canonical(&payload)?;
820        if invitation.domain_seed().epoch() != self.db.epoch() {
821            return Err(DbError::UnauthorizedDevice);
822        }
823        self.pending_invitations
824            .write()
825            .map_err(|_| protocol_error("invitation registry is unavailable"))?
826            .remove(&peer);
827        let result = send_frame(send, &SyncFrame::Invitation { payload }).await;
828        drop(control_reader);
829        result
830    }
831
832    async fn serve_commit_round(
833        &self,
834        peer: AuthorId,
835        send: &mut SendStream,
836        recv: &mut RecvStream,
837        blobs_connection: Connection,
838        max_epoch: u64,
839    ) -> Result<(), DbError> {
840        let remote = recv_commit_pages(recv, FrameKind::Frontier).await?;
841        let local = self.commit_ids_through(max_epoch)?;
842        send_commit_pages(send, FrameKind::Frontier, &local).await?;
843        let remote_need = recv_commit_pages(recv, FrameKind::Need).await?;
844        let expected_remote_need = missing_commits(&remote, &local).map_err(sync_error)?;
845        if remote_need != expected_remote_need {
846            return Err(protocol_error("peer requested an invalid reference set"));
847        }
848        let need = missing_commits(&local, &remote).map_err(sync_error)?;
849        send_commit_pages(send, FrameKind::Need, &need).await?;
850        let pushed = recv_commit_pages(recv, FrameKind::Refs).await?;
851        if pushed != need {
852            return Err(protocol_error("pushed reference set mismatch"));
853        }
854        self.fetch_and_apply_closure(blobs_connection, peer, &pushed)
855            .await?;
856        expect_done(recv).await?;
857        send_frame(send, &SyncFrame::Done).await
858    }
859
860    fn domain_id(&self) -> DomainId {
861        self.db.domain_id()
862    }
863
864    fn blob_engine(&self) -> BlobEngine {
865        self.db.blob_engine()
866    }
867
868    fn require_local_blob_fetch(&self) -> Result<(), DbError> {
869        self.require_local_permission(Permission::Read)?;
870        self.require_local_permission(Permission::FetchBlobs)?;
871        Ok(())
872    }
873
874    fn require_local_permission(&self, permission: Permission) -> Result<(), DbError> {
875        let active_epoch = self.db.inner.store.current_epoch(self.domain_id())?;
876        if self.db.epoch() != active_epoch {
877            return Err(DbError::StaleEpoch {
878                handle: self.db.epoch(),
879                active: active_epoch,
880            });
881        }
882        self.db
883            .capability_for_permission(self.db.author_id(), permission, active_epoch)?;
884        Ok(())
885    }
886
887    fn peer_has_permission(&self, peer: AuthorId, permission: Permission) -> bool {
888        let Ok(active_epoch) = self.db.inner.store.current_epoch(self.domain_id()) else {
889            return false;
890        };
891        self.db.epoch() == active_epoch
892            && self
893                .db
894                .capability_for_permission(peer, permission, active_epoch)
895                .is_ok()
896    }
897
898    fn peer_has_permission_at(&self, peer: AuthorId, permission: Permission, epoch: u64) -> bool {
899        let Ok(Some(capability_id)) = self
900            .db
901            .inner
902            .store
903            .active_capability(self.domain_id(), peer)
904        else {
905            return false;
906        };
907        self.db
908            .inner
909            .store
910            .authorize_capability(self.domain_id(), peer, capability_id, permission, epoch)
911            .unwrap_or(false)
912    }
913
914    async fn authenticate_peer_capability(
915        &self,
916        peer: AuthorId,
917        encoded_chain: &[Vec<u8>],
918    ) -> Result<(), DbError> {
919        let _control_writer = self.db.inner.control_gate.write().await;
920        let Some(descriptor) = self.db.inner.store.domain_descriptor(self.domain_id())? else {
921            return self.require_trusted(peer);
922        };
923        let chain = encoded_chain
924            .iter()
925            .map(|bytes| CapabilityCertificate::decode_canonical(bytes).map_err(DbError::from))
926            .collect::<Result<Vec<_>, _>>()?;
927        let root = chain.first().ok_or(DbError::UnauthorizedDevice)?;
928        if root.domain_id() != self.domain_id()
929            || root.subject() != descriptor.owner()
930            || root.issuer_capability().is_some()
931        {
932            return Err(DbError::UnauthorizedDevice);
933        }
934        for pair in chain.windows(2) {
935            pair[1].validate_delegation(&pair[0])?;
936        }
937        let leaf = chain.last().ok_or(DbError::UnauthorizedDevice)?;
938        if leaf.subject() != peer {
939            return Err(DbError::UnauthorizedDevice);
940        }
941        let current_sequence = self
942            .db
943            .inner
944            .store
945            .current_control_sequence(self.domain_id())?;
946        let anchor_is_ahead = chain
947            .iter()
948            .any(|capability| capability.issued_control_sequence() > current_sequence);
949        if !anchor_is_ahead
950            && (!leaf.permits(Permission::Read, self.db.epoch())
951                || !leaf.permits(Permission::FetchBlobs, self.db.epoch()))
952        {
953            return Err(DbError::UnauthorizedDevice);
954        }
955        let leaf_id = match self.db.inner.store.install_capability_chain_atomic(&chain) {
956            Ok(leaf_id) => leaf_id,
957            Err(iroh_db_store::StoreError::InvalidAuthority(_)) if anchor_is_ahead => {
958                let Some(active) = self
959                    .db
960                    .inner
961                    .store
962                    .active_capability(self.domain_id(), peer)?
963                else {
964                    return Err(DbError::UnauthorizedDevice);
965                };
966                if !self
967                    .db
968                    .inner
969                    .store
970                    .is_author_trusted(self.domain_id(), peer)?
971                    || !self.db.inner.store.authorize_capability(
972                        self.domain_id(),
973                        peer,
974                        active,
975                        Permission::Read,
976                        self.db.epoch(),
977                    )?
978                    || !self.db.inner.store.authorize_capability(
979                        self.domain_id(),
980                        peer,
981                        active,
982                        Permission::FetchBlobs,
983                        self.db.epoch(),
984                    )?
985                {
986                    return Err(DbError::UnauthorizedDevice);
987                }
988                return Ok(());
989            }
990            Err(error) => return Err(DbError::Store(error)),
991        };
992        if !self.db.inner.store.authorize_capability(
993            self.domain_id(),
994            peer,
995            leaf_id,
996            Permission::Read,
997            self.db.epoch(),
998        )? || !self.db.inner.store.authorize_capability(
999            self.domain_id(),
1000            peer,
1001            leaf_id,
1002            Permission::FetchBlobs,
1003            self.db.epoch(),
1004        )? {
1005            return Err(DbError::UnauthorizedDevice);
1006        }
1007        self.db.inner.store.trust_author(self.domain_id(), peer)?;
1008        Ok(())
1009    }
1010
1011    fn can_peer_fetch_blob(&self, peer: AuthorId, hash: BlobHash) -> bool {
1012        self.peer_has_permission(peer, Permission::Read)
1013            && self.peer_has_permission(peer, Permission::FetchBlobs)
1014            && self
1015                .db
1016                .inner
1017                .store
1018                .owns_blob(self.domain_id(), hash)
1019                .unwrap_or(false)
1020    }
1021
1022    fn blob_provider_events(&self) -> EventSender {
1023        let mask = EventMask {
1024            connected: ConnectMode::Intercept,
1025            get: RequestMode::Intercept,
1026            get_many: RequestMode::Disabled,
1027            push: RequestMode::Disabled,
1028            observe: ObserveMode::Intercept,
1029            ..EventMask::DEFAULT
1030        };
1031        let (events, mut receiver) = EventSender::channel(64, mask);
1032        let core = self.clone();
1033        tokio::spawn(async move {
1034            let mut peers = BTreeMap::<u64, AuthorId>::new();
1035            while let Some(message) = receiver.recv().await {
1036                match message {
1037                    ProviderMessage::ClientConnected(message) => {
1038                        let peer = message.endpoint_id.map(author_id);
1039                        let allowed = peer.is_some_and(|peer| {
1040                            core.peer_has_permission(peer, Permission::Read)
1041                                && core.peer_has_permission(peer, Permission::FetchBlobs)
1042                        });
1043                        if let Some(peer) = peer.filter(|_| allowed) {
1044                            peers.insert(message.connection_id, peer);
1045                        }
1046                        let result = allowed.then_some(()).ok_or(AbortReason::Permission);
1047                        let _ignored = message.tx.send(result).await;
1048                    }
1049                    ProviderMessage::ConnectionClosed(message) => {
1050                        peers.remove(&message.connection_id);
1051                    }
1052                    ProviderMessage::GetRequestReceived(message) => {
1053                        let hash = BlobHash::from_bytes(*message.request.hash.as_bytes());
1054                        let allowed = message.request.ranges.is_blob()
1055                            && peers
1056                                .get(&message.connection_id)
1057                                .is_some_and(|peer| core.can_peer_fetch_blob(*peer, hash));
1058                        let result = allowed.then_some(()).ok_or(AbortReason::Permission);
1059                        let _ignored = message.tx.send(result).await;
1060                    }
1061                    ProviderMessage::GetManyRequestReceived(message) => {
1062                        let _ignored = message.tx.send(Err(AbortReason::Permission)).await;
1063                    }
1064                    ProviderMessage::PushRequestReceived(message) => {
1065                        let _ignored = message.tx.send(Err(AbortReason::Permission)).await;
1066                    }
1067                    ProviderMessage::ObserveRequestReceived(message) => {
1068                        let _ignored = message.tx.send(Err(AbortReason::Permission)).await;
1069                    }
1070                    ProviderMessage::Throttle(message) => {
1071                        let _ignored = message.tx.send(Err(AbortReason::Permission)).await;
1072                    }
1073                    ProviderMessage::ClientConnectedNotify(_)
1074                    | ProviderMessage::GetRequestReceivedNotify(_)
1075                    | ProviderMessage::GetManyRequestReceivedNotify(_)
1076                    | ProviderMessage::PushRequestReceivedNotify(_)
1077                    | ProviderMessage::ObserveRequestReceivedNotify(_) => {}
1078                }
1079            }
1080        });
1081        events
1082    }
1083
1084    fn require_trusted(&self, peer: AuthorId) -> Result<(), DbError> {
1085        if self
1086            .db
1087            .inner
1088            .store
1089            .is_author_trusted(self.domain_id(), peer)?
1090        {
1091            Ok(())
1092        } else {
1093            Err(DbError::UnauthorizedDevice)
1094        }
1095    }
1096
1097    fn proof(
1098        &self,
1099        initiator_nonce: [u8; 32],
1100        responder_nonce: [u8; 32],
1101        claimant: AuthorId,
1102        epoch: u64,
1103    ) -> Result<[u8; 32], DbError> {
1104        let key = self
1105            .db
1106            .inner
1107            .credentials
1108            .load_domain_epoch(self.domain_id(), epoch)?
1109            .domain_key()
1110            .derive_subkey(self.domain_id().as_bytes(), b"iroh-db/sync-domain-proof/v1")?;
1111        let mut message = Vec::with_capacity(128);
1112        message.extend_from_slice(b"iroh-db/sync-domain-proof/v1");
1113        message.extend_from_slice(self.domain_id().as_bytes());
1114        message.extend_from_slice(&epoch.to_be_bytes());
1115        message.extend_from_slice(&initiator_nonce);
1116        message.extend_from_slice(&responder_nonce);
1117        message.extend_from_slice(claimant.as_bytes());
1118        Ok(*blake3::keyed_hash(&key, &message).as_bytes())
1119    }
1120
1121    fn authority_state(&self) -> Result<SyncFrame, DbError> {
1122        let head = self.db.inner.store.control_head(self.domain_id())?;
1123        Ok(SyncFrame::AuthorityState {
1124            epoch: self.db.inner.store.current_epoch(self.domain_id())?,
1125            sequence: head.as_ref().map_or(0, ControlTransition::sequence),
1126            head: head.as_ref().map(ControlTransition::id).transpose()?,
1127        })
1128    }
1129
1130    fn local_capability_chain(&self) -> Result<Vec<Vec<u8>>, DbError> {
1131        let Some(capability) = self.db.active_capability()? else {
1132            return Ok(Vec::new());
1133        };
1134        self.db
1135            .inner
1136            .store
1137            .capability_chain(capability.id()?)?
1138            .iter()
1139            .map(|capability| capability.encode_canonical().map_err(DbError::from))
1140            .collect()
1141    }
1142
1143    fn control_grant(&self, sequence: u64, peer: AuthorId) -> Result<SyncFrame, DbError> {
1144        if sequence == 0 {
1145            return Err(protocol_error("control grant sequence cannot be zero"));
1146        }
1147        let control = self
1148            .db
1149            .inner
1150            .store
1151            .accepted_control_at_sequence(self.domain_id(), sequence)?
1152            .ok_or_else(|| protocol_error("requested control transition is unavailable"))?;
1153        if !self
1154            .db
1155            .inner
1156            .store
1157            .is_author_trusted(self.domain_id(), peer)?
1158            || !self.peer_has_permission_at(peer, Permission::Read, control.new_epoch())
1159            || !self.peer_has_permission_at(peer, Permission::FetchBlobs, control.new_epoch())
1160        {
1161            return Err(DbError::UnauthorizedDevice);
1162        }
1163        let capability_chain = self
1164            .db
1165            .inner
1166            .store
1167            .capability_chain(control.controller_capability())?
1168            .iter()
1169            .map(|capability| capability.encode_canonical().map_err(DbError::from))
1170            .collect::<Result<Vec<_>, _>>()?;
1171        let seed = self
1172            .db
1173            .inner
1174            .credentials
1175            .load_domain_epoch(self.domain_id(), control.new_epoch())?;
1176        Ok(SyncFrame::ControlGrant {
1177            transition: control.encode_canonical()?,
1178            epoch_key: seed.epoch_key_bytes(),
1179            capability_chain,
1180        })
1181    }
1182
1183    fn control_id(&self, sequence: u64) -> Result<[u8; 32], DbError> {
1184        if sequence == 0 {
1185            return Err(protocol_error("control sequence cannot be zero"));
1186        }
1187        self.db
1188            .inner
1189            .store
1190            .accepted_control_at_sequence(self.domain_id(), sequence)?
1191            .ok_or_else(|| protocol_error("control transition is unavailable"))?
1192            .id()
1193            .map_err(DbError::from)
1194    }
1195
1196    async fn reconcile_authority_head(
1197        &self,
1198        remote: AuthorityPosition,
1199        send: &mut SendStream,
1200        recv: &mut RecvStream,
1201    ) -> Result<(), DbError> {
1202        let SyncFrame::AuthorityState {
1203            epoch,
1204            sequence,
1205            head,
1206        } = self.authority_state()?
1207        else {
1208            unreachable!("authority_state always returns AuthorityState")
1209        };
1210        if epoch == remote.epoch && sequence == remote.sequence && head == remote.head {
1211            return Ok(());
1212        }
1213        if epoch != remote.epoch || sequence != remote.sequence {
1214            return Err(protocol_error("authority positions do not converge"));
1215        }
1216        let local_head = head.ok_or_else(|| protocol_error("local control head is missing"))?;
1217        let remote_head = remote
1218            .head
1219            .ok_or_else(|| protocol_error("remote control head is missing"))?;
1220        let local = self
1221            .db
1222            .inner
1223            .store
1224            .control_transition(self.domain_id(), local_head)?
1225            .ok_or_else(|| protocol_error("local control object is missing"))?;
1226        send_frame(
1227            send,
1228            &SyncFrame::ControlEvidence {
1229                transition: local.encode_canonical()?,
1230            },
1231        )
1232        .await?;
1233        let SyncFrame::ControlEvidence { transition } = recv_frame(recv).await? else {
1234            return Err(protocol_error("expected control evidence"));
1235        };
1236        let conflicting = ControlTransition::decode_canonical(&transition)?;
1237        if conflicting.id()? != remote_head {
1238            return Err(protocol_error(
1239                "remote control evidence does not match its head",
1240            ));
1241        }
1242        match self.db.inner.store.observe_control(&conflicting) {
1243            Err(error @ iroh_db_store::StoreError::ControlFork { .. }) => {
1244                Err(DbError::Store(error))
1245            }
1246            Err(error) => Err(DbError::Store(error)),
1247            Ok(()) => Err(protocol_error("authority heads do not converge")),
1248        }
1249    }
1250
1251    fn apply_control_grant(&self, frame: SyncFrame) -> Result<(), DbError> {
1252        let SyncFrame::ControlGrant {
1253            transition,
1254            epoch_key,
1255            capability_chain,
1256        } = frame
1257        else {
1258            return Err(protocol_error("expected control grant"));
1259        };
1260        let control = ControlTransition::decode_canonical(&transition)?;
1261        self.validate_incoming_control(&control)?;
1262        let seed = DomainSeed::from_epoch_key(self.domain_id(), control.new_epoch(), epoch_key);
1263        if seed.distribution_digest() != control.key_distribution_digest()
1264            || control.revoked_subjects().contains(&self.db.author_id())
1265        {
1266            return Err(DbError::UnauthorizedDevice);
1267        }
1268        let capability_chain = capability_chain
1269            .iter()
1270            .map(|bytes| CapabilityCertificate::decode_canonical(bytes).map_err(DbError::from))
1271            .collect::<Result<Vec<_>, _>>()?;
1272        for capability in &capability_chain {
1273            if capability.domain_id() != self.domain_id() {
1274                return Err(DbError::UnauthorizedDevice);
1275            }
1276        }
1277        self.db
1278            .apply_epoch_rotation_with_capabilities(&control, &seed, &capability_chain)?;
1279        for subject in [control.author(), control.controller()] {
1280            if !control.revoked_subjects().contains(&subject) {
1281                self.db
1282                    .inner
1283                    .store
1284                    .trust_author(self.domain_id(), subject)?;
1285            }
1286        }
1287        Ok(())
1288    }
1289
1290    fn validate_incoming_control(&self, control: &ControlTransition) -> Result<(), DbError> {
1291        let head = self.db.inner.store.control_head(self.domain_id())?;
1292        let expected_sequence = head
1293            .as_ref()
1294            .map_or(Some(1), |head| head.sequence().checked_add(1))
1295            .ok_or_else(|| protocol_error("control sequence exhausted"))?;
1296        let expected_predecessor = head.as_ref().map(ControlTransition::id).transpose()?;
1297        let can_control = self.db.inner.store.authorize_capability(
1298            self.domain_id(),
1299            control.author(),
1300            control.author_capability(),
1301            Permission::Revoke,
1302            self.db.epoch(),
1303        )? || self.db.inner.store.authorize_capability(
1304            self.domain_id(),
1305            control.author(),
1306            control.author_capability(),
1307            Permission::Admin,
1308            self.db.epoch(),
1309        )?;
1310        if control.domain_id() != self.domain_id()
1311            || control.sequence() != expected_sequence
1312            || control.predecessor() != expected_predecessor
1313            || control.previous_epoch() != self.db.epoch()
1314            || control.cut() != self.db.inner.store.frontier(self.domain_id())?
1315            || control.author() != self.db.inner.store.current_controller(self.domain_id())?
1316            || !can_control
1317        {
1318            return Err(DbError::UnauthorizedDevice);
1319        }
1320        Ok(())
1321    }
1322
1323    fn commit_ids_through(&self, epoch: u64) -> Result<Vec<CommitId>, DbError> {
1324        if epoch == self.db.inner.store.current_epoch(self.domain_id())? {
1325            self.db
1326                .inner
1327                .store
1328                .frontier(self.domain_id())
1329                .map_err(DbError::from)
1330        } else {
1331            self.db
1332                .inner
1333                .store
1334                .list_stored_commit_ids_through(self.domain_id(), epoch)
1335                .map_err(DbError::from)
1336        }
1337    }
1338
1339    fn commit_ids(&self) -> Result<Vec<CommitId>, DbError> {
1340        self.db
1341            .inner
1342            .store
1343            .frontier(self.domain_id())
1344            .map_err(DbError::from)
1345    }
1346
1347    fn gossip_topic(&self) -> Result<TopicId, DbError> {
1348        let topic = self
1349            .db
1350            .domain_key()
1351            .derive_subkey(self.domain_id().as_bytes(), b"iroh-db/gossip-topic/v1")?;
1352        Ok(TopicId::from_bytes(*topic))
1353    }
1354
1355    async fn read_transferred(
1356        &self,
1357        commit_id: CommitId,
1358        deadline: tokio::time::Instant,
1359    ) -> Result<Vec<u8>, DbError> {
1360        let hash = BlobHash::from_bytes(commit_id.to_bytes());
1361        loop {
1362            if let Some(bytes) = self
1363                .db
1364                .inner
1365                .blobs
1366                .read(hash)
1367                .await
1368                .map_err(network_error)?
1369            {
1370                return Ok(bytes);
1371            }
1372            if tokio::time::Instant::now() >= deadline {
1373                return Err(protocol_error("transferred commit is missing"));
1374            }
1375            tokio::time::sleep(Duration::from_millis(5)).await;
1376        }
1377    }
1378}
1379
1380#[derive(Clone, Copy)]
1381enum FrameKind {
1382    Frontier,
1383    Need,
1384    Refs,
1385}
1386
1387fn expect_commits(frame: SyncFrame, kind: FrameKind) -> Result<Vec<CommitId>, DbError> {
1388    match (kind, frame) {
1389        (FrameKind::Frontier, SyncFrame::Frontier { commits })
1390        | (FrameKind::Need, SyncFrame::Need { commits })
1391        | (FrameKind::Refs, SyncFrame::Refs { commits }) => Ok(commits),
1392        _ => Err(protocol_error("unexpected control frame")),
1393    }
1394}
1395
1396async fn send_commit_pages(
1397    send: &mut SendStream,
1398    kind: FrameKind,
1399    commits: &[CommitId],
1400) -> Result<(), DbError> {
1401    if commits.len() > MAX_SESSION_REFERENCES {
1402        return Err(protocol_error("session reference budget exceeded"));
1403    }
1404    for page in commits.chunks(MAX_REFERENCES) {
1405        let frame = match kind {
1406            FrameKind::Frontier => SyncFrame::Frontier {
1407                commits: page.to_vec(),
1408            },
1409            FrameKind::Need => SyncFrame::Need {
1410                commits: page.to_vec(),
1411            },
1412            FrameKind::Refs => SyncFrame::Refs {
1413                commits: page.to_vec(),
1414            },
1415        };
1416        send_frame(send, &frame).await?;
1417    }
1418    let terminator = match kind {
1419        FrameKind::Frontier => SyncFrame::Frontier {
1420            commits: Vec::new(),
1421        },
1422        FrameKind::Need => SyncFrame::Need {
1423            commits: Vec::new(),
1424        },
1425        FrameKind::Refs => SyncFrame::Refs {
1426            commits: Vec::new(),
1427        },
1428    };
1429    send_frame(send, &terminator).await
1430}
1431
1432async fn recv_commit_pages(
1433    recv: &mut RecvStream,
1434    kind: FrameKind,
1435) -> Result<Vec<CommitId>, DbError> {
1436    let mut commits = Vec::new();
1437    loop {
1438        let page = expect_commits(recv_frame(recv).await?, kind)?;
1439        if page.is_empty() {
1440            return Ok(commits);
1441        }
1442        if commits.len().saturating_add(page.len()) > MAX_SESSION_REFERENCES
1443            || commits.last().is_some_and(|last| *last >= page[0])
1444        {
1445            return Err(protocol_error("non-canonical or excessive reference pages"));
1446        }
1447        commits.extend(page);
1448    }
1449}
1450
1451async fn send_frame(send: &mut SendStream, frame: &SyncFrame) -> Result<(), DbError> {
1452    let bytes = frame.encode_canonical().map_err(sync_error)?;
1453    let length = u32::try_from(bytes.len()).map_err(|_| protocol_error("frame is too large"))?;
1454    send.write_all(&length.to_be_bytes())
1455        .await
1456        .map_err(network_error)?;
1457    send.write_all(&bytes).await.map_err(network_error)
1458}
1459
1460async fn recv_frame(recv: &mut RecvStream) -> Result<SyncFrame, DbError> {
1461    let mut length = [0_u8; 4];
1462    recv.read_exact(&mut length).await.map_err(network_error)?;
1463    let length = u32::from_be_bytes(length) as usize;
1464    if length > MAX_CONTROL_FRAME {
1465        return Err(protocol_error("frame exceeds negotiated limit"));
1466    }
1467    let mut bytes = vec![0_u8; length];
1468    recv.read_exact(&mut bytes).await.map_err(network_error)?;
1469    SyncFrame::decode_canonical(&bytes).map_err(sync_error)
1470}
1471
1472async fn expect_done(recv: &mut RecvStream) -> Result<(), DbError> {
1473    if recv_frame(recv).await? == SyncFrame::Done {
1474        Ok(())
1475    } else {
1476        Err(protocol_error("expected done"))
1477    }
1478}
1479
1480fn random_nonce() -> Result<[u8; 32], DbError> {
1481    let mut nonce = [0_u8; 32];
1482    getrandom::fill(&mut nonce).map_err(|_| DbError::RandomnessUnavailable)?;
1483    Ok(nonce)
1484}
1485
1486fn max_control_frame() -> u32 {
1487    u32::try_from(MAX_CONTROL_FRAME).expect("control frame limit fits u32")
1488}
1489
1490fn hello_frame(nonce: [u8; 32]) -> SyncFrame {
1491    SyncFrame::Hello {
1492        nonce,
1493        max_frame: max_control_frame(),
1494        min_protocol: MIN_PROTOCOL_VERSION,
1495        max_protocol: MAX_PROTOCOL_VERSION,
1496    }
1497}
1498
1499fn digest_commits(commits: &[CommitId]) -> [u8; 32] {
1500    let mut hasher = blake3::Hasher::new();
1501    hasher.update(b"iroh-db/frontier-digest/v1");
1502    for commit in commits {
1503        hasher.update(commit.as_bytes());
1504    }
1505    *hasher.finalize().as_bytes()
1506}
1507
1508fn initial_gossip_counter() -> u64 {
1509    std::time::SystemTime::now()
1510        .duration_since(std::time::UNIX_EPOCH)
1511        .map_or(0, |duration| {
1512            u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
1513        })
1514}
1515
1516fn author_id(endpoint_id: iroh::EndpointId) -> AuthorId {
1517    AuthorId::from_bytes(*endpoint_id.as_bytes())
1518}
1519
1520#[allow(clippy::needless_pass_by_value)]
1521fn network_error(error: impl std::fmt::Display) -> DbError {
1522    DbError::Network(error.to_string())
1523}
1524
1525#[allow(clippy::needless_pass_by_value)]
1526fn sync_error(error: impl std::fmt::Display) -> DbError {
1527    DbError::Network(error.to_string())
1528}
1529
1530fn protocol_error(message: &str) -> DbError {
1531    DbError::Network(message.into())
1532}
1533
1534#[cfg(test)]
1535mod tests {
1536    use iroh_db_security::{ConsistencyMode, Permission, PermissionSet};
1537
1538    use super::*;
1539
1540    #[tokio::test]
1541    async fn raw_blob_connection_is_rejected_without_fetch_capability() {
1542        let source_directory = tempfile::tempdir().unwrap();
1543        let target_directory = tempfile::tempdir().unwrap();
1544        let source_device = IrohDb::open(source_directory.path()).await.unwrap();
1545        let target_device = IrohDb::open(target_directory.path()).await.unwrap();
1546        let source = source_device
1547            .create_domain(ConsistencyMode::MultiWriter)
1548            .unwrap();
1549        let invitation = source
1550            .invite(
1551                target_device.author_id(),
1552                PermissionSet::from_permissions(&[Permission::Read]),
1553                PermissionSet::empty(),
1554                None,
1555            )
1556            .unwrap();
1557        let target = target_device.import_invitation(&invitation).unwrap();
1558        let reference = source
1559            .blobs()
1560            .import_bytes(b"provider protected".to_vec(), None)
1561            .await
1562            .unwrap();
1563        let source_node = source.start_sync().await.unwrap();
1564        let target_node = target.start_sync().await.unwrap();
1565        let connection = target_node
1566            .core
1567            .endpoint
1568            .connect(source_node.addr(), iroh_blobs::ALPN)
1569            .await
1570            .unwrap();
1571
1572        assert!(
1573            target
1574                .inner
1575                .blobs
1576                .fetch(connection, reference.manifest_hash())
1577                .await
1578                .is_err()
1579        );
1580        assert!(
1581            target
1582                .inner
1583                .blobs
1584                .read(reference.manifest_hash())
1585                .await
1586                .unwrap()
1587                .is_none()
1588        );
1589
1590        source_node.shutdown().await.unwrap();
1591        target_node.shutdown().await.unwrap();
1592        drop(source);
1593        drop(target);
1594        source_device.close().await.unwrap();
1595        target_device.close().await.unwrap();
1596    }
1597
1598    #[tokio::test]
1599    async fn wrong_epoch_key_is_rejected_before_authority_mutation() {
1600        let owner_directory = tempfile::tempdir().unwrap();
1601        let member_directory = tempfile::tempdir().unwrap();
1602        let owner_device = IrohDb::open(owner_directory.path()).await.unwrap();
1603        let member_device = IrohDb::open(member_directory.path()).await.unwrap();
1604        let owner = owner_device
1605            .create_domain(ConsistencyMode::MultiWriter)
1606            .unwrap();
1607        let invitation = owner
1608            .invite(
1609                member_device.author_id(),
1610                PermissionSet::from_permissions(&[Permission::Read, Permission::FetchBlobs]),
1611                PermissionSet::empty(),
1612                None,
1613            )
1614            .unwrap();
1615        let member = member_device.import_invitation(&invitation).unwrap();
1616        let member_node = member.start_sync().await.unwrap();
1617        let (rotated, control) = owner.rotate_epoch(&[]).unwrap();
1618        let grant = SyncFrame::ControlGrant {
1619            transition: control.encode_canonical().unwrap(),
1620            epoch_key: [0; 32],
1621            capability_chain: Vec::new(),
1622        };
1623
1624        assert!(matches!(
1625            member_node.core.apply_control_grant(grant),
1626            Err(DbError::UnauthorizedDevice)
1627        ));
1628        assert_eq!(member.domain_seed().epoch(), 0);
1629        assert!(
1630            member
1631                .inner
1632                .store
1633                .control_head(member.domain_id())
1634                .unwrap()
1635                .is_none()
1636        );
1637
1638        member_node.shutdown().await.unwrap();
1639        drop(owner);
1640        drop(rotated);
1641        drop(member);
1642        owner_device.close().await.unwrap();
1643        member_device.close().await.unwrap();
1644    }
1645}