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
36pub const SYNC_ALPN: &[u8] = b"/iroh-db/sync/4";
38
39#[derive(Debug, Clone, Copy, PartialEq, Eq)]
41pub struct SyncReport {
42 received: usize,
43 sent: usize,
44}
45
46impl SyncReport {
47 pub const fn received(self) -> usize {
49 self.received
50 }
51
52 pub const fn sent(self) -> usize {
54 self.sent
55 }
56}
57
58pub struct SyncNode {
60 core: SyncCore,
61 gossip: Gossip,
62 lookup: MemoryLookup,
63 router: Router,
64}
65
66impl SyncNode {
67 pub fn addr(&self) -> EndpointAddr {
69 self.router.endpoint().addr()
70 }
71
72 pub async fn sync_with(&self, remote: EndpointAddr) -> Result<SyncReport, DbError> {
74 self.core.sync_with(remote).await
75 }
76
77 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 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 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 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 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 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 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
250pub 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 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 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 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}