1use std::collections::HashSet;
14use std::sync::atomic::{AtomicUsize, Ordering};
15
16use crate::codec::{decode_tree_node, is_tree_node};
17use crate::crypto::decrypt_chk;
18use crate::hashtree::{HashTree, HashTreeError};
19use crate::store::Store;
20use crate::types::{Cid, Hash};
21
22#[derive(Debug, Clone)]
24pub struct TreeDiff {
25 pub added: Vec<Hash>,
27 pub stats: DiffStats,
29}
30
31impl TreeDiff {
32 pub fn empty() -> Self {
34 Self {
35 added: Vec::new(),
36 stats: DiffStats::default(),
37 }
38 }
39
40 pub fn is_empty(&self) -> bool {
42 self.added.is_empty()
43 }
44
45 pub fn added_count(&self) -> usize {
47 self.added.len()
48 }
49}
50
51#[derive(Debug, Clone, Default)]
53pub struct DiffStats {
54 pub old_tree_nodes: usize,
56 pub new_tree_nodes: usize,
58 pub unchanged_subtrees: usize,
60}
61
62fn decrypt_if_keyed(data: Vec<u8>, key: Option<[u8; 32]>) -> Result<Vec<u8>, HashTreeError> {
63 if let Some(k) = key {
64 decrypt_chk(&data, &k).map_err(|e| HashTreeError::Decryption(e.to_string()))
65 } else {
66 Ok(data)
67 }
68}
69
70pub async fn collect_hashes<S: Store>(
83 tree: &HashTree<S>,
84 root: &Cid,
85 concurrency: usize,
86) -> Result<HashSet<Hash>, HashTreeError> {
87 collect_hashes_with_progress(tree, root, concurrency, None).await
88}
89
90pub async fn collect_hashes_with_progress<S: Store>(
92 tree: &HashTree<S>,
93 root: &Cid,
94 concurrency: usize,
95 progress: Option<&AtomicUsize>,
96) -> Result<HashSet<Hash>, HashTreeError> {
97 use futures::stream::{FuturesUnordered, StreamExt};
98 use std::collections::VecDeque;
99
100 let concurrency = concurrency.max(1);
101 let store = tree.get_store();
102 let mut hashes = HashSet::new();
103 let mut pending: VecDeque<(Hash, Option<[u8; 32]>)> = VecDeque::new();
104 let mut active = FuturesUnordered::new();
105
106 pending.push_back((root.hash, root.key));
108
109 loop {
110 while active.len() < concurrency {
112 if let Some((hash, key)) = pending.pop_front() {
113 if !hashes.insert(hash) {
115 continue;
116 }
117
118 let store = &store;
119 let fut = async move {
120 let data = store
121 .get(&hash)
122 .await
123 .map_err(|e| HashTreeError::Store(e.to_string()))?;
124 Ok::<_, HashTreeError>((key, data))
125 };
126 active.push(fut);
127 } else {
128 break;
129 }
130 }
131
132 if active.is_empty() {
134 break;
135 }
136
137 if let Some(result) = active.next().await {
139 let (key, data) = result?;
140
141 if let Some(counter) = progress {
142 counter.fetch_add(1, Ordering::Relaxed);
143 }
144
145 let data = match data {
146 Some(d) => d,
147 None => continue,
148 };
149
150 let plaintext = decrypt_if_keyed(data, key)?;
152
153 if is_tree_node(&plaintext) {
155 if let Ok(node) = decode_tree_node(&plaintext) {
156 for link in node.links {
157 if !hashes.contains(&link.hash) {
158 pending.push_back((link.hash, link.key));
159 }
160 }
161 }
162 }
163 }
164 }
165
166 Ok(hashes)
167}
168
169pub async fn tree_diff<S: Store>(
189 tree: &HashTree<S>,
190 old_root: Option<&Cid>,
191 new_root: &Cid,
192 concurrency: usize,
193) -> Result<TreeDiff, HashTreeError> {
194 let old_hashes = if let Some(old) = old_root {
196 collect_hashes(tree, old, concurrency).await?
197 } else {
198 HashSet::new()
199 };
200
201 tree_diff_with_old_hashes(tree, &old_hashes, new_root, concurrency).await
202}
203
204pub async fn tree_diff_with_old_hashes<S: Store>(
209 tree: &HashTree<S>,
210 old_hashes: &HashSet<Hash>,
211 new_root: &Cid,
212 concurrency: usize,
213) -> Result<TreeDiff, HashTreeError> {
214 let mut added = Vec::new();
215 let stats = tree_diff_streaming(tree, old_hashes, new_root, concurrency, |hash| {
216 added.push(hash);
217 true
218 })
219 .await?;
220 Ok(TreeDiff { added, stats })
221}
222
223pub async fn tree_diff_streaming<S, F>(
228 tree: &HashTree<S>,
229 old_hashes: &HashSet<Hash>,
230 new_root: &Cid,
231 concurrency: usize,
232 mut callback: F,
233) -> Result<DiffStats, HashTreeError>
234where
235 S: Store,
236 F: FnMut(Hash) -> bool, {
238 use futures::stream::{FuturesUnordered, StreamExt};
239 use std::collections::VecDeque;
240
241 let concurrency = concurrency.max(1);
242 let store = tree.get_store();
243 let mut visited: HashSet<Hash> = HashSet::new();
244 let mut pending: VecDeque<(Hash, Option<[u8; 32]>)> = VecDeque::new();
245 let mut active = FuturesUnordered::new();
246
247 let mut stats = DiffStats {
248 old_tree_nodes: old_hashes.len(),
249 new_tree_nodes: 0,
250 unchanged_subtrees: 0,
251 };
252
253 pending.push_back((new_root.hash, new_root.key));
254
255 loop {
256 while active.len() < concurrency {
257 if let Some((hash, key)) = pending.pop_front() {
258 if !visited.insert(hash) {
259 continue;
260 }
261
262 if old_hashes.contains(&hash) {
263 stats.unchanged_subtrees += 1;
264 continue;
265 }
266
267 stats.new_tree_nodes += 1;
268
269 if !callback(hash) {
271 return Ok(stats);
273 }
274
275 let store = &store;
276 let fut = async move {
277 let data = store
278 .get(&hash)
279 .await
280 .map_err(|e| HashTreeError::Store(e.to_string()))?;
281 Ok::<_, HashTreeError>((key, data))
282 };
283 active.push(fut);
284 } else {
285 break;
286 }
287 }
288
289 if active.is_empty() {
290 break;
291 }
292
293 if let Some(result) = active.next().await {
294 let (key, data) = result?;
295
296 let data = match data {
297 Some(d) => d,
298 None => continue,
299 };
300
301 let plaintext = decrypt_if_keyed(data, key)?;
302
303 if is_tree_node(&plaintext) {
304 if let Ok(node) = decode_tree_node(&plaintext) {
305 for link in node.links {
306 if !visited.contains(&link.hash) {
307 pending.push_back((link.hash, link.key));
308 }
309 }
310 }
311 }
312 }
313 }
314
315 Ok(stats)
316}
317
318#[cfg(test)]
319mod tests {
320 use super::*;
321 use crate::store::MemoryStore;
322 use crate::types::{DirEntry, LinkType};
323 use crate::HashTreeConfig;
324 use std::sync::Arc;
325
326 fn make_tree() -> (Arc<MemoryStore>, HashTree<MemoryStore>) {
327 let store = Arc::new(MemoryStore::new());
328 let tree = HashTree::new(HashTreeConfig::new(store.clone()).public());
329 (store, tree)
330 }
331
332 fn make_encrypted_tree() -> (Arc<MemoryStore>, HashTree<MemoryStore>) {
333 let store = Arc::new(MemoryStore::new());
334 let tree = HashTree::new(HashTreeConfig::new(store.clone()));
335 (store, tree)
336 }
337
338 #[tokio::test]
339 async fn test_diff_identical_trees() {
340 let (_store, tree) = make_tree();
341
342 let file1 = tree.put_blob(b"content1").await.unwrap();
344 let file2 = tree.put_blob(b"content2").await.unwrap();
345 let dir_cid = tree
346 .put_directory(vec![
347 DirEntry::new("a.txt", file1).with_size(8),
348 DirEntry::new("b.txt", file2).with_size(8),
349 ])
350 .await
351 .unwrap();
352
353 let diff = tree_diff(&tree, Some(&dir_cid), &dir_cid, 4).await.unwrap();
355
356 assert!(diff.is_empty(), "identical trees should have empty diff");
357 assert_eq!(diff.added_count(), 0);
358 }
359
360 #[tokio::test]
361 async fn test_diff_single_file_change() {
362 let (_store, tree) = make_tree();
363
364 let file1 = tree.put_blob(b"content1").await.unwrap();
366 let file2 = tree.put_blob(b"content2").await.unwrap();
367 let old_dir = tree
368 .put_directory(vec![
369 DirEntry::new("a.txt", file1).with_size(8),
370 DirEntry::new("b.txt", file2).with_size(8),
371 ])
372 .await
373 .unwrap();
374
375 let file1_new = tree.put_blob(b"content1-modified").await.unwrap();
377 let new_dir = tree
378 .put_directory(vec![
379 DirEntry::new("a.txt", file1_new).with_size(17),
380 DirEntry::new("b.txt", file2).with_size(8), ])
382 .await
383 .unwrap();
384
385 let diff = tree_diff(&tree, Some(&old_dir), &new_dir, 4).await.unwrap();
386
387 assert!(!diff.is_empty());
389 assert_eq!(diff.added_count(), 2); assert!(diff.added.contains(&file1_new));
391 assert!(diff.added.contains(&new_dir.hash));
392 assert!(!diff.added.contains(&file2)); }
394
395 #[tokio::test]
396 async fn test_diff_subtree_unchanged() {
397 let (_store, tree) = make_tree();
398
399 let sub_file = tree.put_blob(b"sub content").await.unwrap();
401 let subdir = tree
402 .put_directory(vec![DirEntry::new("sub.txt", sub_file).with_size(11)])
403 .await
404 .unwrap();
405
406 let file1 = tree.put_blob(b"root file").await.unwrap();
408 let old_root = tree
409 .put_directory(vec![
410 DirEntry::new("file.txt", file1).with_size(9),
411 DirEntry::new("subdir", subdir.hash)
412 .with_size(0)
413 .with_link_type(LinkType::Dir),
414 ])
415 .await
416 .unwrap();
417
418 let file1_new = tree.put_blob(b"root file changed").await.unwrap();
420 let new_root = tree
421 .put_directory(vec![
422 DirEntry::new("file.txt", file1_new).with_size(17),
423 DirEntry::new("subdir", subdir.hash)
424 .with_size(0)
425 .with_link_type(LinkType::Dir),
426 ])
427 .await
428 .unwrap();
429
430 let diff = tree_diff(&tree, Some(&old_root), &new_root, 4)
431 .await
432 .unwrap();
433
434 assert!(diff.added.contains(&new_root.hash));
436 assert!(diff.added.contains(&file1_new));
437 assert!(!diff.added.contains(&subdir.hash)); assert!(!diff.added.contains(&sub_file)); assert!(diff.stats.unchanged_subtrees > 0);
442 }
443
444 #[tokio::test]
445 async fn test_diff_new_directory() {
446 let (_store, tree) = make_tree();
447
448 let file1 = tree.put_blob(b"file1").await.unwrap();
450 let old_root = tree
451 .put_directory(vec![DirEntry::new("a.txt", file1).with_size(5)])
452 .await
453 .unwrap();
454
455 let new_file = tree.put_blob(b"new file").await.unwrap();
457 let new_dir = tree
458 .put_directory(vec![DirEntry::new("inner.txt", new_file).with_size(8)])
459 .await
460 .unwrap();
461 let new_root = tree
462 .put_directory(vec![
463 DirEntry::new("a.txt", file1).with_size(5),
464 DirEntry::new("newdir", new_dir.hash)
465 .with_size(0)
466 .with_link_type(LinkType::Dir),
467 ])
468 .await
469 .unwrap();
470
471 let diff = tree_diff(&tree, Some(&old_root), &new_root, 4)
472 .await
473 .unwrap();
474
475 assert!(diff.added.contains(&new_root.hash));
477 assert!(diff.added.contains(&new_dir.hash));
478 assert!(diff.added.contains(&new_file));
479 assert!(!diff.added.contains(&file1));
481 }
482
483 #[tokio::test]
484 async fn test_diff_empty_old_tree() {
485 let (_store, tree) = make_tree();
486
487 let file1 = tree.put_blob(b"content").await.unwrap();
489 let new_root = tree
490 .put_directory(vec![DirEntry::new("file.txt", file1).with_size(7)])
491 .await
492 .unwrap();
493
494 let diff = tree_diff(&tree, None, &new_root, 4).await.unwrap();
496
497 assert_eq!(diff.added_count(), 2); assert!(diff.added.contains(&new_root.hash));
500 assert!(diff.added.contains(&file1));
501 }
502
503 #[tokio::test]
504 async fn test_diff_encrypted_trees() {
505 let (_store, tree) = make_encrypted_tree();
506
507 let (file1_cid, _) = tree.put(b"content1").await.unwrap();
509 let (file2_cid, _) = tree.put(b"content2").await.unwrap();
510 let old_dir = tree
511 .put_directory(vec![
512 DirEntry::from_cid("a.txt", &file1_cid).with_size(8),
513 DirEntry::from_cid("b.txt", &file2_cid).with_size(8),
514 ])
515 .await
516 .unwrap();
517
518 let (file1_new_cid, _) = tree.put(b"content1-modified").await.unwrap();
520 let new_dir = tree
521 .put_directory(vec![
522 DirEntry::from_cid("a.txt", &file1_new_cid).with_size(17),
523 DirEntry::from_cid("b.txt", &file2_cid).with_size(8),
524 ])
525 .await
526 .unwrap();
527
528 let diff = tree_diff(&tree, Some(&old_dir), &new_dir, 4).await.unwrap();
529
530 assert!(!diff.is_empty());
532 assert!(diff.added.contains(&file1_new_cid.hash));
533 assert!(!diff.added.contains(&file2_cid.hash)); }
535
536 #[tokio::test]
537 async fn test_diff_streaming_early_stop() {
538 let (_store, tree) = make_tree();
539
540 let file1 = tree.put_blob(b"f1").await.unwrap();
541 let file2 = tree.put_blob(b"f2").await.unwrap();
542 let file3 = tree.put_blob(b"f3").await.unwrap();
543 let new_root = tree
544 .put_directory(vec![
545 DirEntry::new("a.txt", file1).with_size(2),
546 DirEntry::new("b.txt", file2).with_size(2),
547 DirEntry::new("c.txt", file3).with_size(2),
548 ])
549 .await
550 .unwrap();
551
552 let old_hashes = HashSet::new(); let mut count = 0;
555 let stats = tree_diff_streaming(&tree, &old_hashes, &new_root, 1, |_hash| {
556 count += 1;
557 count < 2 })
559 .await
560 .unwrap();
561
562 assert_eq!(count, 2, "should stop after exactly two callbacks");
563 assert_eq!(stats.new_tree_nodes, 2);
564 }
565
566 #[tokio::test]
567 async fn test_diff_large_tree_structure() {
568 let (_store, tree) = make_tree();
569
570 let mut entries = Vec::new();
572 let mut old_hashes_vec = Vec::new();
573
574 for i in 0..100 {
575 let data = format!("content {}", i);
576 let hash = tree.put_blob(data.as_bytes()).await.unwrap();
577 entries
578 .push(DirEntry::new(format!("file{}.txt", i), hash).with_size(data.len() as u64));
579 old_hashes_vec.push(hash);
580 }
581
582 let old_root = tree.put_directory(entries.clone()).await.unwrap();
583 old_hashes_vec.push(old_root.hash);
584
585 for (i, entry) in entries.iter_mut().enumerate().take(5) {
587 let data = format!("modified content {}", i);
588 let hash = tree.put_blob(data.as_bytes()).await.unwrap();
589 *entry = DirEntry::new(format!("file{}.txt", i), hash).with_size(data.len() as u64);
590 }
591
592 let new_root = tree.put_directory(entries).await.unwrap();
593
594 let diff = tree_diff(&tree, Some(&old_root), &new_root, 8)
595 .await
596 .unwrap();
597
598 assert_eq!(diff.added_count(), 6);
600 assert!(diff.added.contains(&new_root.hash));
601
602 assert!(diff.stats.unchanged_subtrees >= 95);
604 }
605}