1use anyhow::{Context, Result};
2use futures::executor::block_on as sync_block_on;
3use hashtree_core::store::Store;
4use hashtree_core::{to_hex, types::Hash, Cid, HashTree, HashTreeConfig, HashTreeError, LinkType};
5use serde::de::{self, IgnoredAny, MapAccess, SeqAccess, Visitor};
6use serde::{Deserialize, Serialize};
7use std::collections::{BTreeMap, HashSet};
8use std::fs::{File, OpenOptions};
9#[cfg(unix)]
10use std::os::fd::AsRawFd;
11use std::path::{Path, PathBuf};
12use std::sync::Arc;
13use std::time::{SystemTime, UNIX_EPOCH};
14
15use super::quota::{CacheQuotaAdmission, CacheWritePermit};
16use super::{BlobMetadata, HashtreeStore, PRIORITY_FOLLOWED, PRIORITY_OWN};
17
18const MAX_PINNED_TREE_NODES: usize = 10_000_000;
19const MAX_UNBOUNDED_PINNED_TREE_BYTES: u64 = 1 << 50;
20const ORPHAN_SCAN_PAGE_SIZE: usize = 4_096;
21const RETENTION_ROOTS_LOCK_FILE: &str = ".retention-roots.lock";
22const MAX_PROFILE_REPAIR_RETENTION_LEASE_BYTES: u64 = 64 * 1024;
23const MAX_PROFILE_REPAIR_RETENTION_ROOTS: usize = 64;
24
25pub const PROFILE_REPAIR_RETENTION_LEASE_FORMAT: &str =
26 "iris-social/bulk-profile-index-repair-retention@1";
27pub const PROFILE_REPAIR_RETENTION_LEASE_RELATIVE_PATH: &str =
28 "nostr-index/bulk-projection-v2/profile-repair-v1/retention.json";
29
30#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
31#[serde(deny_unknown_fields)]
32pub struct ProfileRepairRetentionLease {
33 pub format: String,
34 pub authority_sha256: String,
35 pub roots: BTreeMap<String, String>,
36}
37
38impl ProfileRepairRetentionLease {
39 pub fn validate(&self) -> Result<()> {
40 if self.format != PROFILE_REPAIR_RETENTION_LEASE_FORMAT {
41 anyhow::bail!(
42 "profile repair retention lease has unsupported format {}",
43 self.format
44 );
45 }
46 if self.authority_sha256.len() != 64
47 || !self
48 .authority_sha256
49 .bytes()
50 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
51 {
52 anyhow::bail!(
53 "profile repair retention lease authority must be 64 lowercase hexadecimal characters"
54 );
55 }
56 if self.roots.is_empty() || self.roots.len() > MAX_PROFILE_REPAIR_RETENTION_ROOTS {
57 anyhow::bail!(
58 "profile repair retention lease must contain between 1 and {} roots",
59 MAX_PROFILE_REPAIR_RETENTION_ROOTS
60 );
61 }
62 for (label, encoded) in &self.roots {
63 if label.is_empty()
64 || label.len() > 128
65 || !label
66 .bytes()
67 .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.'))
68 {
69 anyhow::bail!("profile repair retention lease has invalid root label {label:?}");
70 }
71 let cid = Cid::parse(encoded)
72 .map_err(|error| anyhow::anyhow!("invalid retained root {label}: {error}"))?;
73 if cid.to_string() != *encoded {
74 anyhow::bail!(
75 "profile repair retention lease root {label} is not canonical CID text"
76 );
77 }
78 }
79 Ok(())
80 }
81
82 pub fn canonical_bytes(&self) -> Result<Vec<u8>> {
83 self.validate()?;
84 let mut bytes =
85 serde_json::to_vec(self).map_err(|error| anyhow::anyhow!("encode lease: {error}"))?;
86 bytes.push(b'\n');
87 Ok(bytes)
88 }
89
90 pub fn sha256(&self) -> Result<String> {
91 Ok(to_hex(&hashtree_core::sha256(&self.canonical_bytes()?)))
92 }
93
94 fn root_cids(&self) -> Result<Vec<Cid>> {
95 self.validate()?;
96 self.roots
97 .iter()
98 .map(|(label, encoded)| {
99 Cid::parse(encoded)
100 .map_err(|error| anyhow::anyhow!("invalid retained root {label}: {error}"))
101 })
102 .collect()
103 }
104}
105
106#[derive(Debug, Clone, Copy)]
107enum RetentionRootsLockMode {
108 Shared,
109 Exclusive,
110}
111
112pub struct ProfileRepairRetentionPublicationGuard {
113 file: File,
114}
115
116impl Drop for ProfileRepairRetentionPublicationGuard {
117 fn drop(&mut self) {
118 #[cfg(unix)]
119 unsafe {
120 libc::flock(self.file.as_raw_fd(), libc::LOCK_UN);
121 }
122 }
123}
124
125pub fn acquire_existing_profile_repair_retention_guard(
133 base_path: &Path,
134) -> Result<ProfileRepairRetentionPublicationGuard> {
135 let path = base_path.join(RETENTION_ROOTS_LOCK_FILE);
136 let before = std::fs::symlink_metadata(&path)
137 .with_context(|| format!("inspect existing retention-roots lock {}", path.display()))?;
138 if before.file_type().is_symlink() || !before.file_type().is_file() {
139 anyhow::bail!(
140 "existing retention-roots lock is not a direct regular file: {}",
141 path.display()
142 );
143 }
144 let file = OpenOptions::new()
145 .read(true)
146 .write(true)
147 .open(&path)
148 .with_context(|| format!("open existing retention-roots lock {}", path.display()))?;
149 #[cfg(unix)]
150 {
151 let result = unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX) };
152 if result != 0 {
153 return Err(std::io::Error::last_os_error()).with_context(|| {
154 format!("acquire existing retention-roots lock {}", path.display())
155 });
156 }
157 use std::os::unix::fs::MetadataExt;
158 let opened = file
159 .metadata()
160 .with_context(|| format!("inspect opened retention-roots lock {}", path.display()))?;
161 let current = std::fs::symlink_metadata(&path)
162 .with_context(|| format!("reinspect retention-roots lock {}", path.display()))?;
163 if current.file_type().is_symlink()
164 || !current.file_type().is_file()
165 || opened.dev() != before.dev()
166 || opened.ino() != before.ino()
167 || current.dev() != before.dev()
168 || current.ino() != before.ino()
169 {
170 anyhow::bail!(
171 "retention-roots lock identity changed while it was acquired: {}",
172 path.display()
173 );
174 }
175 }
176 #[cfg(not(unix))]
177 {
178 let _ = before;
179 anyhow::bail!("retention-root transactions require an operating-system advisory file lock");
180 }
181 Ok(ProfileRepairRetentionPublicationGuard { file })
182}
183
184pub(super) struct ActiveRetentionProtection {
185 _guard: ProfileRepairRetentionPublicationGuard,
186 _profile_roots: Option<crate::socialgraph::ProfileRootSnapshotGuard>,
187 hashes: HashSet<Hash>,
188}
189
190impl ActiveRetentionProtection {
191 pub(super) fn contains(&self, hash: &Hash) -> bool {
192 self.hashes.contains(hash)
193 }
194
195 fn hashes(&self) -> &HashSet<Hash> {
196 &self.hashes
197 }
198}
199
200#[derive(Debug, Clone, Copy, PartialEq, Eq)]
201struct OrphanCleanupProgress {
202 freed_bytes: u64,
203 scanned: usize,
204 sweep_complete: bool,
205}
206
207#[derive(Debug, Clone, Copy)]
209pub struct TreeIndexLimits {
210 pub max_nodes: usize,
211 pub max_bytes: u64,
212}
213
214#[derive(Debug, Clone, Copy, PartialEq, Eq)]
216pub struct PinTreeResult {
217 pub indexed_hashes: usize,
218 pub total_size: u64,
219 pub already_pinned: bool,
220}
221
222#[derive(Debug, Clone, Copy, PartialEq, Eq)]
223pub struct RootRetentionReport {
224 pub total_hashes: usize,
225 pub reachable_hashes: usize,
226 pub pinned_hashes: usize,
227 pub candidate_hashes: usize,
228 pub deleted_hashes: usize,
229 pub logical_bytes_before: u64,
230 pub logical_bytes_after: u64,
231}
232
233#[derive(Debug, thiserror::Error)]
234pub enum PinTreeError {
235 #[error("root blob {hash} is missing")]
236 MissingRoot { hash: String },
237 #[error("descendant blob {hash} is missing")]
238 MissingDescendant { hash: String },
239 #[error("invalid DAG node {hash}: {message}")]
240 InvalidDag { hash: String, message: String },
241 #[error("DAG exceeds the {max_nodes} node limit")]
242 NodeLimitExceeded { max_nodes: usize },
243 #[error("DAG exceeds the {max_bytes} byte limit")]
244 ByteLimitExceeded { max_bytes: u64 },
245 #[error("storage error: {0}")]
246 Storage(String),
247}
248
249struct TreeIndexPlan {
250 tracked_hashes: HashSet<Hash>,
251 total_size: u64,
252}
253
254#[derive(Debug, Clone, Serialize)]
256pub struct TreeMeta {
257 pub owner: String,
259 pub name: Option<String>,
261 pub synced_at: u64,
263 pub total_size: u64,
265 pub priority: u8,
267}
268
269impl<'de> Deserialize<'de> for TreeMeta {
270 fn deserialize<D>(deserializer: D) -> std::result::Result<Self, D::Error>
271 where
272 D: serde::Deserializer<'de>,
273 {
274 const FIELDS: &[&str] = &[
275 "owner",
276 "name",
277 "synced_at",
278 "last_accessed_at",
279 "total_size",
280 "priority",
281 ];
282
283 struct TreeMetaVisitor;
284
285 impl<'de> Visitor<'de> for TreeMetaVisitor {
286 type Value = TreeMeta;
287
288 fn expecting(&self, formatter: &mut std::fmt::Formatter) -> std::fmt::Result {
289 formatter.write_str("TreeMeta as current or legacy metadata")
290 }
291
292 fn visit_seq<A>(self, mut seq: A) -> std::result::Result<Self::Value, A::Error>
293 where
294 A: SeqAccess<'de>,
295 {
296 let has_accidental_access_field = matches!(seq.size_hint(), Some(6));
297 let owner = seq
298 .next_element()?
299 .ok_or_else(|| de::Error::invalid_length(0, &self))?;
300 let name = seq
301 .next_element()?
302 .ok_or_else(|| de::Error::invalid_length(1, &self))?;
303 let synced_at = seq
304 .next_element()?
305 .ok_or_else(|| de::Error::invalid_length(2, &self))?;
306
307 if has_accidental_access_field {
308 let _: IgnoredAny = seq
309 .next_element()?
310 .ok_or_else(|| de::Error::invalid_length(3, &self))?;
311 }
312
313 let total_size = seq
314 .next_element()?
315 .ok_or_else(|| de::Error::invalid_length(3, &self))?;
316 let priority = seq
317 .next_element()?
318 .ok_or_else(|| de::Error::invalid_length(4, &self))?;
319
320 Ok(TreeMeta {
321 owner,
322 name,
323 synced_at,
324 total_size,
325 priority,
326 })
327 }
328
329 fn visit_map<A>(self, mut map: A) -> std::result::Result<Self::Value, A::Error>
330 where
331 A: MapAccess<'de>,
332 {
333 let mut owner = None;
334 let mut name = None;
335 let mut synced_at = None;
336 let mut total_size = None;
337 let mut priority = None;
338
339 while let Some(key) = map.next_key::<String>()? {
340 match key.as_str() {
341 "owner" => owner = Some(map.next_value()?),
342 "name" => name = Some(map.next_value()?),
343 "synced_at" => synced_at = Some(map.next_value()?),
344 "last_accessed_at" => {
345 let _: IgnoredAny = map.next_value()?;
346 }
347 "total_size" => total_size = Some(map.next_value()?),
348 "priority" => priority = Some(map.next_value()?),
349 _ => {
350 let _: IgnoredAny = map.next_value()?;
351 }
352 }
353 }
354
355 Ok(TreeMeta {
356 owner: owner.ok_or_else(|| de::Error::missing_field("owner"))?,
357 name: name.unwrap_or(None),
358 synced_at: synced_at.ok_or_else(|| de::Error::missing_field("synced_at"))?,
359 total_size: total_size.ok_or_else(|| de::Error::missing_field("total_size"))?,
360 priority: priority.ok_or_else(|| de::Error::missing_field("priority"))?,
361 })
362 }
363 }
364
365 deserializer.deserialize_struct("TreeMeta", FIELDS, TreeMetaVisitor)
366 }
367}
368
369#[derive(Debug)]
370pub struct StorageStats {
371 pub total_dags: usize,
372 pub pinned_dags: usize,
373 pub total_bytes: u64,
374}
375
376#[derive(Debug, Clone)]
378pub struct StorageByPriority {
379 pub own: u64,
381 pub followed: u64,
383 pub other: u64,
385}
386
387#[derive(Debug, Clone)]
388pub struct PinnedItem {
389 pub cid: String,
390 pub name: String,
391 pub is_directory: bool,
392 pub size_bytes: u64,
393}
394
395#[derive(Debug, Clone)]
396pub struct OwnedBlobStats {
397 pub owner: [u8; 32],
398 pub count: usize,
399 pub total_bytes: u64,
400}
401
402fn pinned_item_name(hash: &Hash, meta: Option<&TreeMeta>) -> String {
403 let Some(meta) = meta else {
404 return to_hex(hash);
405 };
406
407 match (meta.owner.as_str(), meta.name.as_deref()) {
408 ("pinned", Some(name)) => name.to_string(),
409 ("", Some(name)) => name.to_string(),
410 (owner, Some(name)) if !owner.is_empty() => format!("{owner}/{name}"),
411 (owner, None) if !owner.is_empty() && owner != "pinned" => owner.to_string(),
412 _ => to_hex(hash),
413 }
414}
415
416fn unix_timestamp_now() -> u64 {
417 SystemTime::now()
418 .duration_since(UNIX_EPOCH)
419 .unwrap_or_default()
420 .as_secs()
421}
422
423impl HashtreeStore {
424 fn acquire_retention_roots_lock(
425 &self,
426 mode: RetentionRootsLockMode,
427 ) -> Result<ProfileRepairRetentionPublicationGuard> {
428 let path = self.base_path().join(RETENTION_ROOTS_LOCK_FILE);
429 let file = OpenOptions::new()
430 .read(true)
431 .write(true)
432 .create(true)
433 .truncate(false)
434 .open(&path)
435 .with_context(|| format!("open retention-roots lock {}", path.display()))?;
436 #[cfg(unix)]
437 {
438 let operation = match mode {
439 RetentionRootsLockMode::Shared => libc::LOCK_SH,
440 RetentionRootsLockMode::Exclusive => libc::LOCK_EX,
441 };
442 let result = unsafe { libc::flock(file.as_raw_fd(), operation) };
443 if result != 0 {
444 return Err(std::io::Error::last_os_error())
445 .with_context(|| format!("acquire retention-roots lock {}", path.display()));
446 }
447 }
448 #[cfg(not(unix))]
449 {
450 let _ = mode;
451 anyhow::bail!(
452 "retention-root transactions require an operating-system advisory file lock"
453 );
454 }
455 Ok(ProfileRepairRetentionPublicationGuard { file })
456 }
457
458 pub fn acquire_profile_repair_retention_publication_guard(
462 &self,
463 ) -> Result<ProfileRepairRetentionPublicationGuard> {
464 self.acquire_retention_roots_lock(RetentionRootsLockMode::Exclusive)
465 }
466
467 pub fn profile_repair_retention_lease_path(&self) -> PathBuf {
468 self.base_path()
469 .join(PROFILE_REPAIR_RETENTION_LEASE_RELATIVE_PATH)
470 }
471
472 pub fn validate_profile_repair_retention_lease(
473 &self,
474 expected_sha256: &str,
475 ) -> Result<ProfileRepairRetentionLease> {
476 let path = self.profile_repair_retention_lease_path();
477 let metadata = std::fs::symlink_metadata(&path)
478 .with_context(|| format!("inspect {}", path.display()))?;
479 if metadata.file_type().is_symlink() || !metadata.file_type().is_file() {
480 anyhow::bail!(
481 "profile repair retention lease is not a direct regular file: {}",
482 path.display()
483 );
484 }
485 if metadata.len() > MAX_PROFILE_REPAIR_RETENTION_LEASE_BYTES {
486 anyhow::bail!(
487 "profile repair retention lease exceeds {} bytes: {}",
488 MAX_PROFILE_REPAIR_RETENTION_LEASE_BYTES,
489 path.display()
490 );
491 }
492 let bytes = std::fs::read(&path).with_context(|| format!("read {}", path.display()))?;
493 let actual_sha256 = to_hex(&hashtree_core::sha256(&bytes));
494 if actual_sha256 != expected_sha256 {
495 anyhow::bail!(
496 "profile repair retention lease SHA-256 mismatch: expected {expected_sha256}, found {actual_sha256}"
497 );
498 }
499 let lease: ProfileRepairRetentionLease =
500 serde_json::from_slice(&bytes).with_context(|| format!("decode {}", path.display()))?;
501 if lease.canonical_bytes()? != bytes {
502 anyhow::bail!(
503 "profile repair retention lease is not canonical: {}",
504 path.display()
505 );
506 }
507 Ok(lease)
508 }
509
510 pub(super) fn active_retention_protection(&self) -> Result<ActiveRetentionProtection> {
514 let guard = self.acquire_retention_roots_lock(RetentionRootsLockMode::Shared)?;
515 let profile_roots =
516 crate::socialgraph::acquire_profile_root_snapshot_guard(self.base_path())?;
517 let mut roots = match profile_roots.as_ref() {
518 Some(profile_roots) => profile_roots.retention_roots()?,
519 None => Vec::new(),
520 };
521 roots.extend(self.profile_repair_retention_roots()?);
522 roots.sort_by_key(Cid::to_string);
523 roots.dedup();
524 let hashes = self.collect_socialgraph_protected(&roots)?;
525 Ok(ActiveRetentionProtection {
526 _guard: guard,
527 _profile_roots: profile_roots,
528 hashes,
529 })
530 }
531
532 pub fn retain_nostr_root(&self, root: &Cid, apply: bool) -> Result<RootRetentionReport> {
539 let retention = self.active_retention_protection()?;
540 let tree = HashTree::new(HashTreeConfig::new(self.store_arc()));
541 let mut reachable = sync_block_on(async {
542 let mut reachable = HashSet::new();
543 let mut stack = vec![(root.clone(), "root".to_string())];
544 while let Some((cid, path)) = stack.pop() {
545 if !reachable.insert(cid.hash) {
546 continue;
547 }
548 let node = tree
549 .get_tree_node_by_cid(&cid)
550 .await
551 .map_err(|error| anyhow::anyhow!("read retained root DAG: {error}"))?
552 .ok_or_else(|| {
553 anyhow::anyhow!(
554 "retained directory node {} is missing at {path}",
555 to_hex(&cid.hash)
556 )
557 })?;
558 for link in node.links {
559 if link.link_type.is_directory_like() {
560 let name = link.name.as_deref().unwrap_or("<unnamed>");
561 stack.push((link.to_cid(), format!("{path}/{name}")));
562 } else {
563 reachable.insert(link.hash);
564 }
565 }
566 }
567 Ok::<_, anyhow::Error>(reachable)
568 })?;
569 reachable.extend(retention.hashes().iter().copied());
570
571 let rtxn = self.env.read_txn()?;
572 let pinned = self
573 .pins
574 .iter(&rtxn)?
575 .filter_map(std::result::Result::ok)
576 .filter_map(|(hash, _)| hash.try_into().ok())
577 .collect::<HashSet<Hash>>();
578 drop(rtxn);
579
580 let stats_before = self
581 .router
582 .writable_stats()
583 .map_err(|error| anyhow::anyhow!("read writable storage stats: {error}"))?;
584 let all_hashes = self
585 .router
586 .list_writable()
587 .map_err(|error| anyhow::anyhow!("list writable hashes: {error}"))?;
588 let candidates = all_hashes
589 .iter()
590 .filter(|hash| !reachable.contains(*hash) && !pinned.contains(*hash))
591 .copied()
592 .collect::<Vec<_>>();
593
594 let mut deleted = 0usize;
595 if apply {
596 const DELETE_BATCH_SIZE: usize = 16_384;
597 for (batch_index, batch) in candidates.chunks(DELETE_BATCH_SIZE).enumerate() {
598 deleted = deleted.saturating_add(
599 self.router
600 .delete_many_local_only(batch)
601 .map_err(|error| anyhow::anyhow!("delete unreachable batch: {error}"))?,
602 );
603 if (batch_index + 1).is_multiple_of(32) {
604 eprintln!(
605 "Retained-root cleanup: deleted {deleted}/{} unreachable hashes",
606 candidates.len()
607 );
608 }
609 }
610 }
611 let stats_after = self
612 .router
613 .writable_stats()
614 .map_err(|error| anyhow::anyhow!("read writable storage stats: {error}"))?;
615
616 Ok(RootRetentionReport {
617 total_hashes: all_hashes.len(),
618 reachable_hashes: reachable.len(),
619 pinned_hashes: pinned.len(),
620 candidate_hashes: candidates.len(),
621 deleted_hashes: deleted,
622 logical_bytes_before: stats_before.total_bytes,
623 logical_bytes_after: stats_after.total_bytes,
624 })
625 }
626
627 async fn collect_tree_hashes<S: Store>(
628 &self,
629 tree: &HashTree<S>,
630 root: &Cid,
631 require_tree_root: bool,
632 ) -> Result<HashSet<Hash>> {
633 let mut hashes = HashSet::new();
634 let mut stack = vec![(root.clone(), true)];
635
636 while let Some((cid, is_root)) = stack.pop() {
637 if !hashes.insert(cid.hash) {
638 continue;
639 }
640
641 let node = tree
642 .get_node(&cid)
643 .await
644 .map_err(|e| anyhow::anyhow!("Failed to get protected tree node: {}", e))?;
645 if let Some(node) = node {
646 for link in &node.links {
647 stack.push((link.to_cid(), false));
648 }
649 } else if is_root && require_tree_root {
650 anyhow::bail!(
651 "protected retention root {} is missing or is not a tree",
652 cid
653 );
654 }
655 }
656
657 Ok(hashes)
658 }
659
660 fn profile_repair_retention_roots(&self) -> Result<Vec<Cid>> {
661 let path = self.profile_repair_retention_lease_path();
662 let metadata = match std::fs::symlink_metadata(&path) {
663 Ok(metadata) => metadata,
664 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
665 Err(error) => {
666 return Err(error).with_context(|| format!("inspect {}", path.display()));
667 }
668 };
669 if metadata.file_type().is_symlink() || !metadata.file_type().is_file() {
670 anyhow::bail!(
671 "profile repair retention lease is not a direct regular file: {}",
672 path.display()
673 );
674 }
675 if metadata.len() > MAX_PROFILE_REPAIR_RETENTION_LEASE_BYTES {
676 anyhow::bail!(
677 "profile repair retention lease exceeds {} bytes: {}",
678 MAX_PROFILE_REPAIR_RETENTION_LEASE_BYTES,
679 path.display()
680 );
681 }
682 let bytes = std::fs::read(&path).with_context(|| format!("read {}", path.display()))?;
683 let lease: ProfileRepairRetentionLease =
684 serde_json::from_slice(&bytes).with_context(|| format!("decode {}", path.display()))?;
685 if lease.canonical_bytes()? != bytes {
686 anyhow::bail!(
687 "profile repair retention lease is not canonical: {}",
688 path.display()
689 );
690 }
691 lease.root_cids()
692 }
693
694 fn collect_socialgraph_protected(&self, roots: &[Cid]) -> Result<HashSet<Hash>> {
695 let mut protected = HashSet::new();
696 let tree = HashTree::new(HashTreeConfig::new(self.store_arc()).public());
697 for root in roots {
698 protected.extend(sync_block_on(self.collect_tree_hashes(&tree, root, true))?);
699 }
700 Ok(protected)
701 }
702
703 fn socialgraph_protected_for_orphan_sweep(&self, roots: &[Cid]) -> Result<Arc<HashSet<Hash>>> {
710 let root_ids = roots.iter().map(Cid::to_string).collect::<Vec<_>>();
711 {
712 let state = self
713 .orphan_scan
714 .lock()
715 .unwrap_or_else(std::sync::PoisonError::into_inner);
716 if state.socialgraph_roots.as_ref() == Some(&root_ids) {
717 return Ok(Arc::clone(&state.socialgraph_protected));
718 }
719 }
720
721 let newly_protected = self.collect_socialgraph_protected(roots)?;
722 let mut state = self
723 .orphan_scan
724 .lock()
725 .unwrap_or_else(std::sync::PoisonError::into_inner);
726 if state.socialgraph_roots.as_ref() == Some(&root_ids) {
727 return Ok(Arc::clone(&state.socialgraph_protected));
728 }
729
730 let protected = if state.sweep.is_some() && state.socialgraph_roots.is_some() {
731 let mut union = (*state.socialgraph_protected).clone();
732 union.extend(newly_protected);
733 union
734 } else {
735 newly_protected
736 };
737 state.socialgraph_roots = Some(root_ids);
738 state.socialgraph_protected = Arc::new(protected);
739 Ok(Arc::clone(&state.socialgraph_protected))
740 }
741
742 fn metadata_protects_orphan(&self, hash: &Hash) -> Result<bool> {
743 let rtxn = self.env.read_txn()?;
744 if self.pins.get(&rtxn, hash.as_slice())?.is_some() {
745 return Ok(true);
746 }
747 if self
748 .blob_trees
749 .prefix_iter(&rtxn, hash.as_slice())?
750 .next()
751 .transpose()?
752 .is_some()
753 {
754 return Ok(true);
755 }
756 let has_owner = self
757 .blob_owners
758 .prefix_iter(&rtxn, hash.as_slice())?
759 .next()
760 .transpose()?
761 .is_some();
762 Ok(has_owner)
763 }
764
765 fn evict_disposable_orphans_page(
766 &self,
767 target_bytes: u64,
768 additional_protected: &HashSet<Hash>,
769 page_size: usize,
770 ) -> Result<OrphanCleanupProgress> {
771 if page_size == 0 {
772 anyhow::bail!("orphan cleanup page size must be greater than zero");
773 }
774
775 let _retention_roots = self.acquire_retention_roots_lock(RetentionRootsLockMode::Shared)?;
776 let profile_roots =
777 crate::socialgraph::acquire_profile_root_snapshot_guard(self.base_path())?;
778 let mut roots = match profile_roots.as_ref() {
779 Some(profile_roots) => profile_roots.retention_roots()?,
780 None => Vec::new(),
781 };
782 roots.extend(self.profile_repair_retention_roots()?);
783 roots.sort_by_key(Cid::to_string);
784 roots.dedup();
785 let stats = self
786 .router
787 .writable_stats()
788 .map_err(|e| anyhow::anyhow!("Failed to get writable stats: {}", e))?;
789 let mut current_size = stats.total_bytes;
790 if current_size <= target_bytes {
791 let mut state = self
792 .orphan_scan
793 .lock()
794 .unwrap_or_else(std::sync::PoisonError::into_inner);
795 state.sweep = None;
796 state.socialgraph_roots = None;
797 state.socialgraph_protected = Arc::new(HashSet::new());
798 return Ok(OrphanCleanupProgress {
799 freed_bytes: 0,
800 scanned: 0,
801 sweep_complete: false,
802 });
803 }
804
805 let (after, start_after, wrapped) = {
806 let mut state = self
807 .orphan_scan
808 .lock()
809 .unwrap_or_else(std::sync::PoisonError::into_inner);
810 if state.sweep.is_none() {
811 state.sweep = Some(super::OrphanSweep {
812 start_after: state.cursor,
813 wrapped: false,
814 });
815 }
816 let sweep = state.sweep.expect("orphan sweep initialized");
817 (state.cursor, sweep.start_after, sweep.wrapped)
818 };
819 let socialgraph_protected = self.socialgraph_protected_for_orphan_sweep(&roots)?;
820 let mut candidates = self
821 .router
822 .scan_writable_hashes_after(after, page_size)
823 .map_err(|e| anyhow::anyhow!("Failed to scan writable hashes: {}", e))?;
824 let backend_page_len = candidates.len();
825 let mut crossed_start = false;
826 if wrapped {
827 let start_after = start_after.expect("only a non-zero cursor sweep wraps");
828 let keep = candidates.partition_point(|hash| *hash <= start_after);
829 crossed_start = keep < candidates.len();
830 candidates.truncate(keep);
831 }
832
833 let mut freed = 0u64;
834 let mut scanned = 0usize;
835 let mut last_examined = None;
836 for hash in &candidates {
837 if current_size <= target_bytes {
838 break;
839 }
840 let hash = *hash;
841 scanned += 1;
842 last_examined = Some(hash);
843
844 if socialgraph_protected.contains(&hash)
845 || additional_protected.contains(&hash)
846 || self.metadata_protects_orphan(&hash)?
847 {
848 continue;
849 }
850
851 let Some(_delete_guard) = self.cache_quota.begin_retention_delete(hash) else {
852 continue;
853 };
854 if self.metadata_protects_orphan(&hash)? {
858 continue;
859 }
860
861 let Some(size) = self
862 .router
863 .blob_size_sync(&hash)
864 .map_err(|e| anyhow::anyhow!("Failed to get blob size: {}", e))?
865 else {
866 continue;
867 };
868
869 if self
870 .router
871 .delete_local_only(&hash)
872 .map_err(|e| anyhow::anyhow!("Failed to delete orphaned blob: {}", e))?
873 {
874 freed = freed.saturating_add(size);
875 current_size = current_size.saturating_sub(size);
876 tracing::debug!(
877 "Deleted disposable orphaned blob {} ({} bytes)",
878 &to_hex(&hash)[..8],
879 size
880 );
881 }
882 }
883
884 let target_reached = current_size <= target_bytes;
885 let processed_whole_page = scanned == candidates.len();
886 let backend_exhausted = backend_page_len < page_size;
887 let mut sweep_complete = false;
888 let mut state = self
889 .orphan_scan
890 .lock()
891 .unwrap_or_else(std::sync::PoisonError::into_inner);
892 if let Some(last_examined) = last_examined {
893 state.cursor = Some(last_examined);
894 }
895
896 if target_reached {
897 state.sweep = None;
898 state.socialgraph_roots = None;
899 state.socialgraph_protected = Arc::new(HashSet::new());
900 } else if processed_whole_page {
901 let sweep = state.sweep.expect("active orphan sweep");
902 if sweep.wrapped {
903 if crossed_start || backend_exhausted {
904 state.cursor = sweep.start_after;
905 state.sweep = None;
906 sweep_complete = true;
907 }
908 } else if backend_exhausted {
909 if sweep.start_after.is_none() {
910 state.cursor = None;
911 state.sweep = None;
912 sweep_complete = true;
913 } else {
914 state.cursor = None;
915 state.sweep = Some(super::OrphanSweep {
916 wrapped: true,
917 ..sweep
918 });
919 }
920 }
921 if sweep_complete {
922 state.socialgraph_roots = None;
923 state.socialgraph_protected = Arc::new(HashSet::new());
924 }
925 }
926
927 Ok(OrphanCleanupProgress {
928 freed_bytes: freed,
929 scanned,
930 sweep_complete,
931 })
932 }
933
934 fn evict_disposable_orphans_to_target_raw(
935 &self,
936 target_bytes: u64,
937 additional_protected: &HashSet<Hash>,
938 ) -> Result<OrphanCleanupProgress> {
939 self.evict_disposable_orphans_page(
940 target_bytes,
941 additional_protected,
942 ORPHAN_SCAN_PAGE_SIZE,
943 )
944 }
945
946 fn evict_disposable_orphans_to_target(&self, target_bytes: u64) -> Result<u64> {
947 let cleanup = self
948 .cache_quota
949 .begin_standalone_cleanup()
950 .map_err(|denial| anyhow::anyhow!(denial.to_string()))?;
951 let result =
952 self.evict_disposable_orphans_to_target_raw(target_bytes, cleanup.inflight_hashes());
953 if result.is_ok() {
954 let after_usage = self
955 .router
956 .writable_stats()
957 .map_err(|error| {
958 anyhow::anyhow!("Failed to get writable stats after cache cleanup: {error}")
959 })?
960 .total_bytes;
961 cleanup.complete(after_usage);
962 }
963 result.map(|progress| progress.freed_bytes)
964 }
965
966 pub(super) fn prepare_cached_blob_write(
967 &self,
968 incoming_bytes: u64,
969 hashes: Vec<Hash>,
970 force_cleanup: bool,
971 ) -> Result<CacheWritePermit<'_>> {
972 if let Some(denial) = self.cache_quota.quick_denial() {
973 anyhow::bail!(denial.to_string());
974 }
975
976 let observed_usage = if self.max_size_bytes == 0 {
977 0
978 } else {
979 self.router
980 .writable_stats()
981 .map_err(|error| anyhow::anyhow!("Failed to get writable stats: {error}"))?
982 .total_bytes
983 };
984 match self
985 .cache_quota
986 .begin_admission(
987 observed_usage,
988 incoming_bytes,
989 hashes,
990 self.max_size_bytes,
991 force_cleanup,
992 )
993 .map_err(|denial| anyhow::anyhow!(denial.to_string()))?
994 {
995 CacheQuotaAdmission::Admitted(permit) => Ok(permit),
996 CacheQuotaAdmission::Cleanup(cleanup) => {
997 let progress = self.evict_disposable_orphans_to_target_raw(
998 cleanup.target_bytes(),
999 cleanup.inflight_hashes(),
1000 )?;
1001 let after_usage = self
1002 .router
1003 .writable_stats()
1004 .map_err(|error| {
1005 anyhow::anyhow!("Failed to get writable stats after cache cleanup: {error}")
1006 })?
1007 .total_bytes;
1008 cleanup
1009 .complete(after_usage, progress.freed_bytes, progress.sweep_complete)
1010 .map_err(|denial| anyhow::anyhow!(denial.to_string()))
1011 }
1012 }
1013 }
1014
1015 pub fn cache_cleanup_epoch_count(&self) -> u64 {
1020 self.cache_quota.cleanup_epoch_count()
1021 }
1022
1023 pub fn make_room_for_cached_blob(&self, incoming_bytes: u64) -> Result<u64> {
1024 if self.max_size_bytes == 0 {
1025 return Ok(0);
1026 }
1027
1028 let stats = self
1029 .router
1030 .writable_stats()
1031 .map_err(|e| anyhow::anyhow!("Failed to get writable stats: {}", e))?;
1032 if stats.total_bytes.saturating_add(incoming_bytes) <= self.max_size_bytes {
1033 return Ok(0);
1034 }
1035
1036 let target = if incoming_bytes >= self.max_size_bytes {
1037 0
1038 } else {
1039 (self.max_size_bytes.saturating_mul(9) / 10)
1040 .min(self.max_size_bytes.saturating_sub(incoming_bytes))
1041 };
1042 self.evict_disposable_orphans_to_target(target)
1043 }
1044
1045 pub fn enforce_cached_blob_budget_after_insert(&self, inserted_bytes: u64) -> Result<u64> {
1046 if self.max_size_bytes == 0 || inserted_bytes == 0 {
1047 return Ok(0);
1048 }
1049
1050 let stats = self
1051 .router
1052 .writable_stats()
1053 .map_err(|e| anyhow::anyhow!("Failed to get writable stats: {}", e))?;
1054 if stats.total_bytes <= self.max_size_bytes {
1055 return Ok(0);
1056 }
1057
1058 let target = if inserted_bytes >= self.max_size_bytes {
1059 inserted_bytes
1060 } else {
1061 (self.max_size_bytes.saturating_mul(9) / 10)
1062 .saturating_add(inserted_bytes)
1063 .min(self.max_size_bytes)
1064 };
1065 self.evict_disposable_orphans_to_target(target)
1066 }
1067
1068 pub fn make_room_for_durable_blob(&self, incoming_bytes: u64) -> Result<u64> {
1069 if self.max_size_bytes == 0 || incoming_bytes == 0 {
1070 return Ok(0);
1071 }
1072
1073 if incoming_bytes > self.max_size_bytes {
1074 anyhow::bail!(
1075 "storage limit exceeded: incoming blob is {} bytes but limit is {} bytes",
1076 incoming_bytes,
1077 self.max_size_bytes
1078 );
1079 }
1080
1081 let stats = self
1082 .router
1083 .writable_stats()
1084 .map_err(|e| anyhow::anyhow!("Failed to get writable stats: {}", e))?;
1085 if stats.total_bytes.saturating_add(incoming_bytes) <= self.max_size_bytes {
1086 return Ok(0);
1087 }
1088
1089 let target = (self.max_size_bytes.saturating_mul(9) / 10)
1090 .min(self.max_size_bytes.saturating_sub(incoming_bytes));
1091 let freed = self.evict_with_policy_to_target(stats.total_bytes, target)?;
1092
1093 let next_stats = self
1094 .router
1095 .writable_stats()
1096 .map_err(|e| anyhow::anyhow!("Failed to get writable stats after eviction: {}", e))?;
1097 if next_stats.total_bytes.saturating_add(incoming_bytes) > self.max_size_bytes {
1098 anyhow::bail!(
1099 "storage limit exceeded: {} bytes used, {} byte incoming blob, {} byte limit",
1100 next_stats.total_bytes,
1101 incoming_bytes,
1102 self.max_size_bytes
1103 );
1104 }
1105
1106 Ok(freed)
1107 }
1108
1109 pub fn enforce_durable_blob_budget_after_insert(&self, inserted_bytes: u64) -> Result<u64> {
1110 if self.max_size_bytes == 0 || inserted_bytes == 0 {
1111 return Ok(0);
1112 }
1113
1114 if inserted_bytes > self.max_size_bytes {
1115 anyhow::bail!(
1116 "storage limit exceeded: inserted blobs are {} bytes but limit is {} bytes",
1117 inserted_bytes,
1118 self.max_size_bytes
1119 );
1120 }
1121
1122 let stats = self
1123 .router
1124 .writable_stats()
1125 .map_err(|e| anyhow::anyhow!("Failed to get writable stats: {}", e))?;
1126 if stats.total_bytes <= self.max_size_bytes {
1127 return Ok(0);
1128 }
1129
1130 let target = (self.max_size_bytes.saturating_mul(9) / 10)
1131 .saturating_add(inserted_bytes)
1132 .min(self.max_size_bytes);
1133 let freed = self.evict_with_policy_to_target(stats.total_bytes, target)?;
1134
1135 let next_stats = self
1136 .router
1137 .writable_stats()
1138 .map_err(|e| anyhow::anyhow!("Failed to get writable stats after eviction: {}", e))?;
1139 if next_stats.total_bytes > self.max_size_bytes {
1140 anyhow::bail!(
1141 "storage limit exceeded: {} bytes used after inserting {} bytes, {} byte limit",
1142 next_stats.total_bytes,
1143 inserted_bytes,
1144 self.max_size_bytes
1145 );
1146 }
1147
1148 Ok(freed)
1149 }
1150
1151 pub fn relieve_cached_blob_write_pressure(&self, incoming_bytes: u64) -> Result<u64> {
1152 let stats = self
1153 .router
1154 .writable_stats()
1155 .map_err(|e| anyhow::anyhow!("Failed to get writable stats: {}", e))?;
1156 if stats.total_bytes == 0 {
1157 return Ok(0);
1158 }
1159
1160 let headroom = incoming_bytes.max(stats.total_bytes / 10).max(1);
1161 let target = stats.total_bytes.saturating_sub(headroom);
1162 self.evict_disposable_orphans_to_target(target)
1163 }
1164
1165 pub fn pin(&self, hash: &[u8; 32]) -> Result<()> {
1167 let mut wtxn = self.env.write_txn()?;
1168 self.pins.put(&mut wtxn, hash.as_slice(), &())?;
1169 wtxn.commit()?;
1170 Ok(())
1171 }
1172
1173 pub fn unpin(&self, hash: &[u8; 32]) -> Result<()> {
1175 let mut wtxn = self.env.write_txn()?;
1176 self.pins.delete(&mut wtxn, hash.as_slice())?;
1177 wtxn.commit()?;
1178 Ok(())
1179 }
1180
1181 pub fn is_pinned(&self, hash: &[u8; 32]) -> Result<bool> {
1183 let rtxn = self.env.read_txn()?;
1184 Ok(self.pins.get(&rtxn, hash.as_slice())?.is_some())
1185 }
1186
1187 pub fn list_pins_raw(&self) -> Result<Vec<[u8; 32]>> {
1189 let rtxn = self.env.read_txn()?;
1190 let mut pins = Vec::new();
1191
1192 for item in self.pins.iter(&rtxn)? {
1193 let (hash_bytes, _) = item?;
1194 if hash_bytes.len() == 32 {
1195 let mut hash = [0u8; 32];
1196 hash.copy_from_slice(hash_bytes);
1197 pins.push(hash);
1198 }
1199 }
1200
1201 Ok(pins)
1202 }
1203
1204 pub fn list_pins_with_names(&self) -> Result<Vec<PinnedItem>> {
1206 let rtxn = self.env.read_txn()?;
1207 let store = self.store_arc();
1208 let tree = HashTree::new(HashTreeConfig::new(store).public());
1209 let mut pins = Vec::new();
1210
1211 for item in self.pins.iter(&rtxn)? {
1212 let (hash_bytes, _) = item?;
1213 if hash_bytes.len() != 32 {
1214 continue;
1215 }
1216 let mut hash = [0u8; 32];
1217 hash.copy_from_slice(hash_bytes);
1218
1219 let is_directory =
1221 sync_block_on(async { tree.is_directory(&hash).await.unwrap_or(false) });
1222
1223 let meta = self
1224 .tree_meta
1225 .get(&rtxn, hash.as_slice())?
1226 .map(|bytes| {
1227 rmp_serde::from_slice::<TreeMeta>(bytes)
1228 .map_err(|e| anyhow::anyhow!("Failed to deserialize TreeMeta: {}", e))
1229 })
1230 .transpose()?;
1231 let size_bytes = if let Some(meta) = meta.as_ref() {
1232 meta.total_size
1233 } else {
1234 self.router
1235 .blob_size_sync(&hash)
1236 .map_err(|e| anyhow::anyhow!("Failed to get pinned blob size: {}", e))?
1237 .unwrap_or(0)
1238 };
1239
1240 pins.push(PinnedItem {
1241 cid: to_hex(&hash),
1242 name: pinned_item_name(&hash, meta.as_ref()),
1243 is_directory,
1244 size_bytes,
1245 });
1246 }
1247
1248 Ok(pins)
1249 }
1250
1251 pub fn owned_blob_stats(&self) -> Result<Vec<OwnedBlobStats>> {
1252 let rtxn = self.env.read_txn()?;
1253 let mut owners = Vec::new();
1254
1255 for item in self.pubkey_blobs.iter(&rtxn)? {
1256 let (owner_bytes, blobs_bytes) = item?;
1257 if owner_bytes.len() != 32 {
1258 continue;
1259 }
1260
1261 let blobs: Vec<BlobMetadata> = serde_json::from_slice(blobs_bytes)
1262 .map_err(|e| anyhow::anyhow!("Failed to deserialize blob metadata: {}", e))?;
1263 let mut owner = [0u8; 32];
1264 owner.copy_from_slice(owner_bytes);
1265 let total_bytes = blobs
1266 .iter()
1267 .fold(0u64, |total, blob| total.saturating_add(blob.size));
1268 owners.push(OwnedBlobStats {
1269 owner,
1270 count: blobs.len(),
1271 total_bytes,
1272 });
1273 }
1274
1275 owners.sort_by_key(|stats| stats.owner);
1276 Ok(owners)
1277 }
1278
1279 pub fn tree_index_limits(&self) -> TreeIndexLimits {
1285 TreeIndexLimits {
1286 max_nodes: MAX_PINNED_TREE_NODES,
1287 max_bytes: if self.max_size_bytes == 0 {
1288 MAX_UNBOUNDED_PINNED_TREE_BYTES
1289 } else {
1290 self.max_size_bytes
1291 },
1292 }
1293 }
1294
1295 pub fn pin_and_index_tree(
1299 &self,
1300 root: &Cid,
1301 owner: &str,
1302 name: Option<&str>,
1303 priority: u8,
1304 limits: TreeIndexLimits,
1305 ) -> std::result::Result<PinTreeResult, PinTreeError> {
1306 let store = self.store_arc();
1307 let tree = HashTree::new(HashTreeConfig::new(store).public());
1308 let plan = sync_block_on(self.collect_tree_index(&tree, root, limits))?;
1309 let already_pinned = self
1310 .write_tree_index(
1311 &root.hash,
1312 &plan.tracked_hashes,
1313 plan.total_size,
1314 owner,
1315 name,
1316 priority,
1317 None,
1318 true,
1319 )
1320 .map_err(|error| PinTreeError::Storage(error.to_string()))?;
1321
1322 Ok(PinTreeResult {
1323 indexed_hashes: plan.tracked_hashes.len(),
1324 total_size: plan.total_size,
1325 already_pinned,
1326 })
1327 }
1328
1329 pub fn index_tree(
1334 &self,
1335 root_hash: &Hash,
1336 owner: &str,
1337 name: Option<&str>,
1338 priority: u8,
1339 ref_key: Option<&str>,
1340 ) -> Result<()> {
1341 let root_hex = to_hex(root_hash);
1342
1343 if let Some(key) = ref_key {
1345 let rtxn = self.env.read_txn()?;
1346 if let Some(old_hash_bytes) = self.tree_refs.get(&rtxn, key)? {
1347 if old_hash_bytes != root_hash.as_slice() {
1348 let old_hash: Hash = old_hash_bytes
1349 .try_into()
1350 .map_err(|_| anyhow::anyhow!("Invalid hash in tree_refs"))?;
1351 drop(rtxn);
1352 let _ = self.unpin(&old_hash);
1353 let _ = self.unindex_tree(&old_hash);
1355 tracing::debug!("Replaced old tree for ref {}", key);
1356 }
1357 }
1358 }
1359
1360 let store = self.store_arc();
1361 let tree = HashTree::new(HashTreeConfig::new(store).public());
1362
1363 let plan = sync_block_on(self.collect_tree_index(
1364 &tree,
1365 &Cid::public(*root_hash),
1366 TreeIndexLimits {
1367 max_nodes: MAX_PINNED_TREE_NODES,
1368 max_bytes: MAX_UNBOUNDED_PINNED_TREE_BYTES,
1369 },
1370 ))?;
1371 self.write_tree_index(
1372 root_hash,
1373 &plan.tracked_hashes,
1374 plan.total_size,
1375 owner,
1376 name,
1377 priority,
1378 ref_key,
1379 false,
1380 )?;
1381
1382 tracing::debug!(
1383 "Indexed tree {} ({} blobs, {} bytes, priority {})",
1384 &root_hex[..8],
1385 plan.tracked_hashes.len(),
1386 plan.total_size,
1387 priority
1388 );
1389
1390 Ok(())
1391 }
1392
1393 #[allow(clippy::too_many_arguments)]
1394 fn write_tree_index(
1395 &self,
1396 root_hash: &Hash,
1397 tracked_hashes: &HashSet<Hash>,
1398 total_size: u64,
1399 owner: &str,
1400 name: Option<&str>,
1401 priority: u8,
1402 ref_key: Option<&str>,
1403 pin: bool,
1404 ) -> Result<bool> {
1405 let mut wtxn = self.env.write_txn()?;
1406 let already_pinned = self.pins.get(&wtxn, root_hash.as_slice())?.is_some();
1407
1408 for tracked_hash in tracked_hashes {
1409 let mut key = [0u8; 64];
1410 key[..32].copy_from_slice(tracked_hash);
1411 key[32..].copy_from_slice(root_hash);
1412 self.blob_trees.put(&mut wtxn, &key[..], &())?;
1413 }
1414
1415 let meta = TreeMeta {
1416 owner: owner.to_string(),
1417 name: name.map(str::to_string),
1418 synced_at: unix_timestamp_now(),
1419 total_size,
1420 priority,
1421 };
1422 let meta_bytes = rmp_serde::to_vec(&meta)
1423 .map_err(|error| anyhow::anyhow!("Failed to serialize TreeMeta: {error}"))?;
1424 self.tree_meta
1425 .put(&mut wtxn, root_hash.as_slice(), &meta_bytes)?;
1426
1427 if let Some(key) = ref_key {
1428 self.tree_refs.put(&mut wtxn, key, root_hash.as_slice())?;
1429 }
1430 if pin {
1431 self.pins.put(&mut wtxn, root_hash.as_slice(), &())?;
1432 }
1433
1434 wtxn.commit()?;
1435 Ok(already_pinned)
1436 }
1437
1438 async fn collect_tree_index<S: Store>(
1439 &self,
1440 tree: &HashTree<S>,
1441 root: &Cid,
1442 limits: TreeIndexLimits,
1443 ) -> std::result::Result<TreeIndexPlan, PinTreeError> {
1444 let mut hashes = HashSet::new();
1445 let mut visited = HashSet::new();
1446 let mut total_size = 0u64;
1447 let mut stored_size = 0u64;
1448 let mut stack = vec![(root.clone(), true, true, false)];
1450
1451 while let Some((cid, count_bytes, follow_tree, require_tree)) = stack.pop() {
1452 let visit_key = (cid.hash, cid.key, follow_tree);
1453 if !visited.insert(visit_key) {
1454 continue;
1455 }
1456 if visited.len() > limits.max_nodes {
1457 return Err(PinTreeError::NodeLimitExceeded {
1458 max_nodes: limits.max_nodes,
1459 });
1460 }
1461
1462 let size = self
1463 .router
1464 .blob_size_sync(&cid.hash)
1465 .map_err(|error| PinTreeError::Storage(error.to_string()))?
1466 .ok_or_else(|| {
1467 if cid.hash == root.hash {
1468 PinTreeError::MissingRoot {
1469 hash: to_hex(&cid.hash),
1470 }
1471 } else {
1472 PinTreeError::MissingDescendant {
1473 hash: to_hex(&cid.hash),
1474 }
1475 }
1476 })?;
1477 if hashes.insert(cid.hash) {
1478 stored_size = stored_size
1479 .checked_add(size)
1480 .filter(|size| *size <= limits.max_bytes)
1481 .ok_or(PinTreeError::ByteLimitExceeded {
1482 max_bytes: limits.max_bytes,
1483 })?;
1484 }
1485
1486 if !follow_tree {
1487 continue;
1488 }
1489
1490 let node = tree.get_node(&cid).await.map_err(|error| match error {
1491 HashTreeError::Store(message) => PinTreeError::Storage(message),
1492 error => PinTreeError::InvalidDag {
1493 hash: to_hex(&cid.hash),
1494 message: error.to_string(),
1495 },
1496 })?;
1497 let Some(node) = node else {
1498 if require_tree {
1499 return Err(PinTreeError::InvalidDag {
1500 hash: to_hex(&cid.hash),
1501 message: "directory link does not contain a tree node".to_string(),
1502 });
1503 }
1504 if count_bytes {
1505 total_size = total_size
1506 .checked_add(size)
1507 .filter(|size| *size <= limits.max_bytes)
1508 .ok_or(PinTreeError::ByteLimitExceeded {
1509 max_bytes: limits.max_bytes,
1510 })?;
1511 }
1512 continue;
1513 };
1514
1515 if visited
1516 .len()
1517 .saturating_add(stack.len())
1518 .saturating_add(node.links.len())
1519 > limits.max_nodes
1520 {
1521 return Err(PinTreeError::NodeLimitExceeded {
1522 max_nodes: limits.max_nodes,
1523 });
1524 }
1525
1526 for link in &node.links {
1527 match link.link_type {
1528 LinkType::Blob => {
1529 if count_bytes {
1530 total_size = total_size
1531 .checked_add(link.size)
1532 .filter(|size| *size <= limits.max_bytes)
1533 .ok_or(PinTreeError::ByteLimitExceeded {
1534 max_bytes: limits.max_bytes,
1535 })?;
1536 }
1537 stack.push((link.to_cid(), false, false, false));
1538 }
1539 LinkType::File => {
1540 if count_bytes {
1541 total_size = total_size
1542 .checked_add(link.size)
1543 .filter(|size| *size <= limits.max_bytes)
1544 .ok_or(PinTreeError::ByteLimitExceeded {
1545 max_bytes: limits.max_bytes,
1546 })?;
1547 }
1548 stack.push((link.to_cid(), false, true, false));
1549 }
1550 LinkType::Dir | LinkType::Fanout => {
1551 stack.push((link.to_cid(), count_bytes, true, true));
1552 }
1553 }
1554 }
1555 }
1556
1557 Ok(TreeIndexPlan {
1558 tracked_hashes: hashes,
1559 total_size,
1560 })
1561 }
1562
1563 pub fn unindex_tree(&self, root_hash: &Hash) -> Result<u64> {
1566 let cleanup = self
1567 .cache_quota
1568 .begin_standalone_cleanup()
1569 .map_err(|denial| anyhow::anyhow!(denial.to_string()))?;
1570 let retention = self.active_retention_protection()?;
1571 let mut protected = cleanup.inflight_hashes().clone();
1572 protected.extend(retention.hashes().iter().copied());
1573 let result = self.unindex_tree_raw(root_hash, &protected);
1574 if result.is_ok() {
1575 let after_usage = self
1576 .router
1577 .writable_stats()
1578 .map_err(|error| {
1579 anyhow::anyhow!("Failed to get writable stats after tree unindex: {error}")
1580 })?
1581 .total_bytes;
1582 cleanup.complete(after_usage);
1583 }
1584 result
1585 }
1586
1587 fn unindex_tree_raw(
1588 &self,
1589 root_hash: &Hash,
1590 additional_protected: &HashSet<Hash>,
1591 ) -> Result<u64> {
1592 let root_hex = to_hex(root_hash);
1593
1594 let store = self.store_arc();
1595 let tree = HashTree::new(HashTreeConfig::new(store).public());
1596
1597 let tracked_hashes =
1599 sync_block_on(self.collect_tree_hashes(&tree, &Cid::public(*root_hash), false))?;
1600
1601 let mut wtxn = self.env.write_txn()?;
1602 let mut freed = 0u64;
1603
1604 for tracked_hash in &tracked_hashes {
1606 let mut key = [0u8; 64];
1608 key[..32].copy_from_slice(tracked_hash);
1609 key[32..].copy_from_slice(root_hash);
1610 self.blob_trees.delete(&mut wtxn, &key[..])?;
1611
1612 let mut has_other_tree = false;
1614 for item in self.blob_trees.prefix_iter(&wtxn, &tracked_hash[..])? {
1615 if item.is_ok() {
1616 has_other_tree = true;
1617 break;
1618 }
1619 }
1620
1621 let has_owner = self
1622 .blob_owners
1623 .prefix_iter(&wtxn, tracked_hash.as_slice())?
1624 .next()
1625 .transpose()?
1626 .is_some();
1627
1628 if !has_other_tree && !has_owner && !additional_protected.contains(tracked_hash) {
1631 let Some(_delete_guard) = self.cache_quota.begin_retention_delete(*tracked_hash)
1632 else {
1633 continue;
1634 };
1635 if let Some(size) = self
1636 .router
1637 .blob_size_sync(tracked_hash)
1638 .map_err(|e| anyhow::anyhow!("Failed to get blob size: {}", e))?
1639 {
1640 freed += size;
1641 self.router
1643 .delete_local_only(tracked_hash)
1644 .map_err(|e| anyhow::anyhow!("Failed to delete blob: {}", e))?;
1645 }
1646 }
1647 }
1648
1649 self.tree_meta.delete(&mut wtxn, root_hash.as_slice())?;
1651
1652 wtxn.commit()?;
1653
1654 tracing::debug!("Unindexed tree {} ({} bytes freed)", &root_hex[..8], freed);
1655
1656 Ok(freed)
1657 }
1658
1659 pub fn get_tree_meta(&self, root_hash: &Hash) -> Result<Option<TreeMeta>> {
1661 let rtxn = self.env.read_txn()?;
1662 if let Some(bytes) = self.tree_meta.get(&rtxn, root_hash.as_slice())? {
1663 let meta: TreeMeta = rmp_serde::from_slice(bytes)
1664 .map_err(|e| anyhow::anyhow!("Failed to deserialize TreeMeta: {}", e))?;
1665 Ok(Some(meta))
1666 } else {
1667 Ok(None)
1668 }
1669 }
1670
1671 pub fn get_tree_ref(&self, key: &str) -> Result<Option<Hash>> {
1672 let rtxn = self.env.read_txn()?;
1673 let Some(bytes) = self.tree_refs.get(&rtxn, key)? else {
1674 return Ok(None);
1675 };
1676
1677 let hash: Hash = bytes
1678 .try_into()
1679 .map_err(|_| anyhow::anyhow!("Invalid hash in tree_refs"))?;
1680 Ok(Some(hash))
1681 }
1682
1683 pub fn list_indexed_trees(&self) -> Result<Vec<(Hash, TreeMeta)>> {
1685 let rtxn = self.env.read_txn()?;
1686 let mut trees = Vec::new();
1687
1688 for item in self.tree_meta.iter(&rtxn)? {
1689 let (hash_bytes, meta_bytes) = item?;
1690 let hash: Hash = hash_bytes
1691 .try_into()
1692 .map_err(|_| anyhow::anyhow!("Invalid hash in tree_meta"))?;
1693 let meta: TreeMeta = rmp_serde::from_slice(meta_bytes)
1694 .map_err(|e| anyhow::anyhow!("Failed to deserialize TreeMeta: {}", e))?;
1695 trees.push((hash, meta));
1696 }
1697
1698 Ok(trees)
1699 }
1700
1701 pub fn tracked_size(&self) -> Result<u64> {
1703 let rtxn = self.env.read_txn()?;
1704 let mut total = 0u64;
1705
1706 for item in self.tree_meta.iter(&rtxn)? {
1707 let (_, bytes) = item?;
1708 let meta: TreeMeta = rmp_serde::from_slice(bytes)
1709 .map_err(|e| anyhow::anyhow!("Failed to deserialize TreeMeta: {}", e))?;
1710 total += meta.total_size;
1711 }
1712
1713 Ok(total)
1714 }
1715
1716 fn get_evictable_trees(&self) -> Result<Vec<(Hash, TreeMeta)>> {
1722 let mut trees = self.list_indexed_trees()?;
1723
1724 trees.sort_by(|a, b| match a.1.priority.cmp(&b.1.priority) {
1726 std::cmp::Ordering::Equal => a.1.synced_at.cmp(&b.1.synced_at),
1727 other => other,
1728 });
1729
1730 Ok(trees)
1731 }
1732
1733 pub fn evict_if_needed(&self) -> Result<u64> {
1740 let stats = self
1742 .router
1743 .writable_stats()
1744 .map_err(|e| anyhow::anyhow!("Failed to get writable stats: {}", e))?;
1745 let current = stats.total_bytes;
1746
1747 if current <= self.max_size_bytes {
1748 return Ok(0);
1749 }
1750
1751 let target = self.max_size_bytes * 90 / 100;
1753 self.evict_with_policy_to_target(current, target)
1754 }
1755
1756 fn evict_with_policy_to_target(&self, current: u64, target: u64) -> Result<u64> {
1757 let cleanup = self
1758 .cache_quota
1759 .begin_standalone_cleanup()
1760 .map_err(|denial| anyhow::anyhow!(denial.to_string()))?;
1761 let mut freed = 0u64;
1762 let mut current_size = current;
1763
1764 if self.evict_orphans {
1766 let orphan_progress =
1767 self.evict_disposable_orphans_to_target_raw(target, cleanup.inflight_hashes())?;
1768 freed += orphan_progress.freed_bytes;
1769 current_size = current_size.saturating_sub(orphan_progress.freed_bytes);
1770
1771 if orphan_progress.freed_bytes > 0 {
1772 tracing::info!(
1773 "Evicted orphaned blobs: {} bytes freed",
1774 orphan_progress.freed_bytes
1775 );
1776 }
1777
1778 if current_size > target && !orphan_progress.sweep_complete {
1782 let after_usage = self
1783 .router
1784 .writable_stats()
1785 .map_err(|error| {
1786 anyhow::anyhow!(
1787 "Failed to get writable stats after bounded orphan cleanup: {error}"
1788 )
1789 })?
1790 .total_bytes;
1791 cleanup.complete(after_usage);
1792 return Ok(freed);
1793 }
1794 } else {
1795 tracing::debug!("Skipping orphan blob eviction; storage.evict_orphans=false");
1796 }
1797
1798 if current_size <= target {
1800 if freed > 0 {
1801 tracing::info!("Eviction complete: {} bytes freed", freed);
1802 }
1803 let after_usage = self
1804 .router
1805 .writable_stats()
1806 .map_err(|error| {
1807 anyhow::anyhow!("Failed to get writable stats after retention cleanup: {error}")
1808 })?
1809 .total_bytes;
1810 cleanup.complete(after_usage);
1811 return Ok(freed);
1812 }
1813
1814 let retention = self.active_retention_protection()?;
1817 let retention_protected = retention.hashes();
1818 let mut additional_protected = cleanup.inflight_hashes().clone();
1819 additional_protected.extend(retention_protected.iter().copied());
1820 let evictable = self.get_evictable_trees()?;
1821
1822 for (root_hash, meta) in evictable {
1823 if current_size <= target {
1824 break;
1825 }
1826
1827 let root_hex = to_hex(&root_hash);
1828
1829 if self.is_pinned(&root_hash)? {
1831 continue;
1832 }
1833 if retention_protected.contains(&root_hash) {
1834 continue;
1835 }
1836
1837 let tree_freed = self.unindex_tree_raw(&root_hash, &additional_protected)?;
1838 freed += tree_freed;
1839 current_size = current_size.saturating_sub(tree_freed);
1840
1841 tracing::info!(
1842 "Evicted tree {} (owner={}, priority={}, {} bytes)",
1843 &root_hex[..8],
1844 &meta.owner[..8.min(meta.owner.len())],
1845 meta.priority,
1846 tree_freed
1847 );
1848 }
1849
1850 if freed > 0 {
1851 tracing::info!("Eviction complete: {} bytes freed", freed);
1852 }
1853
1854 let after_usage = self
1855 .router
1856 .writable_stats()
1857 .map_err(|error| {
1858 anyhow::anyhow!("Failed to get writable stats after retention cleanup: {error}")
1859 })?
1860 .total_bytes;
1861 cleanup.complete(after_usage);
1862 Ok(freed)
1863 }
1864
1865 pub fn max_size_bytes(&self) -> u64 {
1867 self.max_size_bytes
1868 }
1869
1870 pub fn storage_by_priority(&self) -> Result<StorageByPriority> {
1872 let rtxn = self.env.read_txn()?;
1873 let mut own = 0u64;
1874 let mut followed = 0u64;
1875 let mut other = 0u64;
1876
1877 for item in self.tree_meta.iter(&rtxn)? {
1878 let (_, bytes) = item?;
1879 let meta: TreeMeta = rmp_serde::from_slice(bytes)
1880 .map_err(|e| anyhow::anyhow!("Failed to deserialize TreeMeta: {}", e))?;
1881
1882 if meta.priority == PRIORITY_OWN {
1883 own += meta.total_size;
1884 } else if meta.priority >= PRIORITY_FOLLOWED {
1885 followed += meta.total_size;
1886 } else {
1887 other += meta.total_size;
1888 }
1889 }
1890
1891 Ok(StorageByPriority {
1892 own,
1893 followed,
1894 other,
1895 })
1896 }
1897
1898 pub fn get_storage_stats(&self) -> Result<StorageStats> {
1900 let rtxn = self.env.read_txn()?;
1901 let total_pins = self.pins.len(&rtxn)? as usize;
1902
1903 let stats = self
1904 .router
1905 .stats()
1906 .map_err(|e| anyhow::anyhow!("Failed to get stats: {}", e))?;
1907
1908 Ok(StorageStats {
1909 total_dags: stats.count,
1910 pinned_dags: total_pins,
1911 total_bytes: stats.total_bytes,
1912 })
1913 }
1914}
1915
1916#[cfg(test)]
1917mod tests {
1918 use super::*;
1919 use hashtree_config::StorageBackend;
1920 use hashtree_core::Cid;
1921 use hashtree_index::{BTree, BTreeOptions};
1922 use nostr::{EventBuilder, Keys, Kind, Timestamp};
1923 use std::io::Write;
1924 use std::path::Path;
1925 use tempfile::TempDir;
1926
1927 use crate::storage::PRIORITY_OTHER;
1928
1929 #[cfg(unix)]
1930 #[test]
1931 fn existing_profile_repair_retention_guard_never_creates_its_authority() {
1932 let temp_dir = TempDir::new().expect("temp dir");
1933 let lock_path = temp_dir.path().join(RETENTION_ROOTS_LOCK_FILE);
1934
1935 let error = acquire_existing_profile_repair_retention_guard(temp_dir.path())
1936 .err()
1937 .expect("missing authority must fail");
1938 assert!(error
1939 .to_string()
1940 .contains("inspect existing retention-roots lock"));
1941 assert!(
1942 !lock_path.exists(),
1943 "recovery guard created a missing lock authority"
1944 );
1945
1946 std::fs::write(&lock_path, b"existing-authority").expect("create existing lock authority");
1947 let before = std::fs::metadata(&lock_path).expect("inspect existing lock authority");
1948 let guard = acquire_existing_profile_repair_retention_guard(temp_dir.path())
1949 .expect("acquire existing lock authority");
1950 let after = std::fs::metadata(&lock_path).expect("reinspect existing lock authority");
1951 assert_eq!(before.len(), after.len());
1952 assert_eq!(
1953 std::fs::read(&lock_path).expect("read existing lock authority"),
1954 b"existing-authority"
1955 );
1956 drop(guard);
1957 }
1958
1959 fn write_root_file(path: &Path, cid: &Cid) {
1960 #[derive(Serialize)]
1961 struct StoredCid {
1962 hash: [u8; 32],
1963 key: Option<[u8; 32]>,
1964 }
1965
1966 std::fs::create_dir_all(path.parent().expect("root file parent")).expect("create dir");
1967 let bytes = rmp_serde::to_vec_named(&StoredCid {
1968 hash: cid.hash,
1969 key: cid.key,
1970 })
1971 .expect("encode cid");
1972 std::fs::write(path, bytes).expect("write root file");
1973 }
1974
1975 fn build_test_tree(store: &HashtreeStore) -> Cid {
1976 let index = BTree::new(store.store_arc(), BTreeOptions { order: Some(8) });
1977 sync_block_on(index.build(vec![
1978 ("alpha".to_string(), "one".to_string()),
1979 ("beta".to_string(), "two".to_string()),
1980 ("gamma".to_string(), "three".to_string()),
1981 ]))
1982 .expect("build btree")
1983 .expect("non-empty root")
1984 }
1985
1986 fn build_deep_test_tree(store: &HashtreeStore) -> Cid {
1987 let index = BTree::new(store.store_arc(), BTreeOptions { order: Some(4) });
1988 let entries: Vec<_> = (0..256)
1989 .map(|index| (format!("key-{index:04}"), format!("value-{index:04}")))
1990 .collect();
1991 sync_block_on(index.build(entries))
1992 .expect("build deep btree")
1993 .expect("non-empty deep root")
1994 }
1995
1996 fn build_generated_encrypted_tree(
1997 store: &HashtreeStore,
1998 namespace: &str,
1999 entry_count: usize,
2000 ) -> Cid {
2001 let index = BTree::new(store.store_arc(), BTreeOptions { order: Some(4) });
2002 let entries: Vec<_> = (0..entry_count)
2003 .map(|index| {
2004 (
2005 format!("{namespace}-key-{index:04}"),
2006 format!("{namespace}-value-{index:04}"),
2007 )
2008 })
2009 .collect();
2010 let root = sync_block_on(index.build(entries))
2011 .expect("build generated encrypted btree")
2012 .expect("non-empty generated encrypted root");
2013 assert!(
2014 root.key.is_some(),
2015 "generated retention coverage must exercise a full encrypted CID"
2016 );
2017 root
2018 }
2019
2020 fn generated_tree_hashes(store: &HashtreeStore, root: &Cid) -> HashSet<Hash> {
2021 let tree = HashTree::new(HashTreeConfig::new(store.store_arc()).public());
2022 let hashes = sync_block_on(store.collect_tree_hashes(&tree, root, true))
2023 .expect("collect generated encrypted DAG");
2024 assert!(hashes.len() > 2, "generated DAG must contain descendants");
2025 hashes
2026 }
2027
2028 fn publish_generated_profile_repair_lease(store: &HashtreeStore, label: &str, root: &Cid) {
2029 let lease = ProfileRepairRetentionLease {
2030 format: PROFILE_REPAIR_RETENTION_LEASE_FORMAT.to_string(),
2031 authority_sha256: to_hex(&hashtree_core::sha256(root.to_string().as_bytes())),
2032 roots: BTreeMap::from([(label.to_string(), root.to_string())]),
2033 };
2034 let lease_bytes = lease.canonical_bytes().expect("canonical generated lease");
2035 let publication = store
2036 .acquire_profile_repair_retention_publication_guard()
2037 .expect("exclusive generated retention publication");
2038 let lease_path = store.profile_repair_retention_lease_path();
2039 std::fs::create_dir_all(lease_path.parent().expect("lease parent"))
2040 .expect("create generated lease parent");
2041 let mut lease_file = File::create(&lease_path).expect("create generated lease");
2042 lease_file
2043 .write_all(&lease_bytes)
2044 .expect("write generated lease");
2045 lease_file.sync_all().expect("sync generated lease");
2046 File::open(lease_path.parent().expect("lease parent"))
2047 .expect("open generated lease parent")
2048 .sync_all()
2049 .expect("sync generated lease parent");
2050 drop(publication);
2051 }
2052
2053 #[cfg(feature = "lmdb")]
2054 fn bounded_lmdb_store(path: &Path, max_size_bytes: u64) -> HashtreeStore {
2055 drop(
2059 super::super::LocalStore::new_unbounded_with_lmdb_map_size(
2060 path.join("blobs"),
2061 &StorageBackend::Lmdb,
2062 Some(16 * 1024 * 1024),
2063 )
2064 .expect("seed single LMDB"),
2065 );
2066 HashtreeStore::with_options_and_backend(
2067 path,
2068 None,
2069 max_size_bytes,
2070 true,
2071 &StorageBackend::Lmdb,
2072 )
2073 .expect("LMDB store")
2074 }
2075
2076 #[cfg(feature = "lmdb")]
2077 fn put_ordered_hashes(store: &HashtreeStore, count: u8) -> Vec<Hash> {
2078 let mut hashes = Vec::new();
2079 for value in 1..=count {
2080 let data = [value];
2081 let hash = hashtree_core::sha256(&data);
2082 store
2083 .router
2084 .put_sync(hash, &data)
2085 .expect("put ordered test hash");
2086 hashes.push(hash);
2087 }
2088 hashes.sort_unstable();
2089 hashes
2090 }
2091
2092 #[cfg(feature = "lmdb")]
2093 #[test]
2094 fn bounded_orphan_sweep_progresses_and_preserves_all_durable_classes() {
2095 let temp_dir = TempDir::new().expect("temp dir");
2096 let store = bounded_lmdb_store(temp_dir.path(), 1024 * 1024);
2097 let hashes = put_ordered_hashes(&store, 7);
2098
2099 store.pin(&hashes[0]).expect("pin first hash");
2100 let mut tree_key = [0u8; 64];
2101 tree_key[..32].copy_from_slice(&hashes[1]);
2102 tree_key[32..].fill(99);
2103 let mut wtxn = store.env.write_txn().expect("metadata write");
2104 store
2105 .blob_trees
2106 .put(&mut wtxn, &tree_key, &())
2107 .expect("index second hash");
2108 wtxn.commit().expect("commit tree index");
2109 store
2110 .set_blob_owner(&hashes[2], &[77; 32])
2111 .expect("own third hash");
2112 let socialgraph_root = build_test_tree(&store);
2113 let socialgraph_hashes = generated_tree_hashes(&store, &socialgraph_root);
2114 write_root_file(
2115 &temp_dir.path().join("socialgraph/events-root.msgpack"),
2116 &socialgraph_root,
2117 );
2118
2119 let mut progress = Vec::new();
2120 loop {
2121 let page = store
2122 .evict_disposable_orphans_page(0, &HashSet::new(), 2)
2123 .expect("bounded orphan page");
2124 assert!(page.scanned <= 2, "page exceeded its candidate bound");
2125 progress.push(page);
2126 if page.sweep_complete {
2127 break;
2128 }
2129 assert!(progress.len() < 50, "bounded sweep did not converge");
2130 }
2131
2132 assert_eq!(progress.iter().map(|page| page.freed_bytes).sum::<u64>(), 4);
2133 assert!(progress.iter().all(|page| page.scanned <= 2));
2134 for hash in &hashes[..3] {
2135 assert!(store.blob_exists(hash).expect("protected blob lookup"));
2136 }
2137 for hash in &hashes[3..] {
2138 assert!(!store.blob_exists(hash).expect("orphan blob lookup"));
2139 }
2140 for hash in socialgraph_hashes {
2141 assert!(
2142 store.blob_exists(&hash).expect("socialgraph blob lookup"),
2143 "bounded sweep deleted generated socialgraph DAG hash {}",
2144 to_hex(&hash)
2145 );
2146 }
2147 }
2148
2149 #[cfg(feature = "lmdb")]
2150 #[test]
2151 fn socialgraph_root_change_unions_protection_until_sweep_boundary() {
2152 let temp_dir = TempDir::new().expect("temp dir");
2153 let store = bounded_lmdb_store(temp_dir.path(), 1024 * 1024);
2154 let old_root = build_generated_encrypted_tree(&store, "old-live-root", 24);
2155 let old_hashes = generated_tree_hashes(&store, &old_root);
2156 let root_path = temp_dir.path().join("socialgraph/events-root.msgpack");
2157 write_root_file(&root_path, &old_root);
2158
2159 let first = store
2160 .evict_disposable_orphans_page(0, &HashSet::new(), 1)
2161 .expect("first root page");
2162 assert_eq!(first.scanned, 1);
2163 assert_eq!(first.freed_bytes, 0);
2164 let new_root = build_generated_encrypted_tree(&store, "new-live-root", 24);
2165 let new_hashes = generated_tree_hashes(&store, &new_root);
2166 write_root_file(&root_path, &new_root);
2167
2168 let second = store
2169 .evict_disposable_orphans_page(0, &HashSet::new(), 1)
2170 .expect("changed root page");
2171 assert_eq!(second.scanned, 1);
2172 assert_eq!(second.freed_bytes, 0);
2173 {
2174 let state = store
2175 .orphan_scan
2176 .lock()
2177 .unwrap_or_else(std::sync::PoisonError::into_inner);
2178 assert!(old_hashes
2179 .iter()
2180 .all(|hash| state.socialgraph_protected.contains(hash)));
2181 assert!(new_hashes
2182 .iter()
2183 .all(|hash| state.socialgraph_protected.contains(hash)));
2184 }
2185
2186 let mut sweep_complete = false;
2187 for _ in 0..256 {
2188 let page = store
2189 .evict_disposable_orphans_page(0, &HashSet::new(), 1)
2190 .expect("finish unioned socialgraph sweep");
2191 assert_eq!(page.freed_bytes, 0);
2192 if page.sweep_complete {
2193 sweep_complete = true;
2194 break;
2195 }
2196 }
2197 assert!(sweep_complete, "unioned socialgraph sweep did not converge");
2198 for hash in old_hashes.iter().chain(new_hashes.iter()) {
2199 assert!(
2200 store.blob_exists(hash).expect("unioned root lookup"),
2201 "active sweep deleted union-protected DAG hash {}",
2202 to_hex(hash)
2203 );
2204 }
2205 }
2206
2207 #[cfg(feature = "lmdb")]
2208 #[test]
2209 fn immutable_profile_repair_lease_preserves_complete_generated_dag() {
2210 let temp_dir = TempDir::new().expect("temp dir");
2211 let store = bounded_lmdb_store(temp_dir.path(), 64 * 1024 * 1024);
2212 let protected_root = build_deep_test_tree(&store);
2213 let tree = HashTree::new(HashTreeConfig::new(store.store_arc()).public());
2214 let protected_hashes =
2215 sync_block_on(store.collect_tree_hashes(&tree, &protected_root, true))
2216 .expect("collect generated repair DAG");
2217 assert!(protected_hashes.len() > 2);
2218
2219 let orphan_bytes = b"unleased generated orphan";
2220 let orphan = hashtree_core::sha256(orphan_bytes);
2221 store
2222 .router
2223 .put_sync(orphan, orphan_bytes)
2224 .expect("put generated orphan");
2225
2226 let lease = ProfileRepairRetentionLease {
2227 format: PROFILE_REPAIR_RETENTION_LEASE_FORMAT.to_string(),
2228 authority_sha256: "11".repeat(32),
2229 roots: BTreeMap::from([("profile-search".to_string(), protected_root.to_string())]),
2230 };
2231 let lease_bytes = lease.canonical_bytes().expect("canonical lease");
2232 let publication = store
2233 .acquire_profile_repair_retention_publication_guard()
2234 .expect("exclusive retention publication");
2235 let lease_path = store.profile_repair_retention_lease_path();
2236 std::fs::create_dir_all(lease_path.parent().expect("lease parent"))
2237 .expect("create lease parent");
2238 let mut lease_file = File::create(&lease_path).expect("create lease");
2239 lease_file.write_all(&lease_bytes).expect("write lease");
2240 lease_file.sync_all().expect("sync lease");
2241 File::open(lease_path.parent().expect("lease parent"))
2242 .expect("open lease parent")
2243 .sync_all()
2244 .expect("sync lease parent");
2245 drop(publication);
2246
2247 loop {
2248 let page = store
2249 .evict_disposable_orphans_page(0, &HashSet::new(), 31)
2250 .expect("lease-protected orphan sweep");
2251 if page.sweep_complete {
2252 break;
2253 }
2254 }
2255
2256 for hash in protected_hashes {
2257 assert!(
2258 store.blob_exists(&hash).expect("protected blob lookup"),
2259 "lease lost generated DAG blob {}",
2260 to_hex(&hash)
2261 );
2262 }
2263 assert!(!store.blob_exists(&orphan).expect("orphan lookup"));
2264 }
2265
2266 #[cfg(feature = "lmdb")]
2267 #[test]
2268 fn garbage_collection_preserves_leased_generated_encrypted_dag() {
2269 let temp_dir = TempDir::new().expect("temp dir");
2270 let store = bounded_lmdb_store(temp_dir.path(), 64 * 1024 * 1024);
2271 let leased_root = build_generated_encrypted_tree(&store, "gc-leased", 96);
2272 let leased_hashes = generated_tree_hashes(&store, &leased_root);
2273 let current_root = build_generated_encrypted_tree(&store, "gc-current", 96);
2274 let current_hashes = generated_tree_hashes(&store, ¤t_root);
2275 write_root_file(
2276 &temp_dir.path().join("socialgraph/events-root.msgpack"),
2277 ¤t_root,
2278 );
2279 let orphan_bytes = format!("gc-unleased-orphan:{}", leased_root);
2280 let orphan_hash = hashtree_core::sha256(orphan_bytes.as_bytes());
2281 store
2282 .router
2283 .put_sync(orphan_hash, orphan_bytes.as_bytes())
2284 .expect("put generated GC orphan");
2285 publish_generated_profile_repair_lease(&store, "gc-leased", &leased_root);
2286
2287 let report = store.gc().expect("lease-aware garbage collection");
2288
2289 assert_eq!(report.deleted_dags, 1);
2290 assert!(!store.blob_exists(&orphan_hash).expect("GC orphan lookup"));
2291 for hash in leased_hashes.into_iter().chain(current_hashes) {
2292 assert!(
2293 store.blob_exists(&hash).expect("leased GC hash lookup"),
2294 "GC deleted active encrypted DAG hash {}",
2295 to_hex(&hash)
2296 );
2297 }
2298 }
2299
2300 #[test]
2301 fn garbage_collection_preserves_generated_pending_profile_root_pair_until_recovery() {
2302 let temp_dir = TempDir::new().expect("temp dir");
2303 let store = HashtreeStore::with_embedded_options(temp_dir.path(), None, 64 * 1024 * 1024)
2304 .expect("embedded generated store");
2305 let graph = crate::socialgraph::open_test_social_graph_store_with_storage(
2306 temp_dir.path(),
2307 store.store_arc(),
2308 None,
2309 )
2310 .expect("open generated social graph");
2311
2312 let mut profiles = Vec::new();
2313 let mut decisions = BTreeMap::new();
2314 for index in 0..80 {
2315 let keys = Keys::generate();
2316 let event = EventBuilder::new(
2317 Kind::Metadata,
2318 serde_json::json!({
2319 "display_name": format!("pending-profile-{index:04}")
2320 })
2321 .to_string(),
2322 )
2323 .custom_created_at(Timestamp::from_secs(index + 1))
2324 .sign_with_keys(&keys)
2325 .expect("sign generated profile");
2326 decisions.insert(event.pubkey.to_hex(), Some(1));
2327 profiles.push(event);
2328 }
2329 let expected_profile = profiles[0].clone();
2330 let prepared = graph
2331 .build_unpublished_profile_index_repair_with_frozen_distances(&profiles, &decisions)
2332 .expect("build generated replacement profile roots");
2333 let mut pending_hashes = HashSet::new();
2334 for root in [
2335 prepared
2336 .new_roots()
2337 .by_pubkey
2338 .as_ref()
2339 .expect("generated by-pubkey root"),
2340 prepared
2341 .new_roots()
2342 .search
2343 .as_ref()
2344 .expect("generated search root"),
2345 ] {
2346 pending_hashes.extend(generated_tree_hashes(&store, root));
2347 }
2348
2349 let crash = graph
2350 .crash_after_prepared_profile_root_pair_intent(&prepared)
2351 .expect_err("generated commit must stop after its durable intent");
2352 assert!(format!("{crash:#}").contains("injected interruption after durable"));
2353 let commit_path = temp_dir
2354 .path()
2355 .join("socialgraph/profile-root-pair.commit.json");
2356 assert!(commit_path.is_file(), "generated commit intent is missing");
2357 assert!(
2358 !temp_dir
2359 .path()
2360 .join("socialgraph/profiles-by-pubkey-root.msgpack")
2361 .exists(),
2362 "crash injection unexpectedly installed by-pubkey root"
2363 );
2364 assert!(
2365 !temp_dir
2366 .path()
2367 .join("socialgraph/profile-search-root.msgpack")
2368 .exists(),
2369 "crash injection unexpectedly installed search root"
2370 );
2371 drop(graph);
2372
2373 let orphan_bytes = b"pending-root-pair-unrelated-orphan";
2374 let orphan_hash = hashtree_core::sha256(orphan_bytes);
2375 store
2376 .router
2377 .put_sync(orphan_hash, orphan_bytes)
2378 .expect("put generated orphan");
2379
2380 let report = store.gc().expect("pending-commit-aware garbage collection");
2381 assert_eq!(report.deleted_dags, 1);
2382 assert!(!store.blob_exists(&orphan_hash).expect("orphan lookup"));
2383 for hash in &pending_hashes {
2384 assert!(
2385 store.blob_exists(hash).expect("pending DAG lookup"),
2386 "GC deleted pending root-pair DAG hash {}",
2387 to_hex(hash)
2388 );
2389 }
2390
2391 let recovered = crate::socialgraph::open_test_social_graph_store_with_storage(
2392 temp_dir.path(),
2393 store.store_arc(),
2394 None,
2395 )
2396 .expect("recover generated pending root pair");
2397 assert!(!commit_path.exists(), "pending commit was not recovered");
2398 assert_eq!(
2399 recovered
2400 .latest_profile_event(&expected_profile.pubkey.to_hex())
2401 .expect("read recovered profile")
2402 .expect("recovered profile is missing")
2403 .id,
2404 expected_profile.id
2405 );
2406 }
2407
2408 #[cfg(feature = "lmdb")]
2409 #[test]
2410 fn applied_root_retention_preserves_leased_generated_encrypted_dag() {
2411 let temp_dir = TempDir::new().expect("temp dir");
2412 let store = bounded_lmdb_store(temp_dir.path(), 64 * 1024 * 1024);
2413 let retained_root = build_generated_encrypted_tree(&store, "retained-target", 96);
2414 let leased_root = build_generated_encrypted_tree(&store, "retention-leased", 96);
2415 let leased_hashes = generated_tree_hashes(&store, &leased_root);
2416 let current_root = build_generated_encrypted_tree(&store, "retention-current", 96);
2417 let current_hashes = generated_tree_hashes(&store, ¤t_root);
2418 write_root_file(
2419 &temp_dir
2420 .path()
2421 .join("socialgraph/profile-search-root.msgpack"),
2422 ¤t_root,
2423 );
2424 let orphan_bytes = format!("retention-unleased-orphan:{}", leased_root);
2425 let orphan_hash = hashtree_core::sha256(orphan_bytes.as_bytes());
2426 store
2427 .router
2428 .put_sync(orphan_hash, orphan_bytes.as_bytes())
2429 .expect("put generated retention orphan");
2430 publish_generated_profile_repair_lease(&store, "retention-leased", &leased_root);
2431
2432 let report = store
2433 .retain_nostr_root(&retained_root, true)
2434 .expect("lease-aware apply retention");
2435
2436 assert_eq!(report.candidate_hashes, 1);
2437 assert_eq!(report.deleted_hashes, 1);
2438 assert!(!store
2439 .blob_exists(&orphan_hash)
2440 .expect("retention orphan lookup"));
2441 for hash in leased_hashes.into_iter().chain(current_hashes) {
2442 assert!(
2443 store
2444 .blob_exists(&hash)
2445 .expect("leased retention hash lookup"),
2446 "apply retention deleted active encrypted DAG hash {}",
2447 to_hex(&hash)
2448 );
2449 }
2450 let leased_index = BTree::new(store.store_arc(), BTreeOptions { order: Some(4) });
2451 assert_eq!(
2452 sync_block_on(leased_index.get(Some(&leased_root), "retention-leased-key-0095"))
2453 .expect("read leased encrypted DAG after retention"),
2454 Some("retention-leased-value-0095".to_string())
2455 );
2456 }
2457
2458 #[cfg(feature = "lmdb")]
2459 #[test]
2460 fn poolstore_orphan_cleanup_fails_closed_without_deleting_catalog_data() {
2461 let temp_dir = TempDir::new().expect("temp dir");
2462 let store = HashtreeStore::with_options_and_backend(
2463 temp_dir.path(),
2464 None,
2465 1024 * 1024,
2466 true,
2467 &StorageBackend::Lmdb,
2468 )
2469 .expect("fresh shared store");
2470 assert!(matches!(
2471 store.router.local_store().as_ref(),
2472 super::super::LocalStore::Pool(_)
2473 ));
2474 let data = b"durable pool catalog data";
2475 let hash = hashtree_core::sha256(data);
2476 store.router.put_sync(hash, data).expect("put pool blob");
2477
2478 let error = store
2479 .evict_disposable_orphans_to_target(0)
2480 .expect_err("PoolStore orphan deletion must fail closed");
2481 assert!(
2482 error.to_string().contains("tier-aware deletion"),
2483 "unexpected error: {error}"
2484 );
2485 assert!(store.blob_exists(&hash).expect("pool blob lookup"));
2486 }
2487
2488 #[cfg(feature = "lmdb")]
2489 #[test]
2490 fn hash_inserted_before_cursor_is_seen_after_wrap() {
2491 let temp_dir = TempDir::new().expect("temp dir");
2492 let store = bounded_lmdb_store(temp_dir.path(), 1024 * 1024);
2493 let mut entries = (1u8..=8)
2494 .map(|value| {
2495 let data = vec![value];
2496 (hashtree_core::sha256(&data), data)
2497 })
2498 .collect::<Vec<_>>();
2499 entries.sort_unstable_by_key(|(hash, _)| *hash);
2500 let (behind_cursor_hash, behind_cursor_data) = entries[0].clone();
2501 for (hash, data) in &entries[1..=3] {
2502 store
2503 .router
2504 .put_sync(*hash, data)
2505 .expect("put initial cursor fixture");
2506 }
2507
2508 let first = store
2509 .evict_disposable_orphans_page(0, &HashSet::new(), 1)
2510 .expect("first page");
2511 assert_eq!(first.scanned, 1);
2512 assert_eq!(first.freed_bytes, 1);
2513 store
2514 .router
2515 .put_sync(behind_cursor_hash, &behind_cursor_data)
2516 .expect("insert behind active cursor");
2517
2518 for _ in 0..8 {
2519 store
2520 .evict_disposable_orphans_page(0, &HashSet::new(), 1)
2521 .expect("continue wrapped sweep");
2522 if !store
2523 .blob_exists(&behind_cursor_hash)
2524 .expect("behind-cursor lookup")
2525 {
2526 return;
2527 }
2528 }
2529 panic!("hash inserted behind the active cursor was not seen after wrap");
2530 }
2531
2532 #[cfg(feature = "lmdb")]
2533 #[test]
2534 fn orphan_cleanup_keeps_indexed_tree_hashes() {
2535 let temp_dir = TempDir::new().expect("temp dir");
2536 let store = bounded_lmdb_store(temp_dir.path(), 1024);
2537 let cid = build_test_tree(&store);
2538
2539 store
2540 .index_tree(
2541 &cid.hash,
2542 "owner",
2543 Some("tree"),
2544 PRIORITY_OTHER,
2545 Some("owner/tree"),
2546 )
2547 .expect("index tree");
2548 let freed = store
2549 .evict_disposable_orphans_to_target(0)
2550 .expect("orphan cleanup");
2551
2552 assert!(freed < 1024);
2553 assert!(store.blob_exists(&cid.hash).expect("root exists"));
2554 }
2555
2556 #[test]
2557 fn list_pins_with_names_uses_indexed_tree_metadata() {
2558 let temp_dir = TempDir::new().expect("temp dir");
2559 let store = HashtreeStore::with_options(temp_dir.path(), None, 1024 * 1024).expect("store");
2560 let cid = build_test_tree(&store);
2561
2562 store.pin(&cid.hash).expect("pin tree");
2563 store
2564 .index_tree(
2565 &cid.hash,
2566 "npub1example",
2567 Some("playlist"),
2568 PRIORITY_OTHER,
2569 Some("npub1example/playlist"),
2570 )
2571 .expect("index tree");
2572
2573 let pins = store.list_pins_with_names().expect("list pins");
2574
2575 assert_eq!(pins.len(), 1);
2576 assert_eq!(pins[0].name, "npub1example/playlist");
2577 assert!(pins[0].size_bytes > 0);
2578 }
2579
2580 #[test]
2581 fn index_tree_records_multilevel_file_size_from_links() {
2582 let temp_dir = TempDir::new().expect("temp dir");
2583 let store = HashtreeStore::with_options(temp_dir.path(), None, 1024 * 1024).expect("store");
2584 let tree = HashTree::new(
2585 HashTreeConfig::new(store.store_arc())
2586 .public()
2587 .with_chunk_size(4)
2588 .with_max_links(2),
2589 );
2590 let data = (0u8..31).collect::<Vec<_>>();
2591 let (cid, size) = sync_block_on(tree.put(&data)).expect("put file");
2592
2593 store
2594 .index_tree(
2595 &cid.hash,
2596 "npub1example",
2597 Some("large-file"),
2598 PRIORITY_OTHER,
2599 Some("npub1example/large-file"),
2600 )
2601 .expect("index tree");
2602
2603 let meta = store
2604 .get_tree_meta(&cid.hash)
2605 .expect("tree meta")
2606 .expect("indexed meta");
2607 assert_eq!(size, data.len() as u64);
2608 assert_eq!(meta.total_size, data.len() as u64);
2609 }
2610
2611 #[test]
2612 fn get_tree_ref_returns_stored_root() {
2613 let temp_dir = TempDir::new().expect("temp dir");
2614 let store = HashtreeStore::with_options(temp_dir.path(), None, 1024 * 1024).expect("store");
2615 let cid = build_test_tree(&store);
2616
2617 store
2618 .index_tree(
2619 &cid.hash,
2620 "npub1example",
2621 Some("playlist"),
2622 PRIORITY_OTHER,
2623 Some("npub1example/playlist"),
2624 )
2625 .expect("index tree");
2626
2627 assert_eq!(
2628 store
2629 .get_tree_ref("npub1example/playlist")
2630 .expect("tree ref lookup"),
2631 Some(cid.hash)
2632 );
2633 }
2634
2635 #[test]
2636 fn tree_meta_deserializes_metadata_without_tree_access_field() {
2637 #[derive(Serialize)]
2638 struct LegacyTreeMeta {
2639 owner: String,
2640 name: Option<String>,
2641 synced_at: u64,
2642 total_size: u64,
2643 priority: u8,
2644 }
2645
2646 let bytes = rmp_serde::to_vec(&LegacyTreeMeta {
2647 owner: "owner".to_string(),
2648 name: Some("tree".to_string()),
2649 synced_at: 123,
2650 total_size: 456,
2651 priority: PRIORITY_OTHER,
2652 })
2653 .expect("serialize legacy metadata");
2654 let meta: TreeMeta = rmp_serde::from_slice(&bytes).expect("deserialize tree metadata");
2655
2656 assert_eq!(meta.owner, "owner");
2657 assert_eq!(meta.name.as_deref(), Some("tree"));
2658 assert_eq!(meta.synced_at, 123);
2659 assert_eq!(meta.total_size, 456);
2660 assert_eq!(meta.priority, PRIORITY_OTHER);
2661 }
2662
2663 #[test]
2664 fn tree_meta_deserializes_accidental_access_field_but_drops_it_on_write() {
2665 #[derive(Serialize)]
2666 struct AccidentalTreeMeta {
2667 owner: String,
2668 name: Option<String>,
2669 synced_at: u64,
2670 last_accessed_at: u64,
2671 total_size: u64,
2672 priority: u8,
2673 }
2674
2675 let bytes = rmp_serde::to_vec(&AccidentalTreeMeta {
2676 owner: "owner".to_string(),
2677 name: Some("tree".to_string()),
2678 synced_at: 123,
2679 last_accessed_at: 999,
2680 total_size: 456,
2681 priority: PRIORITY_OTHER,
2682 })
2683 .expect("serialize accidental metadata");
2684 let meta: TreeMeta = rmp_serde::from_slice(&bytes).expect("deserialize tree metadata");
2685 let encoded = rmp_serde::to_vec(&meta).expect("serialize current metadata");
2686 let reparsed: (String, Option<String>, u64, u64, u8) =
2687 rmp_serde::from_slice(&encoded).expect("parse current metadata shape");
2688
2689 assert_eq!(meta.owner, "owner");
2690 assert_eq!(meta.name.as_deref(), Some("tree"));
2691 assert_eq!(meta.synced_at, 123);
2692 assert_eq!(meta.total_size, 456);
2693 assert_eq!(meta.priority, PRIORITY_OTHER);
2694 assert_eq!(reparsed.0, "owner");
2695 assert_eq!(reparsed.3, 456);
2696 assert_eq!(reparsed.4, PRIORITY_OTHER);
2697 }
2698
2699 #[cfg(feature = "lmdb")]
2700 #[test]
2701 fn eviction_prefers_oldest_tree_within_priority() {
2702 let temp_dir = TempDir::new().expect("temp dir");
2703 let store = bounded_lmdb_store(temp_dir.path(), 500);
2704
2705 let hash1 = hashtree_core::sha256(&[1u8; 200]);
2706 let hash2 = hashtree_core::sha256(&[2u8; 200]);
2707 let hash3 = hashtree_core::sha256(&[3u8; 200]);
2708 store.put_blob(&[1u8; 200]).expect("put blob 1");
2709 store.put_blob(&[2u8; 200]).expect("put blob 2");
2710 store.put_blob(&[3u8; 200]).expect("put blob 3");
2711 store
2712 .index_tree(&hash1, "owner1", Some("tree1"), PRIORITY_OTHER, None)
2713 .expect("index tree 1");
2714 store
2715 .index_tree(&hash2, "owner2", Some("tree2"), PRIORITY_OTHER, None)
2716 .expect("index tree 2");
2717 store
2718 .index_tree(&hash3, "owner3", Some("tree3"), PRIORITY_OTHER, None)
2719 .expect("index tree 3");
2720
2721 let freed = store.evict_if_needed().expect("evict");
2722
2723 assert!(freed > 0);
2724 assert!(
2725 store.get_tree_meta(&hash3).expect("tree meta").is_some(),
2726 "newest tree should survive before older peers at the same priority"
2727 );
2728 }
2729
2730 #[cfg(feature = "lmdb")]
2731 #[test]
2732 fn orphan_cleanup_keeps_socialgraph_root_hashes() {
2733 let temp_dir = TempDir::new().expect("temp dir");
2734 let store = bounded_lmdb_store(temp_dir.path(), 1024);
2735 let cid = build_test_tree(&store);
2736 write_root_file(
2737 &temp_dir.path().join("socialgraph/events-root.msgpack"),
2738 &cid,
2739 );
2740
2741 let freed = store
2742 .evict_disposable_orphans_to_target(0)
2743 .expect("orphan cleanup");
2744
2745 assert!(freed < 1024);
2746 assert!(store.blob_exists(&cid.hash).expect("root exists"));
2747 }
2748
2749 #[test]
2750 fn retained_nostr_root_cleanup_is_dry_run_first_and_keeps_the_dag() {
2751 let temp_dir = TempDir::new().expect("temp dir");
2752 let store = HashtreeStore::with_options(temp_dir.path(), None, 1024 * 1024).expect("store");
2753 let root = build_deep_test_tree(&store);
2754 let orphan_bytes = b"unreachable historical index node";
2755 let orphan = hashtree_core::sha256(orphan_bytes);
2756 store.put_blob(orphan_bytes).expect("put orphan");
2757
2758 let dry_run = store
2759 .retain_nostr_root(&root, false)
2760 .expect("retention dry run");
2761 assert_eq!(dry_run.deleted_hashes, 0);
2762 assert_eq!(dry_run.candidate_hashes, 1);
2763 assert_eq!(dry_run.total_hashes, dry_run.reachable_hashes + 1);
2764 assert!(dry_run.reachable_hashes > 256);
2765 assert!(store.blob_exists(&orphan).expect("orphan exists"));
2766
2767 let applied = store
2768 .retain_nostr_root(&root, true)
2769 .expect("apply retention");
2770 assert_eq!(applied.deleted_hashes, applied.candidate_hashes);
2771 assert!(!store.blob_exists(&orphan).expect("orphan deleted"));
2772 let index = BTree::new(store.store_arc(), BTreeOptions { order: Some(8) });
2773 assert_eq!(
2774 sync_block_on(index.get(Some(&root), "key-0255")).expect("read retained index"),
2775 Some("value-0255".to_string())
2776 );
2777 }
2778}