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 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 #[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}