Skip to main content

hashtree_cli/storage/
retention.rs

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
125/// Acquire the existing retention-root publication lock without creating any
126/// path or opening the application blob writer.
127///
128/// Recovery commands use this before they snapshot profile roots or inspect
129/// retained event DAGs. Requiring an existing direct regular file keeps a
130/// typo, symlink swap, or incomplete Social namespace from silently creating a
131/// different lock authority.
132pub 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/// Resource limits for validating and indexing a complete pinned DAG.
208#[derive(Debug, Clone, Copy)]
209pub struct TreeIndexLimits {
210    pub max_nodes: usize,
211    pub max_bytes: u64,
212}
213
214/// Result of atomically indexing and pinning a complete DAG.
215#[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/// Metadata for a synced tree (for eviction tracking)
255#[derive(Debug, Clone, Serialize)]
256pub struct TreeMeta {
257    /// Pubkey of tree owner
258    pub owner: String,
259    /// Tree name if known (from nostr key like "npub.../name")
260    pub name: Option<String>,
261    /// Unix timestamp when this tree was synced
262    pub synced_at: u64,
263    /// Total size of all blobs in this tree
264    pub total_size: u64,
265    /// Eviction priority: 255=own/pinned, 128=followed, 64=other
266    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/// Storage usage broken down by priority tier
377#[derive(Debug, Clone)]
378pub struct StorageByPriority {
379    /// Own/pinned trees (priority 255)
380    pub own: u64,
381    /// Followed users' trees (priority 128)
382    pub followed: u64,
383    /// Other trees (priority 64)
384    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    /// Serialize publication of the immutable Social repair retention lease
459    /// against every local retention deletion pass. Persist and fsync the
460    /// canonical lease while this guard is alive.
461    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    /// Freeze the active immutable repair lease and resolve every full-CID DAG
511    /// it names while holding the shared side of the lease-publication lock.
512    /// Keep this value alive through the complete deletion pass.
513    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    /// Retain only one Nostr event-index DAG plus explicit pins in the writable
533    /// blob store. This is intended for a dedicated, closed-writer index store.
534    ///
535    /// Nostr B-tree links mark event values as files even though the values are
536    /// direct blobs. Traversing directory/fanout links only therefore visits
537    /// every index node while avoiding millions of unnecessary event reads.
538    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    /// Return the socialgraph protection snapshot for the active bounded sweep.
704    ///
705    /// Root changes during a sweep are unioned with the existing snapshot.
706    /// This can temporarily over-protect an old DAG, but can never make either
707    /// the old or new root disposable. The union is discarded at the sweep
708    /// boundary and rebuilt from the current roots.
709    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            // Recheck all durable metadata after the deletion claim. Ownership
855            // writers serialize with the claim; the point checks also narrow
856            // the pin/index publication window without materializing either DB.
857            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    /// Number of cache/orphan cleanup leadership epochs started by this store.
1016    ///
1017    /// This is intentionally a cheap diagnostic so overload tests and operators
1018    /// can verify that concurrent cache pressure coalesces into one scanner.
1019    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    /// Pin a hash (prevent garbage collection)
1166    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    /// Unpin a hash (allow garbage collection)
1174    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    /// Check if hash is pinned
1182    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    /// List all pinned hashes (raw bytes)
1188    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    /// List all pinned hashes with names
1205    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            // Try to determine if it's a directory
1220            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    // === Tree indexing for eviction ===
1280
1281    /// Bounds complete-DAG indexing by both the configured storage budget and
1282    /// an absolute traversal-node ceiling. A zero storage budget means the raw
1283    /// store is unbounded, but authenticated requests still retain a hard cap.
1284    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    /// Validate every referenced blob, then index all descendants and pin the
1296    /// root in one LMDB transaction. No pin or index metadata is written if
1297    /// traversal, decryption, decoding, or resource validation fails.
1298    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    /// Index a tree after sync - tracks all blobs in the tree for eviction
1330    ///
1331    /// If `ref_key` is provided (e.g. "npub.../name"), it will replace any existing
1332    /// tree with that ref, allowing old versions to be evicted.
1333    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 ref_key provided, check for and unindex old version
1344        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                    // Unindex old tree (will delete orphaned blobs)
1354                    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        // (cid, count logical bytes, decode/follow tree nodes, require a tree node)
1449        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    /// Unindex a tree - removes blob-tree mappings and deletes orphaned blobs.
1564    /// Returns the number of bytes freed.
1565    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        // Walk tree and collect all blob hashes
1598        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 each blob, remove the blob-tree entry and check if orphaned
1605        for tracked_hash in &tracked_hashes {
1606            // Delete blob-tree entry (64-byte key: blob_hash ++ tree_hash)
1607            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            // Check if blob is in any other tree (prefix scan on first 32 bytes)
1613            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            // Tree retention must not delete committed Blossom data or a body
1629            // whose owner/index transaction is still in flight.
1630            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                    // Delete locally only - keep S3 as archive
1642                    self.router
1643                        .delete_local_only(tracked_hash)
1644                        .map_err(|e| anyhow::anyhow!("Failed to delete blob: {}", e))?;
1645                }
1646            }
1647        }
1648
1649        // Delete tree metadata
1650        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    /// Get tree metadata
1660    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    /// List all indexed trees
1684    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    /// Get total tracked storage size (sum of all tree_meta.total_size)
1702    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    /// Get evictable trees sorted by (priority ASC, synced_at ASC).
1717    ///
1718    /// Blob-level access and raw LRU order live in the storage adapter. Indexed
1719    /// tree metadata stays cheap and does not try to summarize all descendant
1720    /// blob access on every stats or eviction pass.
1721    fn get_evictable_trees(&self) -> Result<Vec<(Hash, TreeMeta)>> {
1722        let mut trees = self.list_indexed_trees()?;
1723
1724        // Sort by priority (lower first), then by age.
1725        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    /// Run eviction if storage is over quota
1734    /// Returns bytes freed
1735    ///
1736    /// Eviction order:
1737    /// 1. Orphaned blobs (not in any indexed tree and not pinned)
1738    /// 2. Trees by priority (lowest first) and access age (least recent first)
1739    pub fn evict_if_needed(&self) -> Result<u64> {
1740        // Get storage used by the canonical writable store.
1741        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        // Target 90% of max to avoid constant eviction
1752        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        // Phase 1: Evict orphaned blobs (not in any tree and not pinned)
1765        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            // Do not evict indexed trees merely because the current bounded
1779            // orphan page was protected. Finish one complete orphan sweep
1780            // before escalating to durable tree policy.
1781            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        // Check if we're now under target
1799        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        // Phase 2: Evict trees by priority (lowest first) and access age (least recent first)
1815        // Own trees CAN be evicted (just last), but PINNED trees are never evicted
1816        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            // Never evict pinned trees
1830            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    /// Get the maximum storage size in bytes
1866    pub fn max_size_bytes(&self) -> u64 {
1867        self.max_size_bytes
1868    }
1869
1870    /// Get storage usage by priority tier
1871    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    /// Get storage statistics
1899    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        // Seed the legacy single-store path so the shared-layout opener does
2056        // not create a fresh PoolStore. Bounded orphan deletion is
2057        // intentionally LMDB-only until PoolStore has a remove-hot-copy API.
2058        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, &current_root);
2275        write_root_file(
2276            &temp_dir.path().join("socialgraph/events-root.msgpack"),
2277            &current_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, &current_root);
2418        write_root_file(
2419            &temp_dir
2420                .path()
2421                .join("socialgraph/profile-search-root.msgpack"),
2422            &current_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}