Skip to main content

deaddrop_net/sync/
mod.rs

1use crate::routing::strategy_from_kind;
2use deaddrop_core::crypto::PrivateIdentity;
3use deaddrop_core::event::{Bus, Event, Metrics};
4use deaddrop_core::identity::ContactCard;
5use deaddrop_core::protocol::{
6    BloomFilter, ChunkWant, Message, handshake_initiator, handshake_responder, read_msg,
7    verify_envelope, write_msg,
8};
9use deaddrop_core::store::{Store, unix_now};
10use deaddrop_core::{DdError, NodeMode, ObjectId, Ownership, PeerId, Result, RoutingPolicyKind};
11use std::collections::{HashSet, VecDeque};
12use std::sync::atomic::Ordering;
13use tokio::io::{AsyncRead, AsyncWrite};
14
15#[derive(Debug, Default, Clone)]
16pub struct SyncStats {
17    pub peer: Option<PeerId>,
18    pub envelopes: u32,
19    pub chunks: u32,
20}
21
22pub struct SyncOpts<'a> {
23    pub store: &'a Store,
24    pub identity: &'a PrivateIdentity,
25    pub mode: NodeMode,
26    pub default_strategy: RoutingPolicyKind,
27    pub metrics: Option<&'a Metrics>,
28    pub bus: Option<&'a Bus>,
29    pub initiator: bool,
30}
31
32pub async fn sync_session<S>(opts: SyncOpts<'_>, mut stream: S) -> Result<SyncStats>
33where
34    S: AsyncRead + AsyncWrite + Unpin,
35{
36    let session = if opts.initiator {
37        handshake_initiator(opts.identity, &mut stream).await?
38    } else {
39        handshake_responder(opts.identity, &mut stream).await?
40    };
41    let (mut reader, mut writer) = tokio::io::split(stream);
42    if let Some(b) = opts.bus {
43        b.emit(Event::PeerAuthenticated {
44            peer: session.peer.to_string(),
45        });
46    }
47    let _ = opts.store.put_contact(
48        ContactCard::from_public(None, session.peer_identity.clone())?,
49        deaddrop_core::TrustState::Observed,
50    );
51    let now = unix_now();
52    let local_ids = opts.store.inventory(now)?;
53    if session.caps.prefer_inventory() == "bloom-v1" && local_ids.len() > 64 {
54        let bloom = BloomFilter::from_ids(&local_ids, 10);
55        write_msg(
56            &mut writer,
57            &Message::InventoryBloom {
58                n: bloom.n,
59                k: bloom.k,
60                bits: bloom.bits,
61            },
62        )
63        .await?;
64    } else {
65        write_msg(
66            &mut writer,
67            &Message::InventorySorted {
68                object_ids: local_ids.clone(),
69            },
70        )
71        .await?;
72    }
73    let inv = read_msg(&mut reader).await?;
74    let mut want_ids: Vec<ObjectId> = match inv {
75        Message::InventorySorted { object_ids } => object_ids
76            .into_iter()
77            .filter(|id| !local_ids.contains(id))
78            .collect(),
79        Message::InventoryBloom { n: _, k, bits } => {
80            let bloom = BloomFilter { n: 0, k, bits };
81            local_ids
82                .iter()
83                .copied()
84                .filter(|id| !bloom.may_contain(id))
85                .collect()
86        }
87        other => {
88            return Err(DdError::invalid_frame(format!(
89                "expected inventory, got {}",
90                other.name()
91            )));
92        }
93    };
94    let mut chunk_wants = Vec::new();
95    for id in &local_ids {
96        if let Ok(mask) = opts.store.present_mask(id) {
97            let missing: Vec<u32> = mask
98                .iter()
99                .enumerate()
100                .filter(|(_, p)| !**p)
101                .map(|(i, _)| i as u32)
102                .collect();
103            if !missing.is_empty() {
104                chunk_wants.push(ChunkWant {
105                    object_id: *id,
106                    indices: missing,
107                });
108            }
109        }
110    }
111    write_msg(
112        &mut writer,
113        &Message::Want {
114            object_ids: want_ids.clone(),
115            chunks: chunk_wants.clone(),
116        },
117    )
118    .await?;
119    let their_want = match read_msg(&mut reader).await? {
120        Message::Want { object_ids, chunks } => (object_ids, chunks),
121        other => {
122            return Err(DdError::invalid_frame(format!(
123                "expected want, got {}",
124                other.name()
125            )));
126        }
127    };
128
129    let mut to_send_env: VecDeque<ObjectId> = VecDeque::new();
130    let mut to_send_chunks: VecDeque<(ObjectId, u32)> = VecDeque::new();
131    let strat = strategy_from_kind(opts.default_strategy);
132    for id in &their_want.0 {
133        if let Ok(Some(env)) = opts.store.get_envelope(id) {
134            let copies = opts.store.replication(id).unwrap_or(1);
135            let d = strat.decide(
136                &env,
137                session.peer,
138                opts.identity.peer_id,
139                now,
140                opts.store,
141                opts.mode,
142                session.caps.peer_relay,
143                copies,
144            )?;
145            if let Some(m) = opts.metrics {
146                m.route_decisions.fetch_add(1, Ordering::Relaxed);
147            }
148            if d.forward {
149                to_send_env.push_back(*id);
150                if let Ok(mask) = opts.store.present_mask(id) {
151                    for (i, present) in mask.iter().enumerate() {
152                        if *present {
153                            to_send_chunks.push_back((*id, i as u32));
154                        }
155                    }
156                }
157            }
158        }
159    }
160    for w in their_want.1 {
161        for idx in w.indices {
162            to_send_chunks.push_back((w.object_id, idx));
163        }
164    }
165
166    let mut needed: HashSet<ObjectId> = want_ids.drain(..).collect();
167    let mut we_done = false;
168    let mut they_done = false;
169    let mut stats = SyncStats {
170        peer: Some(session.peer),
171        ..Default::default()
172    };
173
174    loop {
175        if to_send_env.is_empty() && to_send_chunks.is_empty() && !we_done {
176            write_msg(&mut writer, &Message::Done).await?;
177            we_done = true;
178        }
179        if we_done && they_done {
180            break;
181        }
182        tokio::select! {
183            r = send_one(opts.store, &mut writer, &mut to_send_env, &mut to_send_chunks, &mut stats, now), if !to_send_env.is_empty() || !to_send_chunks.is_empty() => {
184                r?;
185            }
186            incoming = read_msg(&mut reader) => {
187                match incoming? {
188                    Message::Envelope { envelope, manifest } => {
189                        if verify_envelope(&envelope, now).is_ok() {
190                            let oid = envelope.object_id;
191                            let own = if envelope.destination.includes(&opts.identity.peer_id) {
192                                Ownership::Incoming
193                            } else {
194                                Ownership::Relay
195                            };
196                            let _ = opts.store.put_object(&envelope, &manifest, &[], own, now);
197                            needed.remove(&oid);
198                            stats.envelopes += 1;
199                        }
200                    }
201                    Message::ChunkData { object_id, index, data, .. } => {
202                        if opts.store.put_chunk(object_id, index, &data).is_ok() {
203                            stats.chunks += 1;
204                            if let Some(m) = opts.metrics {
205                                m.chunks_transferred.fetch_add(1, Ordering::Relaxed);
206                                m.bytes_transferred.fetch_add(data.len() as u64, Ordering::Relaxed);
207                            }
208                        }
209                    }
210                    Message::Done => they_done = true,
211                    Message::Error { code, message } => {
212                        return Err(DdError::invalid_frame(format!("{code}: {message}")));
213                    }
214                    _ => {}
215                }
216            }
217        }
218    }
219    let ok = true;
220    let _ = opts.store.record_encounter(session.peer, now, 0, 0, ok);
221    let _ = strat;
222    Ok(stats)
223}
224
225async fn send_one<S: AsyncWrite + Unpin>(
226    store: &Store,
227    stream: &mut S,
228    envs: &mut VecDeque<ObjectId>,
229    chunks: &mut VecDeque<(ObjectId, u32)>,
230    stats: &mut SyncStats,
231    now: u64,
232) -> Result<()> {
233    if let Some(id) = envs.pop_front() {
234        if let (Ok(Some(env)), Ok(Some(man))) = (store.get_envelope(&id), store.get_manifest(&id)) {
235            if verify_envelope(&env, now).is_ok() {
236                write_msg(
237                    stream,
238                    &Message::Envelope {
239                        envelope: env,
240                        manifest: man,
241                    },
242                )
243                .await?;
244                stats.envelopes += 1;
245            }
246        }
247        return Ok(());
248    }
249    if let Some((id, idx)) = chunks.pop_front() {
250        if let Ok(data) = store.load_chunk(&id, idx) {
251            if let Ok(refer) = {
252                // chunk id from store via load
253                use deaddrop_core::HashAlgorithm;
254                use deaddrop_core::crypto::{CryptoProvider, DefaultProvider};
255                let d = DefaultProvider.hash(HashAlgorithm::Blake3, &data);
256                Ok::<_, DdError>(deaddrop_core::ChunkId::blake3(d.0))
257            } {
258                write_msg(
259                    stream,
260                    &Message::ChunkData {
261                        object_id: id,
262                        chunk_id: refer,
263                        index: idx,
264                        data,
265                    },
266                )
267                .await?;
268                stats.chunks += 1;
269            }
270        }
271    }
272    Ok(())
273}
274
275#[cfg(test)]
276mod tests {
277    use super::*;
278    use deaddrop_core::identity::ContactCard;
279    use deaddrop_core::protocol::{CreateDrop, build_drop, open_drop};
280    use deaddrop_core::store::Store;
281    use deaddrop_core::{Destination, Ownership, Priority, RoutingPolicy};
282
283    fn tmp(name: &str) -> std::path::PathBuf {
284        let p = std::env::temp_dir().join(format!("ddp2-{name}-{}", std::process::id()));
285        let _ = std::fs::remove_dir_all(&p);
286        std::fs::create_dir_all(&p).unwrap();
287        p
288    }
289
290    #[tokio::test]
291    async fn multi_hop_encrypted_carry() {
292        let alice_id = PrivateIdentity::generate();
293        let bob_id = PrivateIdentity::generate();
294        let charlie_id = PrivateIdentity::generate();
295        let sa = Store::open(tmp("a"), Default::default()).unwrap();
296        let sb = Store::open(tmp("b"), Default::default()).unwrap();
297        let sc = Store::open(tmp("c"), Default::default()).unwrap();
298        sa.put_contact(
299            ContactCard::from_public(None, charlie_id.public.clone()).unwrap(),
300            deaddrop_core::TrustState::Known,
301        )
302        .unwrap();
303        let built = build_drop(CreateDrop {
304            author: &alice_id,
305            recipients: vec![(charlie_id.peer_id, charlie_id.public.clone())],
306            destination: Destination::One {
307                peer: charlie_id.peer_id,
308            },
309            plaintext: b"phase-two-secret".to_vec(),
310            now: unix_now(),
311            ttl_secs: Some(3600),
312            priority: Priority::Normal,
313            routing: RoutingPolicy {
314                kind: deaddrop_core::RoutingPolicyKind::Epidemic,
315                replication_budget: 8,
316                trusted_only: false,
317            },
318            application: "dd.file".into(),
319            topic: None,
320            chunking: deaddrop_core::chunk::default_fixed(),
321            compress: false,
322            hop_limit: 8,
323            public: false,
324            seal_until: None,
325            seal_quorum: None,
326            erasure: None,
327        })
328        .unwrap();
329        let chunks: Vec<_> = built
330            .chunks
331            .iter()
332            .cloned()
333            .enumerate()
334            .map(|(i, d)| (i as u32, d))
335            .collect();
336        sa.put_object(
337            &built.envelope,
338            &built.manifest,
339            &chunks,
340            Ownership::Local,
341            unix_now(),
342        )
343        .unwrap();
344
345        let (ab, ba) = tokio::io::duplex(1 << 20);
346        let oa = SyncOpts {
347            store: &sa,
348            identity: &alice_id,
349            mode: NodeMode::Balanced,
350            default_strategy: deaddrop_core::RoutingPolicyKind::Epidemic,
351            metrics: None,
352            bus: None,
353            initiator: true,
354        };
355        let ob = SyncOpts {
356            store: &sb,
357            identity: &bob_id,
358            mode: NodeMode::Balanced,
359            default_strategy: deaddrop_core::RoutingPolicyKind::Epidemic,
360            metrics: None,
361            bus: None,
362            initiator: false,
363        };
364        let (r1, r2) = tokio::join!(sync_session(oa, ab), sync_session(ob, ba));
365        r1.unwrap();
366        r2.unwrap();
367        assert!(
368            sb.get_envelope(&built.envelope.object_id)
369                .unwrap()
370                .is_some()
371        );
372        let man = sb.get_manifest(&built.envelope.object_id).unwrap().unwrap();
373        let bchunks = sb.load_chunks(&built.envelope.object_id).unwrap();
374        assert!(open_drop(&bob_id, &built.envelope, &man, &bchunks).is_err());
375
376        let (bc, cb) = tokio::io::duplex(1 << 20);
377        let ob2 = SyncOpts {
378            store: &sb,
379            identity: &bob_id,
380            mode: NodeMode::Balanced,
381            default_strategy: deaddrop_core::RoutingPolicyKind::Epidemic,
382            metrics: None,
383            bus: None,
384            initiator: true,
385        };
386        let oc = SyncOpts {
387            store: &sc,
388            identity: &charlie_id,
389            mode: NodeMode::Balanced,
390            default_strategy: deaddrop_core::RoutingPolicyKind::Epidemic,
391            metrics: None,
392            bus: None,
393            initiator: false,
394        };
395        let (r1, r2) = tokio::join!(sync_session(ob2, bc), sync_session(oc, cb));
396        r1.unwrap();
397        r2.unwrap();
398        let man = sc.get_manifest(&built.envelope.object_id).unwrap().unwrap();
399        let cchunks = sc.load_chunks(&built.envelope.object_id).unwrap();
400        let pt = open_drop(
401            &charlie_id,
402            &sc.get_envelope(&built.envelope.object_id).unwrap().unwrap(),
403            &man,
404            &cchunks,
405        )
406        .unwrap();
407        assert_eq!(pt, b"phase-two-secret");
408    }
409
410    /// Two machines on a LAN: TCP listen + one-shot sync, like home PC ↔ laptop.
411    #[tokio::test]
412    async fn two_computers_over_tcp() {
413        use crate::node::Node;
414        use deaddrop_core::config::Config;
415        use tokio::net::TcpListener;
416
417        let home = tmp("home-pc");
418        let laptop = tmp("laptop");
419        let sh = Store::open(&home, Default::default()).unwrap();
420        let sl = Store::open(&laptop, Default::default()).unwrap();
421        let home_id = sh.init_identity(true).unwrap();
422        let laptop_id = sl.init_identity(true).unwrap();
423        sh.put_contact(
424            ContactCard::from_public(Some("laptop".into()), laptop_id.public.clone()).unwrap(),
425            deaddrop_core::TrustState::Known,
426        )
427        .unwrap();
428        sl.put_contact(
429            ContactCard::from_public(Some("home".into()), home_id.public.clone()).unwrap(),
430            deaddrop_core::TrustState::Known,
431        )
432        .unwrap();
433
434        let home_node = Node::open(&home, Config::default()).unwrap();
435        let laptop_node = Node::open(&laptop, Config::default()).unwrap();
436        home_node
437            .send_payload(
438                "laptop",
439                b"family-photo".to_vec(),
440                false,
441                None,
442                Priority::Normal,
443                "dd.file",
444                None,
445                false,
446            )
447            .unwrap();
448
449        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
450        let addr = listener.local_addr().unwrap();
451        let store = laptop_node.store.clone();
452        let secrets = (
453            laptop_node.identity.ed25519_bytes(),
454            laptop_node.identity.x25519_bytes(),
455        );
456        tokio::spawn(async move {
457            let identity = PrivateIdentity::from_secrets(secrets.0, secrets.1);
458            let (s, _) = listener.accept().await.unwrap();
459            let opts = SyncOpts {
460                store: store.as_ref(),
461                identity: &identity,
462                mode: NodeMode::Balanced,
463                default_strategy: deaddrop_core::RoutingPolicyKind::Epidemic,
464                metrics: None,
465                bus: None,
466                initiator: false,
467            };
468            let _ = sync_session(opts, s).await;
469        });
470        tokio::time::sleep(std::time::Duration::from_millis(50)).await;
471        home_node.sync_once(addr).await.unwrap();
472        tokio::time::sleep(std::time::Duration::from_millis(80)).await;
473        let inbox = laptop_node.inbox().unwrap();
474        assert!(
475            inbox.iter().any(|i| i.body == b"family-photo"),
476            "laptop should decrypt the Drop sent from the home PC"
477        );
478    }
479}