Skip to main content

hashtree_lmdb/
lib.rs

1//! LMDB-backed content-addressed blob storage.
2
3use async_trait::async_trait;
4use hashtree_core::store::{PutManyReport, Store, StoreError, StoreStats};
5use hashtree_core::{to_hex, types::Hash};
6use heed::types::*;
7use heed::{Database, EnvFlags, EnvOpenOptions, Error as HeedError, MdbError, PutFlags};
8use std::collections::HashSet;
9use std::fs::{self, File};
10use std::io::{Read, Seek, SeekFrom, Write};
11use std::path::{Path, PathBuf};
12use std::sync::atomic::{AtomicU64, Ordering};
13use std::time::{SystemTime, UNIX_EPOCH};
14
15// Re-export sha256 for convenience
16pub use hashtree_core::hash::sha256 as compute_sha256;
17
18#[cfg(target_pointer_width = "64")]
19const DEFAULT_MAP_SIZE: usize = 10 * 1024 * 1024 * 1024;
20#[cfg(target_pointer_width = "32")]
21const DEFAULT_MAP_SIZE: usize = 1024 * 1024 * 1024;
22const DEFAULT_MAX_READERS: u32 = 1024;
23const DATABASE_COUNT: u32 = 5;
24const REOPEN_HEADROOM_BYTES: u64 = 64 * 1024 * 1024;
25const EVICTION_BATCH_TARGET_BYTES: u64 = 256 * 1024 * 1024;
26const EVICTION_BATCH_MAX_ITEMS: usize = 4096;
27const LEGACY_BLOB_META_BYTES: usize = 16;
28const BLOB_META_BYTES: usize = 24;
29const ORDER_KEY_BYTES: usize = 40;
30const PIN_COUNT_BYTES: usize = 4;
31const STORE_TOTALS_BYTES: usize = 32;
32const STORE_TOTALS_KEY: &[u8] = b"totals";
33const NEXT_ORDER_BYTES: usize = 8;
34const NEXT_ORDER_KEY: &[u8] = b"next_order";
35const NEXT_ORDER_MIGRATION_FLOOR: u64 = 1 << 48;
36const EXTERNAL_BLOB_MARKER_PREFIX: &[u8] = b"\0hashtree-lmdb-external-blob-v1\0";
37const EXTERNAL_PACK_MARKER_PREFIX: &[u8] = b"\0hashtree-lmdb-external-pack-v1\0";
38const EXTERNAL_PACK_RESERVED_MARKER_PREFIX: &[u8] = b"\0hashtree-lmdb-external-pack-reserved-v1\0";
39const LMDB_NO_READ_AHEAD_ENV: &str = "HTREE_LMDB_NO_READ_AHEAD";
40const LMDB_NO_SYNC_ENV: &str = "HTREE_LMDB_NO_SYNC";
41const LMDB_NO_META_SYNC_ENV: &str = "HTREE_LMDB_NO_META_SYNC";
42const LMDB_EXTERNAL_BLOB_MIN_BYTES_ENV: &str = "HTREE_LMDB_EXTERNAL_BLOB_MIN_BYTES";
43const LMDB_EXTERNAL_BLOB_DIR_ENV: &str = "HTREE_LMDB_EXTERNAL_BLOB_DIR";
44const LMDB_EXTERNAL_BLOB_SYNC_ENV: &str = "HTREE_LMDB_EXTERNAL_BLOB_SYNC";
45const LMDB_EXTERNAL_BLOB_PACK_TARGET_BYTES_ENV: &str = "HTREE_LMDB_EXTERNAL_BLOB_PACK_TARGET_BYTES";
46static EXTERNAL_PACK_COUNTER: AtomicU64 = AtomicU64::new(0);
47
48#[derive(Debug, Clone, Copy)]
49struct BlobMeta {
50    order: u64,
51    size: u64,
52    last_accessed_at: u64,
53}
54
55#[derive(Debug, Default, Clone, Copy)]
56struct StoreTotals {
57    count: u64,
58    total_bytes: u64,
59    pinned_count: u64,
60    pinned_bytes: u64,
61}
62
63#[derive(Debug, Clone)]
64struct ExternalBlobConfig {
65    base_path: PathBuf,
66    min_bytes: usize,
67    sync: bool,
68    pack_target_bytes: Option<usize>,
69}
70
71/// Options for storing larger LMDB blobs in normal files.
72#[derive(Debug, Clone)]
73pub struct ExternalBlobOptions {
74    pub base_path: PathBuf,
75    pub min_bytes: usize,
76    pub sync: bool,
77    pub pack_target_bytes: Option<usize>,
78}
79
80impl ExternalBlobOptions {
81    /// Build external-blob options from the standard hashtree LMDB environment.
82    pub fn from_env(env_path: &Path) -> Option<Self> {
83        ExternalBlobConfig::from_env(env_path).map(Into::into)
84    }
85
86    /// Override only the base directory, keeping the remaining env-derived knobs.
87    pub fn with_base_path(mut self, base_path: PathBuf) -> Self {
88        self.base_path = base_path;
89        self
90    }
91}
92
93impl ExternalBlobConfig {
94    fn from_env(env_path: &Path) -> Option<Self> {
95        let min_bytes = std::env::var(LMDB_EXTERNAL_BLOB_MIN_BYTES_ENV)
96            .ok()
97            .and_then(|value| value.parse::<usize>().ok())
98            .filter(|value| *value > 0)?;
99        let base_path = std::env::var(LMDB_EXTERNAL_BLOB_DIR_ENV)
100            .ok()
101            .filter(|value| !value.trim().is_empty())
102            .map(PathBuf::from)
103            .unwrap_or_else(|| env_path.join("external-blobs"));
104        let sync = env_bool(LMDB_EXTERNAL_BLOB_SYNC_ENV).unwrap_or(true);
105        let pack_target_bytes = std::env::var(LMDB_EXTERNAL_BLOB_PACK_TARGET_BYTES_ENV)
106            .ok()
107            .and_then(|value| value.parse::<usize>().ok())
108            .filter(|value| *value > 0);
109        Some(Self {
110            base_path,
111            min_bytes,
112            sync,
113            pack_target_bytes,
114        })
115    }
116}
117
118impl From<ExternalBlobConfig> for ExternalBlobOptions {
119    fn from(config: ExternalBlobConfig) -> Self {
120        Self {
121            base_path: config.base_path,
122            min_bytes: config.min_bytes,
123            sync: config.sync,
124            pack_target_bytes: config.pack_target_bytes,
125        }
126    }
127}
128
129impl From<ExternalBlobOptions> for ExternalBlobConfig {
130    fn from(options: ExternalBlobOptions) -> Self {
131        Self {
132            base_path: options.base_path,
133            min_bytes: options.min_bytes,
134            sync: options.sync,
135            pack_target_bytes: options.pack_target_bytes,
136        }
137    }
138}
139
140#[derive(Debug, Clone)]
141struct ExternalPackRef {
142    name: String,
143    offset: u64,
144    len: u64,
145}
146
147/// LMDB-backed blob store implementing hashtree's Store trait.
148pub struct LmdbBlobStore {
149    env: heed::Env,
150    /// Maps SHA256 hash (32 bytes) → blob data
151    blobs: Database<Bytes, Bytes>,
152    /// Maps SHA256 hash (32 bytes) → [order: u64][size: u64]
153    metadata: Database<Bytes, Bytes>,
154    /// Maps [order: u64][hash: 32 bytes] → ()
155    eviction_order: Database<Bytes, Unit>,
156    /// Maps SHA256 hash (32 bytes) → pin count (u32)
157    pins: Database<Bytes, Bytes>,
158    /// Small aggregate counters used by quota checks and stats.
159    stats: Database<Bytes, Bytes>,
160    max_bytes: AtomicU64,
161    current_bytes: AtomicU64,
162    next_order: AtomicU64,
163    external_blobs: Option<ExternalBlobConfig>,
164}
165
166impl LmdbBlobStore {
167    /// Open or create an LMDB blob store at the given path.
168    pub fn new<P: AsRef<Path>>(path: P) -> Result<Self, StoreError> {
169        Self::with_map_size(path, DEFAULT_MAP_SIZE)
170    }
171
172    /// Open or create with a maximum logical storage size.
173    pub fn with_max_bytes<P: AsRef<Path>>(path: P, max_bytes: u64) -> Result<Self, StoreError> {
174        Self::with_max_bytes_and_external_blob_options(
175            path,
176            max_bytes,
177            ExternalBlobOptions::from_env,
178        )
179    }
180
181    /// Open or create with a maximum logical storage size and custom external-blob options.
182    pub fn with_max_bytes_and_external_blob_options<P, F>(
183        path: P,
184        max_bytes: u64,
185        external_blobs: F,
186    ) -> Result<Self, StoreError>
187    where
188        P: AsRef<Path>,
189        F: FnOnce(&Path) -> Option<ExternalBlobOptions>,
190    {
191        let path_ref = path.as_ref();
192        let existing_map_size = std::fs::metadata(path_ref.join("data.mdb"))
193            .map(|metadata| metadata.len())
194            .unwrap_or(0);
195        let existing_headroom = if existing_map_size == 0 {
196            0
197        } else {
198            existing_map_size
199                .saturating_div(10)
200                .max(REOPEN_HEADROOM_BYTES)
201        };
202        let requested_map_size = usize::try_from(Self::align_map_size_bytes(
203            existing_map_size
204                .saturating_add(existing_headroom)
205                .max(max_bytes),
206        ))
207        .unwrap_or(usize::MAX)
208        .max(DEFAULT_MAP_SIZE);
209        let external_blobs = external_blobs(path_ref);
210        let store = Self::with_map_size_and_external_blob_options(
211            path_ref,
212            requested_map_size,
213            external_blobs,
214        )?;
215        store.max_bytes.store(max_bytes, Ordering::Relaxed);
216        let current = store.current_bytes.load(Ordering::Relaxed);
217        if max_bytes > 0 && current > max_bytes {
218            let target = max_bytes.saturating_mul(9) / 10;
219            store.evict_to_target(current, target)?;
220        }
221        Ok(store)
222    }
223
224    /// Open or create with the default map size and explicit external-blob options.
225    pub fn with_external_blob_options<P: AsRef<Path>>(
226        path: P,
227        external_blobs: Option<ExternalBlobOptions>,
228    ) -> Result<Self, StoreError> {
229        Self::with_map_size_and_external_blob_options(path, DEFAULT_MAP_SIZE, external_blobs)
230    }
231
232    fn align_map_size_bytes(bytes: u64) -> u64 {
233        if bytes == 0 {
234            return 0;
235        }
236        let page_size = (page_size::get() as u64).max(4096);
237        let remainder = bytes % page_size;
238        if remainder == 0 {
239            bytes
240        } else {
241            bytes.saturating_add(page_size - remainder)
242        }
243    }
244
245    /// Open or create with custom map size.
246    pub fn with_map_size<P: AsRef<Path>>(path: P, map_size: usize) -> Result<Self, StoreError> {
247        Self::with_map_size_and_settings(
248            path,
249            map_size,
250            env_flags_from_env(),
251            ExternalBlobConfig::from_env,
252        )
253    }
254
255    /// Open or create with custom map size and explicit external-blob options.
256    pub fn with_map_size_and_external_blob_options<P: AsRef<Path>>(
257        path: P,
258        map_size: usize,
259        external_blobs: Option<ExternalBlobOptions>,
260    ) -> Result<Self, StoreError> {
261        Self::with_map_size_and_settings(path, map_size, env_flags_from_env(), |_| {
262            external_blobs.map(Into::into)
263        })
264    }
265
266    fn with_map_size_and_settings<P, F>(
267        path: P,
268        map_size: usize,
269        flags: EnvFlags,
270        external_blobs: F,
271    ) -> Result<Self, StoreError>
272    where
273        P: AsRef<Path>,
274        F: FnOnce(&Path) -> Option<ExternalBlobConfig>,
275    {
276        let path_ref = path.as_ref();
277        std::fs::create_dir_all(path_ref).map_err(StoreError::Io)?;
278        let existing_map_size = std::fs::metadata(path_ref.join("data.mdb"))
279            .map(|metadata| metadata.len())
280            .unwrap_or(0);
281        let existing_headroom = if existing_map_size == 0 {
282            0
283        } else {
284            existing_map_size
285                .saturating_div(10)
286                .max(REOPEN_HEADROOM_BYTES)
287        };
288        let requested_map_size = u64::try_from(map_size).unwrap_or(u64::MAX);
289        let map_size = usize::try_from(Self::align_map_size_bytes(
290            requested_map_size.max(existing_map_size.saturating_add(existing_headroom)),
291        ))
292        .unwrap_or(usize::MAX);
293
294        let mut env_options = EnvOpenOptions::new();
295        env_options
296            .map_size(map_size)
297            .max_dbs(DATABASE_COUNT)
298            .max_readers(DEFAULT_MAX_READERS);
299        unsafe {
300            env_options.flags(flags);
301        }
302        let env = unsafe { env_options.open(path_ref).map_err(map_heed_error)? };
303        let _ = env.clear_stale_readers();
304        if env.info().map_size < map_size {
305            unsafe { env.resize(map_size) }.map_err(map_heed_error)?;
306        }
307
308        let mut wtxn = env.write_txn().map_err(map_heed_error)?;
309        let blobs = env
310            .create_database(&mut wtxn, Some("blobs"))
311            .map_err(map_heed_error)?;
312        let metadata = env
313            .create_database(&mut wtxn, Some("metadata"))
314            .map_err(map_heed_error)?;
315        let eviction_order = env
316            .create_database(&mut wtxn, Some("eviction_order"))
317            .map_err(map_heed_error)?;
318        let pins = env
319            .create_database(&mut wtxn, Some("pins"))
320            .map_err(map_heed_error)?;
321        let stats = env
322            .create_database(&mut wtxn, Some("stats"))
323            .map_err(map_heed_error)?;
324        wtxn.commit().map_err(map_heed_error)?;
325
326        let totals = Self::load_or_rebuild_totals(&env, metadata, pins, stats)?;
327        let next_order = Self::load_or_initialize_next_order(&env, stats, totals.count > 0)?;
328
329        Ok(Self {
330            env,
331            blobs,
332            metadata,
333            eviction_order,
334            pins,
335            stats,
336            max_bytes: AtomicU64::new(0),
337            current_bytes: AtomicU64::new(totals.total_bytes),
338            next_order: AtomicU64::new(next_order),
339            external_blobs: external_blobs(path_ref),
340        })
341    }
342
343    fn load_or_rebuild_totals(
344        env: &heed::Env,
345        metadata: Database<Bytes, Bytes>,
346        pins: Database<Bytes, Bytes>,
347        stats: Database<Bytes, Bytes>,
348    ) -> Result<StoreTotals, StoreError> {
349        {
350            let rtxn = env.read_txn().map_err(map_heed_error)?;
351            if let Some(totals) = Self::read_store_totals(stats, &rtxn)? {
352                return Ok(totals);
353            }
354        }
355
356        let mut totals = StoreTotals::default();
357        {
358            let rtxn = env.read_txn().map_err(map_heed_error)?;
359            for item in metadata.iter(&rtxn).map_err(map_heed_error)? {
360                let (_, bytes) = item.map_err(map_heed_error)?;
361                let meta = Self::decode_blob_meta(bytes)?;
362                totals.count = totals.count.saturating_add(1);
363                totals.total_bytes = totals.total_bytes.saturating_add(meta.size);
364            }
365            for item in pins.iter(&rtxn).map_err(map_heed_error)? {
366                let (hash, pin_bytes) = item.map_err(map_heed_error)?;
367                let count = Self::decode_pin_count(pin_bytes)?;
368                if count == 0 {
369                    continue;
370                }
371                let Some(meta_bytes) = metadata.get(&rtxn, hash).map_err(map_heed_error)? else {
372                    continue;
373                };
374                let meta = Self::decode_blob_meta(meta_bytes)?;
375                totals.pinned_count = totals.pinned_count.saturating_add(1);
376                totals.pinned_bytes = totals.pinned_bytes.saturating_add(meta.size);
377            }
378        }
379
380        let mut wtxn = env.write_txn().map_err(map_heed_error)?;
381        Self::write_store_totals(stats, &mut wtxn, totals).map_err(map_heed_error)?;
382        wtxn.commit().map_err(map_heed_error)?;
383        Ok(totals)
384    }
385
386    fn load_or_initialize_next_order(
387        env: &heed::Env,
388        stats: Database<Bytes, Bytes>,
389        has_existing_blobs: bool,
390    ) -> Result<u64, StoreError> {
391        {
392            let rtxn = env.read_txn().map_err(map_heed_error)?;
393            if let Some(next_order) = Self::read_next_order(stats, &rtxn)? {
394                return Ok(next_order);
395            }
396        }
397
398        let next_order = if has_existing_blobs {
399            Self::legacy_next_order_floor()
400        } else {
401            0
402        };
403        let mut wtxn = env.write_txn().map_err(map_heed_error)?;
404        let next_order = Self::read_next_order_lossy(stats, &wtxn)
405            .map_err(map_heed_error)?
406            .unwrap_or(next_order);
407        Self::write_next_order(stats, &mut wtxn, next_order).map_err(map_heed_error)?;
408        wtxn.commit().map_err(map_heed_error)?;
409        Ok(next_order)
410    }
411
412    fn legacy_next_order_floor() -> u64 {
413        let micros = SystemTime::now()
414            .duration_since(UNIX_EPOCH)
415            .ok()
416            .and_then(|duration| u64::try_from(duration.as_micros()).ok())
417            .unwrap_or(NEXT_ORDER_MIGRATION_FLOOR);
418        micros.max(NEXT_ORDER_MIGRATION_FLOOR)
419    }
420
421    fn read_store_totals(
422        stats: Database<Bytes, Bytes>,
423        txn: &heed::RoTxn,
424    ) -> Result<Option<StoreTotals>, StoreError> {
425        stats
426            .get(txn, STORE_TOTALS_KEY)
427            .map_err(map_heed_error)?
428            .map(Self::decode_store_totals)
429            .transpose()
430    }
431
432    fn read_store_totals_lossy(
433        stats: Database<Bytes, Bytes>,
434        txn: &heed::RoTxn,
435    ) -> std::result::Result<StoreTotals, HeedError> {
436        Ok(stats
437            .get(txn, STORE_TOTALS_KEY)?
438            .and_then(Self::decode_store_totals_lossy)
439            .unwrap_or_default())
440    }
441
442    fn write_store_totals(
443        stats: Database<Bytes, Bytes>,
444        txn: &mut heed::RwTxn,
445        totals: StoreTotals,
446    ) -> std::result::Result<(), HeedError> {
447        let encoded = Self::encode_store_totals(totals);
448        stats.put(txn, STORE_TOTALS_KEY, &encoded)
449    }
450
451    fn read_next_order(
452        stats: Database<Bytes, Bytes>,
453        txn: &heed::RoTxn,
454    ) -> Result<Option<u64>, StoreError> {
455        stats
456            .get(txn, NEXT_ORDER_KEY)
457            .map_err(map_heed_error)?
458            .map(Self::decode_next_order)
459            .transpose()
460    }
461
462    fn read_next_order_lossy(
463        stats: Database<Bytes, Bytes>,
464        txn: &heed::RoTxn,
465    ) -> std::result::Result<Option<u64>, HeedError> {
466        Ok(stats
467            .get(txn, NEXT_ORDER_KEY)?
468            .and_then(Self::decode_next_order_lossy))
469    }
470
471    fn write_next_order(
472        stats: Database<Bytes, Bytes>,
473        txn: &mut heed::RwTxn,
474        next_order: u64,
475    ) -> std::result::Result<(), HeedError> {
476        let encoded = Self::encode_next_order(next_order);
477        stats.put(txn, NEXT_ORDER_KEY, &encoded)
478    }
479
480    fn allocate_order_range_in_txn(
481        &self,
482        txn: &mut heed::RwTxn,
483        count: usize,
484    ) -> std::result::Result<u64, HeedError> {
485        let start = Self::read_next_order_lossy(self.stats, txn)?
486            .unwrap_or_else(|| self.next_order.load(Ordering::Relaxed));
487        let end = start.saturating_add(count as u64);
488        Self::write_next_order(self.stats, txn, end)?;
489        self.next_order.store(end, Ordering::Relaxed);
490        Ok(start)
491    }
492
493    fn increment_totals_in_txn(
494        &self,
495        txn: &mut heed::RwTxn,
496        count: u64,
497        total_bytes: u64,
498        pinned_count: u64,
499        pinned_bytes: u64,
500    ) -> std::result::Result<(), HeedError> {
501        let mut totals = Self::read_store_totals_lossy(self.stats, txn)?;
502        totals.count = totals.count.saturating_add(count);
503        totals.total_bytes = totals.total_bytes.saturating_add(total_bytes);
504        totals.pinned_count = totals.pinned_count.saturating_add(pinned_count);
505        totals.pinned_bytes = totals.pinned_bytes.saturating_add(pinned_bytes);
506        Self::write_store_totals(self.stats, txn, totals)
507    }
508
509    fn decrement_totals_in_txn(
510        &self,
511        txn: &mut heed::RwTxn,
512        count: u64,
513        total_bytes: u64,
514        pinned_count: u64,
515        pinned_bytes: u64,
516    ) -> std::result::Result<(), HeedError> {
517        let mut totals = Self::read_store_totals_lossy(self.stats, txn)?;
518        totals.count = totals.count.saturating_sub(count);
519        totals.total_bytes = totals.total_bytes.saturating_sub(total_bytes);
520        totals.pinned_count = totals.pinned_count.saturating_sub(pinned_count);
521        totals.pinned_bytes = totals.pinned_bytes.saturating_sub(pinned_bytes);
522        Self::write_store_totals(self.stats, txn, totals)
523    }
524
525    /// Check if a hash exists (sync version for internal use).
526    pub fn exists(&self, hash: &Hash) -> Result<bool, StoreError> {
527        let rtxn = self
528            .env
529            .read_txn()
530            .map_err(|e| StoreError::Other(e.to_string()))?;
531
532        if self
533            .metadata
534            .get(&rtxn, hash)
535            .map_err(|e| StoreError::Other(e.to_string()))?
536            .is_some()
537        {
538            return Ok(true);
539        }
540
541        Ok(self
542            .blobs
543            .get(&rtxn, hash)
544            .map_err(|e| StoreError::Other(e.to_string()))?
545            .is_some())
546    }
547
548    fn mark_existing_hashes_in_db(
549        db: Database<Bytes, Bytes>,
550        rtxn: &heed::RoTxn,
551        sorted_hashes: &[Hash],
552        existing: &mut [bool],
553    ) -> Result<(), StoreError> {
554        debug_assert_eq!(sorted_hashes.len(), existing.len());
555        debug_assert!(sorted_hashes.windows(2).all(|pair| pair[0] <= pair[1]));
556
557        if sorted_hashes.is_empty() {
558            return Ok(());
559        }
560
561        for (index, hash) in sorted_hashes.iter().enumerate() {
562            if existing[index] {
563                continue;
564            }
565            if db.get(rtxn, hash).map_err(map_heed_error)?.is_some() {
566                existing[index] = true;
567            }
568        }
569
570        Ok(())
571    }
572
573    /// Mark which sorted hashes already exist using bounded LMDB range scans.
574    pub fn existing_hashes_in_sorted_candidates(
575        &self,
576        sorted_hashes: &[Hash],
577    ) -> Result<Vec<bool>, StoreError> {
578        let mut existing = vec![false; sorted_hashes.len()];
579        if sorted_hashes.is_empty() {
580            return Ok(existing);
581        }
582
583        let rtxn = self.env.read_txn().map_err(map_heed_error)?;
584        Self::mark_existing_hashes_in_db(self.metadata, &rtxn, sorted_hashes, &mut existing)?;
585        if existing.iter().all(|exists| *exists) {
586            return Ok(existing);
587        }
588        Self::mark_existing_hashes_in_db(self.blobs, &rtxn, sorted_hashes, &mut existing)?;
589        Ok(existing)
590    }
591
592    pub fn blob_size_sync(&self, hash: &Hash) -> Result<Option<u64>, StoreError> {
593        let rtxn = self
594            .env
595            .read_txn()
596            .map_err(|e| StoreError::Other(e.to_string()))?;
597        self.blob_size_in_txn(&rtxn, hash)
598    }
599
600    pub fn map_size_bytes(&self) -> usize {
601        self.env.info().map_size
602    }
603
604    pub fn force_sync(&self) -> Result<(), StoreError> {
605        self.env.force_sync().map_err(map_heed_error)
606    }
607
608    fn evict_for_write_pressure(&self, incoming_bytes: u64) -> Result<u64, StoreError> {
609        let current = self.current_bytes.load(Ordering::Relaxed);
610        if current == 0 {
611            return Ok(0);
612        }
613
614        let headroom = incoming_bytes.max(current / 10).max(1);
615        let target = current.saturating_sub(headroom);
616        self.evict_to_target(current, target)
617    }
618
619    fn enforce_max_bytes_after_insert(&self, inserted_bytes: u64) -> Result<u64, StoreError> {
620        let max = self.max_bytes.load(Ordering::Relaxed);
621        if max == 0 || inserted_bytes == 0 {
622            return Ok(0);
623        }
624
625        let current = self.current_bytes.load(Ordering::Relaxed);
626        if current <= max {
627            return Ok(0);
628        }
629
630        let target = if inserted_bytes >= max {
631            inserted_bytes
632        } else {
633            max.saturating_mul(9)
634                .saturating_div(10)
635                .saturating_add(inserted_bytes)
636                .min(max)
637        };
638        self.evict_to_target(current, target)
639    }
640
641    fn write_metadata_for_inserted_blobs(
642        &self,
643        wtxn: &mut heed::RwTxn,
644        inserted_entries: &[(Hash, u64)],
645    ) -> Result<(), StoreError> {
646        if inserted_entries.is_empty() {
647            return Ok(());
648        }
649
650        let now = unix_timestamp_now();
651        let mut inserted_bytes = 0u64;
652        let mut inserted_pinned = 0u64;
653        let mut inserted_pinned_bytes = 0u64;
654        let mut order = self
655            .allocate_order_range_in_txn(wtxn, inserted_entries.len())
656            .map_err(map_heed_error)?;
657
658        for (hash, data_len) in inserted_entries {
659            let meta = Self::encode_blob_meta(BlobMeta {
660                order,
661                size: *data_len,
662                last_accessed_at: now,
663            });
664            let order_key = Self::encode_order_key(order, hash);
665            order = order.saturating_add(1);
666            self.metadata
667                .put(wtxn, hash, &meta)
668                .map_err(map_heed_error)?;
669            self.eviction_order
670                .put(wtxn, &order_key, &())
671                .map_err(map_heed_error)?;
672            inserted_bytes = inserted_bytes.saturating_add(*data_len);
673            if self
674                .read_pin_count_lossy(wtxn, hash)
675                .map_err(map_heed_error)?
676                > 0
677            {
678                inserted_pinned = inserted_pinned.saturating_add(1);
679                inserted_pinned_bytes = inserted_pinned_bytes.saturating_add(*data_len);
680            }
681        }
682
683        self.increment_totals_in_txn(
684            wtxn,
685            inserted_entries.len() as u64,
686            inserted_bytes,
687            inserted_pinned,
688            inserted_pinned_bytes,
689        )
690        .map_err(map_heed_error)?;
691        Ok(())
692    }
693
694    fn put_sync_attempt(&self, hash: Hash, data: &[u8]) -> Result<bool, StoreError> {
695        let mut wtxn = self.env.write_txn().map_err(map_heed_error)?;
696        let external_config = self
697            .external_blobs
698            .as_ref()
699            .filter(|config| data.len() >= config.min_bytes);
700        let external_marker = external_config.map(|_| Self::external_blob_ref(&hash));
701        let value = external_marker.as_deref().unwrap_or(data);
702
703        match self
704            .blobs
705            .put_with_flags(&mut wtxn, PutFlags::NO_OVERWRITE, &hash, value)
706        {
707            Ok(()) => {}
708            Err(HeedError::Mdb(MdbError::KeyExist)) => return Ok(false),
709            Err(err) => return Err(map_heed_error(err)),
710        }
711
712        if let Some(config) = external_config {
713            self.write_external_blob(&hash, data, config)?;
714        }
715        self.write_metadata_for_inserted_blobs(&mut wtxn, &[(hash, data.len() as u64)])?;
716        wtxn.commit().map_err(map_heed_error)?;
717        Ok(true)
718    }
719
720    fn put_many_sync_attempt(
721        &self,
722        total: usize,
723        items: &[(Hash, &[u8])],
724    ) -> Result<PutManyReport, StoreError> {
725        let mut wtxn = self.env.write_txn().map_err(map_heed_error)?;
726        let mut report = PutManyReport {
727            total,
728            ..PutManyReport::default()
729        };
730        let mut inserted_entries: Vec<(Hash, u64)> = Vec::new();
731        let mut external_blobs: Vec<(Hash, &[u8])> = Vec::new();
732        let mut external_pack_entries: Vec<(usize, Hash, &[u8])> = Vec::new();
733        let external_config = self.external_blobs.as_ref();
734
735        for (hash, data) in items {
736            let external = external_config.filter(|config| data.len() >= config.min_bytes);
737            let pack_external = external.is_some_and(|config| config.pack_target_bytes.is_some());
738            let reserved_marker;
739            let external_marker;
740            let value = if pack_external {
741                reserved_marker = Self::external_pack_reserved_ref(hash);
742                reserved_marker.as_slice()
743            } else if external.is_some() {
744                external_marker = Self::external_blob_ref(hash);
745                external_marker.as_slice()
746            } else {
747                *data
748            };
749
750            match self
751                .blobs
752                .put_with_flags(&mut wtxn, PutFlags::NO_OVERWRITE, hash, value)
753            {
754                Ok(()) => {}
755                Err(HeedError::Mdb(MdbError::KeyExist)) => continue,
756                Err(err) => return Err(map_heed_error(err)),
757            }
758
759            let inserted_index = inserted_entries.len();
760            let data_len = data.len() as u64;
761            inserted_entries.push((*hash, data_len));
762            report.inserted = report.inserted.saturating_add(1);
763            report.inserted_bytes = report.inserted_bytes.saturating_add(data_len);
764            report.inserted_hashes.push(*hash);
765
766            if pack_external {
767                external_pack_entries.push((inserted_index, *hash, *data));
768            } else if external.is_some() {
769                external_blobs.push((*hash, *data));
770            }
771        }
772
773        if inserted_entries.is_empty() {
774            return Ok(report);
775        }
776
777        if let Some(config) = external_config {
778            for (hash, data) in external_blobs {
779                self.write_external_blob(&hash, data, config)?;
780            }
781            if let Some(pack_target_bytes) = config.pack_target_bytes {
782                for (inserted_index, marker) in self.write_external_blob_packs(
783                    &external_pack_entries,
784                    config,
785                    pack_target_bytes,
786                )? {
787                    let hash = inserted_entries[inserted_index].0;
788                    self.blobs
789                        .put(&mut wtxn, &hash, &marker)
790                        .map_err(map_heed_error)?;
791                }
792            }
793        }
794
795        self.write_metadata_for_inserted_blobs(&mut wtxn, &inserted_entries)?;
796
797        wtxn.commit().map_err(map_heed_error)?;
798        Ok(report)
799    }
800
801    /// Get storage statistics.
802    pub fn stats(&self) -> Result<LmdbStats, StoreError> {
803        let rtxn = self
804            .env
805            .read_txn()
806            .map_err(|e| StoreError::Other(e.to_string()))?;
807        let totals = Self::read_store_totals(self.stats, &rtxn)?.unwrap_or_default();
808
809        Ok(LmdbStats {
810            count: totals.count as usize,
811            total_bytes: totals.total_bytes,
812            pinned_count: totals.pinned_count as usize,
813            pinned_bytes: totals.pinned_bytes,
814        })
815    }
816
817    /// List all hashes in the store.
818    pub fn list(&self) -> Result<Vec<Hash>, StoreError> {
819        let rtxn = self
820            .env
821            .read_txn()
822            .map_err(|e| StoreError::Other(e.to_string()))?;
823
824        let mut hashes = Vec::new();
825        for item in self.metadata.iter(&rtxn).map_err(map_heed_error)? {
826            let (hash, _) = item.map_err(|e| StoreError::Other(e.to_string()))?;
827            let hash_arr: Hash = hash
828                .try_into()
829                .map_err(|_| StoreError::Other("invalid hash length".into()))?;
830            hashes.push(hash_arr);
831        }
832
833        if hashes.is_empty() {
834            for item in self
835                .blobs
836                .iter(&rtxn)
837                .map_err(|e| StoreError::Other(e.to_string()))?
838            {
839                let (hash, _) = item.map_err(|e| StoreError::Other(e.to_string()))?;
840                let hash_arr: Hash = hash
841                    .try_into()
842                    .map_err(|_| StoreError::Other("invalid hash length".into()))?;
843                hashes.push(hash_arr);
844            }
845        }
846
847        Ok(hashes)
848    }
849
850    /// Sync put operation (for use in sync contexts).
851    pub fn put_sync(&self, hash: Hash, data: &[u8]) -> Result<bool, StoreError> {
852        let incoming_bytes = data.len() as u64;
853
854        let mut retried_after_eviction = false;
855        loop {
856            match self.put_sync_attempt(hash, data) {
857                Ok(inserted) => {
858                    if inserted {
859                        self.current_bytes
860                            .fetch_add(incoming_bytes, Ordering::Relaxed);
861                        self.enforce_max_bytes_after_insert(incoming_bytes)?;
862                    }
863                    return Ok(inserted);
864                }
865                Err(err) if is_map_full_store_error(&err) && !retried_after_eviction => {
866                    let freed = self.evict_for_write_pressure(incoming_bytes)?;
867                    if freed == 0 {
868                        return Err(err);
869                    }
870                    retried_after_eviction = true;
871                }
872                Err(err) => return Err(err),
873            }
874        }
875    }
876
877    /// Sync batch put operation with exact insert accounting.
878    pub fn put_many_report_sync(
879        &self,
880        items: &[(Hash, Vec<u8>)],
881    ) -> Result<PutManyReport, StoreError> {
882        let total = items.len();
883        if items.is_empty() {
884            return Ok(PutManyReport::default());
885        }
886
887        let mut seen_missing = HashSet::new();
888        let write_items = items
889            .iter()
890            .filter_map(|(hash, data)| {
891                if !seen_missing.insert(*hash) {
892                    None
893                } else {
894                    Some((*hash, data.as_slice()))
895                }
896            })
897            .collect::<Vec<_>>();
898        let incoming_bytes = write_items
899            .iter()
900            .map(|(_, data)| data.len() as u64)
901            .fold(0u64, |total, size| total.saturating_add(size));
902
903        let mut retried_after_eviction = false;
904        loop {
905            match self.put_many_sync_attempt(total, &write_items) {
906                Ok(report) => {
907                    if report.inserted_bytes > 0 {
908                        self.current_bytes
909                            .fetch_add(report.inserted_bytes, Ordering::Relaxed);
910                        self.enforce_max_bytes_after_insert(report.inserted_bytes)?;
911                    }
912                    return Ok(report);
913                }
914                Err(err) if is_map_full_store_error(&err) && !retried_after_eviction => {
915                    let freed = self.evict_for_write_pressure(incoming_bytes)?;
916                    if freed == 0 {
917                        return Err(err);
918                    }
919                    retried_after_eviction = true;
920                }
921                Err(err) => return Err(err),
922            }
923        }
924    }
925
926    /// Sync batch put operation (for use in sync contexts).
927    pub fn put_many_sync(&self, items: &[(Hash, Vec<u8>)]) -> Result<usize, StoreError> {
928        self.put_many_report_sync(items)
929            .map(|report| report.inserted)
930    }
931
932    /// Sync get operation (for use in sync contexts).
933    pub fn get_sync(&self, hash: &Hash) -> Result<Option<Vec<u8>>, StoreError> {
934        let rtxn = self
935            .env
936            .read_txn()
937            .map_err(|e| StoreError::Other(e.to_string()))?;
938
939        let Some(blob) = self
940            .blobs
941            .get(&rtxn, hash)
942            .map_err(|e| StoreError::Other(e.to_string()))?
943        else {
944            return Ok(None);
945        };
946
947        self.decode_blob_value(hash, blob).map(Some)
948    }
949
950    pub fn get_range_sync(
951        &self,
952        hash: &Hash,
953        start: u64,
954        end_inclusive: u64,
955    ) -> Result<Option<Vec<u8>>, StoreError> {
956        let rtxn = self
957            .env
958            .read_txn()
959            .map_err(|e| StoreError::Other(e.to_string()))?;
960        let Some(blob) = self
961            .blobs
962            .get(&rtxn, hash)
963            .map_err(|e| StoreError::Other(e.to_string()))?
964        else {
965            return Ok(None);
966        };
967
968        if self.is_external_blob_ref(hash, blob) {
969            return self.read_external_blob_range(hash, start, end_inclusive);
970        }
971        if let Some(pack_ref) = Self::decode_external_pack_ref(blob)? {
972            return self.read_external_pack_range(&pack_ref, start, end_inclusive);
973        }
974
975        if blob.is_empty() || end_inclusive < start {
976            return Ok(Some(Vec::new()));
977        }
978        let len = blob.len() as u64;
979        if start >= len {
980            return Ok(Some(Vec::new()));
981        }
982
983        let actual_end = end_inclusive.min(len - 1);
984        let start = usize::try_from(start)
985            .map_err(|_| StoreError::Other("blob range start is too large".to_string()))?;
986        let end_exclusive = usize::try_from(actual_end.saturating_add(1))
987            .map_err(|_| StoreError::Other("blob range end is too large".to_string()))?;
988        Ok(Some(blob[start..end_exclusive].to_vec()))
989    }
990
991    pub fn touch_accessed_sync(&self, hash: &Hash, now: u64) -> Result<bool, StoreError> {
992        self.touch_many_accessed_sync(std::slice::from_ref(hash), now)
993            .map(|updated| updated > 0)
994    }
995
996    pub fn touch_many_accessed_sync(&self, hashes: &[Hash], now: u64) -> Result<usize, StoreError> {
997        if hashes.is_empty() {
998            return Ok(0);
999        }
1000
1001        let mut wtxn = self
1002            .env
1003            .write_txn()
1004            .map_err(|e| StoreError::Other(e.to_string()))?;
1005        let mut to_touch = Vec::new();
1006
1007        for hash in hashes {
1008            let meta = self
1009                .metadata
1010                .get(&wtxn, hash)
1011                .map_err(|e| StoreError::Other(e.to_string()))?
1012                .map(Self::decode_blob_meta)
1013                .transpose()?;
1014            let Some(meta) = meta else {
1015                continue;
1016            };
1017
1018            if meta.last_accessed_at >= now {
1019                continue;
1020            }
1021            to_touch.push((*hash, meta));
1022        }
1023
1024        let updated = to_touch.len();
1025        if updated == 0 {
1026            wtxn.commit()
1027                .map_err(|e| StoreError::Other(e.to_string()))?;
1028            return Ok(0);
1029        }
1030
1031        let mut order = self
1032            .allocate_order_range_in_txn(&mut wtxn, updated)
1033            .map_err(|e| StoreError::Other(e.to_string()))?;
1034
1035        for (hash, mut meta) in to_touch {
1036            let old_order_key = Self::encode_order_key(meta.order, &hash);
1037            let _ = self
1038                .eviction_order
1039                .delete(&mut wtxn, &old_order_key)
1040                .map_err(|e| StoreError::Other(e.to_string()))?;
1041
1042            meta.order = order;
1043            order = order.saturating_add(1);
1044            meta.last_accessed_at = now;
1045            let meta_bytes = Self::encode_blob_meta(meta);
1046            let new_order_key = Self::encode_order_key(meta.order, &hash);
1047            self.metadata
1048                .put(&mut wtxn, &hash, &meta_bytes)
1049                .map_err(|e| StoreError::Other(e.to_string()))?;
1050            self.eviction_order
1051                .put(&mut wtxn, &new_order_key, &())
1052                .map_err(|e| StoreError::Other(e.to_string()))?;
1053        }
1054
1055        wtxn.commit()
1056            .map_err(|e| StoreError::Other(e.to_string()))?;
1057        Ok(updated)
1058    }
1059
1060    pub fn last_accessed_at_sync(&self, hash: &Hash) -> Result<Option<u64>, StoreError> {
1061        let rtxn = self
1062            .env
1063            .read_txn()
1064            .map_err(|e| StoreError::Other(e.to_string()))?;
1065        self.metadata
1066            .get(&rtxn, hash)
1067            .map_err(|e| StoreError::Other(e.to_string()))?
1068            .map(Self::decode_blob_meta)
1069            .transpose()
1070            .map(|meta| meta.map(|meta| meta.last_accessed_at))
1071    }
1072
1073    pub fn many_last_accessed_at_sync(
1074        &self,
1075        hashes: &[Hash],
1076    ) -> Result<Vec<(Hash, u64)>, StoreError> {
1077        let rtxn = self
1078            .env
1079            .read_txn()
1080            .map_err(|e| StoreError::Other(e.to_string()))?;
1081        let mut results = Vec::new();
1082        for hash in hashes {
1083            let Some(meta) = self
1084                .metadata
1085                .get(&rtxn, hash)
1086                .map_err(|e| StoreError::Other(e.to_string()))?
1087                .map(Self::decode_blob_meta)
1088                .transpose()?
1089            else {
1090                continue;
1091            };
1092            results.push((*hash, meta.last_accessed_at));
1093        }
1094        Ok(results)
1095    }
1096
1097    /// Sync delete operation (for use in sync contexts).
1098    pub fn delete_sync(&self, hash: &Hash) -> Result<bool, StoreError> {
1099        let mut wtxn = self
1100            .env
1101            .write_txn()
1102            .map_err(|e| StoreError::Other(e.to_string()))?;
1103        let external_path = self.external_blob_path_in_txn(&wtxn, hash)?;
1104        let (existed, freed) = self.delete_blob_in_txn(&mut wtxn, hash)?;
1105
1106        wtxn.commit()
1107            .map_err(|e| StoreError::Other(e.to_string()))?;
1108        if freed > 0 {
1109            self.current_bytes.fetch_sub(freed, Ordering::Relaxed);
1110        }
1111        if existed {
1112            self.remove_external_blob_file(external_path);
1113        }
1114
1115        Ok(existed)
1116    }
1117
1118    fn pin_sync(&self, hash: &Hash) -> Result<(), StoreError> {
1119        let mut wtxn = self
1120            .env
1121            .write_txn()
1122            .map_err(|e| StoreError::Other(e.to_string()))?;
1123        let previous = self.read_pin_count(&wtxn, hash)?;
1124        let count = previous.saturating_add(1);
1125        let encoded = count.to_be_bytes();
1126        self.pins
1127            .put(&mut wtxn, hash, &encoded)
1128            .map_err(|e| StoreError::Other(e.to_string()))?;
1129        if previous == 0 {
1130            if let Some(size) = self.blob_size_in_txn(&wtxn, hash)? {
1131                self.increment_totals_in_txn(&mut wtxn, 0, 0, 1, size)
1132                    .map_err(|e| StoreError::Other(e.to_string()))?;
1133            }
1134        }
1135        wtxn.commit()
1136            .map_err(|e| StoreError::Other(e.to_string()))?;
1137        Ok(())
1138    }
1139
1140    fn unpin_sync(&self, hash: &Hash) -> Result<(), StoreError> {
1141        let mut wtxn = self
1142            .env
1143            .write_txn()
1144            .map_err(|e| StoreError::Other(e.to_string()))?;
1145        let count = self.read_pin_count(&wtxn, hash)?;
1146        if count == 1 {
1147            if let Some(size) = self.blob_size_in_txn(&wtxn, hash)? {
1148                self.decrement_totals_in_txn(&mut wtxn, 0, 0, 1, size)
1149                    .map_err(|e| StoreError::Other(e.to_string()))?;
1150            }
1151        }
1152        if count <= 1 {
1153            let _ = self
1154                .pins
1155                .delete(&mut wtxn, hash)
1156                .map_err(|e| StoreError::Other(e.to_string()))?;
1157        } else {
1158            let encoded = (count - 1).to_be_bytes();
1159            self.pins
1160                .put(&mut wtxn, hash, &encoded)
1161                .map_err(|e| StoreError::Other(e.to_string()))?;
1162        }
1163        wtxn.commit()
1164            .map_err(|e| StoreError::Other(e.to_string()))?;
1165        Ok(())
1166    }
1167
1168    fn evict_to_target(&self, current_bytes: u64, target_bytes: u64) -> Result<u64, StoreError> {
1169        if current_bytes <= target_bytes {
1170            return Ok(0);
1171        }
1172
1173        let rtxn = self
1174            .env
1175            .read_txn()
1176            .map_err(|e| StoreError::Other(e.to_string()))?;
1177        let order_keys: Vec<Vec<u8>> = self
1178            .eviction_order
1179            .iter(&rtxn)
1180            .map_err(|e| StoreError::Other(e.to_string()))?
1181            .map(|item| {
1182                item.map(|(key, _)| key.to_vec())
1183                    .map_err(|e| StoreError::Other(e.to_string()))
1184            })
1185            .collect::<Result<_, _>>()?;
1186        drop(rtxn);
1187
1188        let mut freed_total = 0u64;
1189        let to_free = current_bytes - target_bytes;
1190        let mut index = 0usize;
1191
1192        while freed_total < to_free && index < order_keys.len() {
1193            let mut wtxn = self
1194                .env
1195                .write_txn()
1196                .map_err(|e| StoreError::Other(e.to_string()))?;
1197            let mut batch_freed = 0u64;
1198            let mut batch_items = 0usize;
1199            let mut batch_external_paths = Vec::new();
1200
1201            while freed_total + batch_freed < to_free && index < order_keys.len() {
1202                let order_key = &order_keys[index];
1203                index += 1;
1204
1205                let hash = Self::decode_hash_from_order_key(order_key)?;
1206                if self.read_pin_count(&wtxn, &hash)? > 0 {
1207                    continue;
1208                }
1209
1210                let external_path = self.external_blob_path_in_txn(&wtxn, &hash)?;
1211                let (_, bytes_freed) = self.delete_blob_in_txn(&mut wtxn, &hash)?;
1212                if bytes_freed == 0 {
1213                    let _ = self
1214                        .eviction_order
1215                        .delete(&mut wtxn, order_key)
1216                        .map_err(|e| StoreError::Other(e.to_string()))?;
1217                    continue;
1218                }
1219
1220                batch_freed = batch_freed.saturating_add(bytes_freed);
1221                batch_items += 1;
1222                if let Some(path) = external_path {
1223                    batch_external_paths.push(path);
1224                }
1225
1226                if batch_freed >= EVICTION_BATCH_TARGET_BYTES
1227                    || batch_items >= EVICTION_BATCH_MAX_ITEMS
1228                {
1229                    break;
1230                }
1231            }
1232
1233            wtxn.commit()
1234                .map_err(|e| StoreError::Other(e.to_string()))?;
1235            for path in batch_external_paths {
1236                let _ = fs::remove_file(path);
1237            }
1238            if batch_freed > 0 {
1239                self.current_bytes.fetch_sub(batch_freed, Ordering::Relaxed);
1240                freed_total = freed_total.saturating_add(batch_freed);
1241            }
1242        }
1243
1244        Ok(freed_total)
1245    }
1246
1247    fn delete_blob_in_txn(
1248        &self,
1249        wtxn: &mut heed::RwTxn,
1250        hash: &Hash,
1251    ) -> Result<(bool, u64), StoreError> {
1252        let pin_count = self.read_pin_count(wtxn, hash)?;
1253        let data_len = self
1254            .blobs
1255            .get(wtxn, hash)
1256            .map_err(|e| StoreError::Other(e.to_string()))?
1257            .map(|data| data.len() as u64);
1258        let meta = self
1259            .metadata
1260            .get(wtxn, hash)
1261            .map_err(|e| StoreError::Other(e.to_string()))?
1262            .map(Self::decode_blob_meta)
1263            .transpose()?;
1264
1265        let existed = self
1266            .blobs
1267            .delete(wtxn, hash)
1268            .map_err(|e| StoreError::Other(e.to_string()))?;
1269        let _ = self
1270            .metadata
1271            .delete(wtxn, hash)
1272            .map_err(|e| StoreError::Other(e.to_string()))?;
1273        let _ = self
1274            .pins
1275            .delete(wtxn, hash)
1276            .map_err(|e| StoreError::Other(e.to_string()))?;
1277        if let Some(meta) = meta {
1278            let order_key = Self::encode_order_key(meta.order, hash);
1279            let _ = self
1280                .eviction_order
1281                .delete(wtxn, &order_key)
1282                .map_err(|e| StoreError::Other(e.to_string()))?;
1283        }
1284        let bytes_freed = meta.map(|m| m.size).or(data_len).unwrap_or(0);
1285        if existed || meta.is_some() {
1286            self.decrement_totals_in_txn(
1287                wtxn,
1288                1,
1289                bytes_freed,
1290                u64::from(pin_count > 0),
1291                if pin_count > 0 { bytes_freed } else { 0 },
1292            )
1293            .map_err(|e| StoreError::Other(e.to_string()))?;
1294        }
1295
1296        Ok((existed || meta.is_some(), bytes_freed))
1297    }
1298
1299    fn blob_size_in_txn(&self, txn: &heed::RoTxn, hash: &Hash) -> Result<Option<u64>, StoreError> {
1300        if let Some(meta) = self
1301            .metadata
1302            .get(txn, hash)
1303            .map_err(|e| StoreError::Other(e.to_string()))?
1304            .map(Self::decode_blob_meta)
1305            .transpose()?
1306        {
1307            return Ok(Some(meta.size));
1308        }
1309
1310        self.blobs
1311            .get(txn, hash)
1312            .map_err(|e| StoreError::Other(e.to_string()))
1313            .map(|data| data.map(|bytes| bytes.len() as u64))
1314    }
1315
1316    fn read_pin_count(&self, txn: &heed::RoTxn, hash: &[u8]) -> Result<u32, StoreError> {
1317        self.pins
1318            .get(txn, hash)
1319            .map_err(|e| StoreError::Other(e.to_string()))?
1320            .map(Self::decode_pin_count)
1321            .transpose()?
1322            .map_or(Ok(0), Ok)
1323    }
1324
1325    fn read_pin_count_lossy(
1326        &self,
1327        txn: &heed::RoTxn,
1328        hash: &[u8],
1329    ) -> std::result::Result<u32, HeedError> {
1330        Ok(self
1331            .pins
1332            .get(txn, hash)?
1333            .and_then(Self::decode_pin_count_lossy)
1334            .unwrap_or(0))
1335    }
1336
1337    fn encode_blob_meta(meta: BlobMeta) -> [u8; BLOB_META_BYTES] {
1338        let mut encoded = [0u8; BLOB_META_BYTES];
1339        encoded[..8].copy_from_slice(&meta.order.to_be_bytes());
1340        encoded[8..16].copy_from_slice(&meta.size.to_be_bytes());
1341        encoded[16..].copy_from_slice(&meta.last_accessed_at.to_be_bytes());
1342        encoded
1343    }
1344
1345    fn decode_blob_meta(bytes: &[u8]) -> Result<BlobMeta, StoreError> {
1346        if bytes.len() != LEGACY_BLOB_META_BYTES && bytes.len() != BLOB_META_BYTES {
1347            return Err(StoreError::Other(format!(
1348                "invalid blob metadata length: {}",
1349                bytes.len()
1350            )));
1351        }
1352        Ok(BlobMeta {
1353            order: Self::decode_order(&bytes[..8])?,
1354            size: u64::from_be_bytes(
1355                bytes[8..16]
1356                    .try_into()
1357                    .map_err(|_| StoreError::Other("invalid blob size bytes".into()))?,
1358            ),
1359            last_accessed_at: if bytes.len() >= BLOB_META_BYTES {
1360                u64::from_be_bytes(
1361                    bytes[16..24]
1362                        .try_into()
1363                        .map_err(|_| StoreError::Other("invalid blob access time bytes".into()))?,
1364                )
1365            } else {
1366                0
1367            },
1368        })
1369    }
1370
1371    fn encode_store_totals(totals: StoreTotals) -> [u8; STORE_TOTALS_BYTES] {
1372        let mut encoded = [0u8; STORE_TOTALS_BYTES];
1373        encoded[0..8].copy_from_slice(&totals.count.to_be_bytes());
1374        encoded[8..16].copy_from_slice(&totals.total_bytes.to_be_bytes());
1375        encoded[16..24].copy_from_slice(&totals.pinned_count.to_be_bytes());
1376        encoded[24..32].copy_from_slice(&totals.pinned_bytes.to_be_bytes());
1377        encoded
1378    }
1379
1380    fn decode_store_totals(bytes: &[u8]) -> Result<StoreTotals, StoreError> {
1381        Self::decode_store_totals_lossy(bytes).ok_or_else(|| {
1382            StoreError::Other(format!("invalid store totals length: {}", bytes.len()))
1383        })
1384    }
1385
1386    fn decode_store_totals_lossy(bytes: &[u8]) -> Option<StoreTotals> {
1387        if bytes.len() != STORE_TOTALS_BYTES {
1388            return None;
1389        }
1390        Some(StoreTotals {
1391            count: u64::from_be_bytes(bytes[0..8].try_into().ok()?),
1392            total_bytes: u64::from_be_bytes(bytes[8..16].try_into().ok()?),
1393            pinned_count: u64::from_be_bytes(bytes[16..24].try_into().ok()?),
1394            pinned_bytes: u64::from_be_bytes(bytes[24..32].try_into().ok()?),
1395        })
1396    }
1397
1398    fn encode_next_order(next_order: u64) -> [u8; NEXT_ORDER_BYTES] {
1399        next_order.to_be_bytes()
1400    }
1401
1402    fn decode_next_order(bytes: &[u8]) -> Result<u64, StoreError> {
1403        Self::decode_next_order_lossy(bytes)
1404            .ok_or_else(|| StoreError::Other(format!("invalid next order length: {}", bytes.len())))
1405    }
1406
1407    fn decode_next_order_lossy(bytes: &[u8]) -> Option<u64> {
1408        if bytes.len() != NEXT_ORDER_BYTES {
1409            return None;
1410        }
1411        Some(u64::from_be_bytes(bytes.try_into().ok()?))
1412    }
1413
1414    fn encode_order_key(order: u64, hash: &Hash) -> [u8; ORDER_KEY_BYTES] {
1415        let mut key = [0u8; ORDER_KEY_BYTES];
1416        key[..8].copy_from_slice(&order.to_be_bytes());
1417        key[8..].copy_from_slice(hash);
1418        key
1419    }
1420
1421    fn decode_order(bytes: &[u8]) -> Result<u64, StoreError> {
1422        if bytes.len() != 8 {
1423            return Err(StoreError::Other(format!(
1424                "invalid order length: {}",
1425                bytes.len()
1426            )));
1427        }
1428        Ok(u64::from_be_bytes(bytes.try_into().map_err(|_| {
1429            StoreError::Other("invalid order bytes".into())
1430        })?))
1431    }
1432
1433    fn decode_hash_from_order_key(bytes: &[u8]) -> Result<Hash, StoreError> {
1434        if bytes.len() != ORDER_KEY_BYTES {
1435            return Err(StoreError::Other(format!(
1436                "invalid order key length: {}",
1437                bytes.len()
1438            )));
1439        }
1440        let mut hash = [0u8; 32];
1441        hash.copy_from_slice(&bytes[8..]);
1442        Ok(hash)
1443    }
1444
1445    fn decode_pin_count(bytes: &[u8]) -> Result<u32, StoreError> {
1446        Self::decode_pin_count_lossy(bytes)
1447            .ok_or_else(|| StoreError::Other(format!("invalid pin count length: {}", bytes.len())))
1448    }
1449
1450    fn decode_pin_count_lossy(bytes: &[u8]) -> Option<u32> {
1451        if bytes.len() != PIN_COUNT_BYTES {
1452            return None;
1453        }
1454        Some(u32::from_be_bytes(bytes.try_into().ok()?))
1455    }
1456
1457    fn external_blob_ref(hash: &Hash) -> Vec<u8> {
1458        let mut marker = Vec::with_capacity(EXTERNAL_BLOB_MARKER_PREFIX.len() + hash.len());
1459        marker.extend_from_slice(EXTERNAL_BLOB_MARKER_PREFIX);
1460        marker.extend_from_slice(hash);
1461        marker
1462    }
1463
1464    fn external_pack_reserved_ref(hash: &Hash) -> Vec<u8> {
1465        let mut marker =
1466            Vec::with_capacity(EXTERNAL_PACK_RESERVED_MARKER_PREFIX.len() + hash.len());
1467        marker.extend_from_slice(EXTERNAL_PACK_RESERVED_MARKER_PREFIX);
1468        marker.extend_from_slice(hash);
1469        marker
1470    }
1471
1472    fn external_pack_blob_ref(
1473        pack_name: &str,
1474        offset: u64,
1475        len: u64,
1476    ) -> Result<Vec<u8>, StoreError> {
1477        let name_len = u16::try_from(pack_name.len())
1478            .map_err(|_| StoreError::Other("external pack file name is too long".to_string()))?;
1479        let mut marker =
1480            Vec::with_capacity(EXTERNAL_PACK_MARKER_PREFIX.len() + 2 + pack_name.len() + 16);
1481        marker.extend_from_slice(EXTERNAL_PACK_MARKER_PREFIX);
1482        marker.extend_from_slice(&name_len.to_be_bytes());
1483        marker.extend_from_slice(pack_name.as_bytes());
1484        marker.extend_from_slice(&offset.to_be_bytes());
1485        marker.extend_from_slice(&len.to_be_bytes());
1486        Ok(marker)
1487    }
1488
1489    fn is_external_blob_ref(&self, hash: &Hash, value: &[u8]) -> bool {
1490        value.len() == EXTERNAL_BLOB_MARKER_PREFIX.len() + hash.len()
1491            && value.starts_with(EXTERNAL_BLOB_MARKER_PREFIX)
1492            && &value[EXTERNAL_BLOB_MARKER_PREFIX.len()..] == hash
1493    }
1494
1495    fn decode_external_pack_ref(value: &[u8]) -> Result<Option<ExternalPackRef>, StoreError> {
1496        if !value.starts_with(EXTERNAL_PACK_MARKER_PREFIX) {
1497            return Ok(None);
1498        }
1499
1500        let rest = &value[EXTERNAL_PACK_MARKER_PREFIX.len()..];
1501        if rest.len() < 2 + 8 + 8 {
1502            return Err(StoreError::Other(
1503                "invalid external LMDB pack marker length".to_string(),
1504            ));
1505        }
1506        let name_len = u16::from_be_bytes(
1507            rest[..2]
1508                .try_into()
1509                .map_err(|_| StoreError::Other("invalid external pack name length".into()))?,
1510        ) as usize;
1511        let expected_len = 2usize.saturating_add(name_len).saturating_add(16);
1512        if rest.len() != expected_len {
1513            return Err(StoreError::Other(format!(
1514                "invalid external LMDB pack marker length: {}",
1515                value.len()
1516            )));
1517        }
1518        let name_bytes = &rest[2..2 + name_len];
1519        let name = std::str::from_utf8(name_bytes)
1520            .map_err(|_| StoreError::Other("external pack name is not UTF-8".into()))?;
1521        if name.len() < 2
1522            || name.contains('/')
1523            || name.contains('\\')
1524            || name.starts_with('.')
1525            || name.contains("..")
1526        {
1527            return Err(StoreError::Other(
1528                "external pack name is not a safe relative file name".into(),
1529            ));
1530        }
1531        let offset_start = 2 + name_len;
1532        let offset = u64::from_be_bytes(
1533            rest[offset_start..offset_start + 8]
1534                .try_into()
1535                .map_err(|_| StoreError::Other("invalid external pack offset".into()))?,
1536        );
1537        let len = u64::from_be_bytes(
1538            rest[offset_start + 8..offset_start + 16]
1539                .try_into()
1540                .map_err(|_| StoreError::Other("invalid external pack length".into()))?,
1541        );
1542        Ok(Some(ExternalPackRef {
1543            name: name.to_string(),
1544            offset,
1545            len,
1546        }))
1547    }
1548
1549    fn external_blob_path_for_config(config: &ExternalBlobConfig, hash: &Hash) -> PathBuf {
1550        let hex = to_hex(hash);
1551        config
1552            .base_path
1553            .join(&hex[..2])
1554            .join(&hex[2..4])
1555            .join(&hex[4..])
1556    }
1557
1558    fn external_blob_path(&self, hash: &Hash) -> Option<PathBuf> {
1559        self.external_blobs
1560            .as_ref()
1561            .map(|config| Self::external_blob_path_for_config(config, hash))
1562    }
1563
1564    fn external_pack_path_for_config(config: &ExternalBlobConfig, pack_name: &str) -> PathBuf {
1565        config
1566            .base_path
1567            .join("packs")
1568            .join(&pack_name[..2])
1569            .join(pack_name)
1570    }
1571
1572    fn external_pack_path(&self, pack_ref: &ExternalPackRef) -> Option<PathBuf> {
1573        self.external_blobs
1574            .as_ref()
1575            .map(|config| Self::external_pack_path_for_config(config, &pack_ref.name))
1576    }
1577
1578    fn external_pack_name(first_hash: &Hash) -> String {
1579        let hash_hex = to_hex(first_hash);
1580        let nanos = SystemTime::now()
1581            .duration_since(UNIX_EPOCH)
1582            .unwrap_or_default()
1583            .as_nanos();
1584        let counter = EXTERNAL_PACK_COUNTER.fetch_add(1, Ordering::Relaxed);
1585        format!(
1586            "{}-{:032x}-{}-{:016x}.pack",
1587            &hash_hex[..12],
1588            nanos,
1589            std::process::id(),
1590            counter
1591        )
1592    }
1593
1594    fn write_external_blob(
1595        &self,
1596        hash: &Hash,
1597        data: &[u8],
1598        config: &ExternalBlobConfig,
1599    ) -> Result<(), StoreError> {
1600        let path = Self::external_blob_path_for_config(config, hash);
1601        if path.exists() {
1602            return Ok(());
1603        }
1604
1605        let parent = path
1606            .parent()
1607            .ok_or_else(|| StoreError::Other("external blob path has no parent".to_string()))?;
1608        fs::create_dir_all(parent)?;
1609        let temp_path = unique_temp_path(&path);
1610        {
1611            let mut file = File::options()
1612                .write(true)
1613                .create_new(true)
1614                .open(&temp_path)?;
1615            file.write_all(data)?;
1616            if config.sync {
1617                file.sync_all()?;
1618            }
1619        }
1620
1621        if let Err(error) = fs::rename(&temp_path, &path) {
1622            let _ = fs::remove_file(&temp_path);
1623            return Err(error.into());
1624        }
1625        if config.sync {
1626            File::open(parent)?.sync_all()?;
1627        }
1628        Ok(())
1629    }
1630
1631    fn write_external_blob_packs(
1632        &self,
1633        entries: &[(usize, Hash, &[u8])],
1634        config: &ExternalBlobConfig,
1635        pack_target_bytes: usize,
1636    ) -> Result<Vec<(usize, Vec<u8>)>, StoreError> {
1637        let mut markers = Vec::with_capacity(entries.len());
1638        let mut pack_entries = Vec::new();
1639        let mut pack_bytes = 0usize;
1640
1641        for entry in entries {
1642            let data_len = entry.2.len();
1643            let would_exceed_target =
1644                !pack_entries.is_empty() && pack_bytes.saturating_add(data_len) > pack_target_bytes;
1645            if would_exceed_target {
1646                markers.extend(self.write_external_blob_pack(&pack_entries, config)?);
1647                pack_entries.clear();
1648                pack_bytes = 0;
1649            }
1650
1651            pack_entries.push(*entry);
1652            pack_bytes = pack_bytes.saturating_add(data_len);
1653        }
1654
1655        if !pack_entries.is_empty() {
1656            markers.extend(self.write_external_blob_pack(&pack_entries, config)?);
1657        }
1658
1659        Ok(markers)
1660    }
1661
1662    fn write_external_blob_pack(
1663        &self,
1664        entries: &[(usize, Hash, &[u8])],
1665        config: &ExternalBlobConfig,
1666    ) -> Result<Vec<(usize, Vec<u8>)>, StoreError> {
1667        if entries.is_empty() {
1668            return Ok(Vec::new());
1669        }
1670
1671        let pack_name = Self::external_pack_name(&entries[0].1);
1672        let path = Self::external_pack_path_for_config(config, &pack_name);
1673        let parent = path
1674            .parent()
1675            .ok_or_else(|| StoreError::Other("external pack path has no parent".to_string()))?;
1676        fs::create_dir_all(parent)?;
1677        let temp_path = unique_temp_path(&path);
1678        let mut markers = Vec::with_capacity(entries.len());
1679        let write_result = (|| -> Result<(), StoreError> {
1680            let mut file = File::options()
1681                .write(true)
1682                .create_new(true)
1683                .open(&temp_path)?;
1684            let mut offset = 0u64;
1685            for (index, _, data) in entries {
1686                let len = data.len() as u64;
1687                file.write_all(data)?;
1688                markers.push((
1689                    *index,
1690                    Self::external_pack_blob_ref(&pack_name, offset, len)?,
1691                ));
1692                offset = offset.saturating_add(len);
1693            }
1694            if config.sync {
1695                file.sync_all()?;
1696            }
1697            Ok(())
1698        })();
1699
1700        if let Err(error) = write_result {
1701            let _ = fs::remove_file(&temp_path);
1702            return Err(error);
1703        }
1704        if let Err(error) = fs::rename(&temp_path, &path) {
1705            let _ = fs::remove_file(&temp_path);
1706            return Err(error.into());
1707        }
1708        if config.sync {
1709            File::open(parent)?.sync_all()?;
1710        }
1711        Ok(markers)
1712    }
1713
1714    fn decode_blob_value(&self, hash: &Hash, value: &[u8]) -> Result<Vec<u8>, StoreError> {
1715        if self.is_external_blob_ref(hash, value) {
1716            let path = self.external_blob_path(hash).ok_or_else(|| {
1717                StoreError::Other(
1718                    "external LMDB blob marker found but external blobs are disabled".into(),
1719                )
1720            })?;
1721            return fs::read(path).map_err(StoreError::Io);
1722        }
1723        if let Some(pack_ref) = Self::decode_external_pack_ref(value)? {
1724            return self.read_external_pack_blob(&pack_ref);
1725        }
1726        Ok(value.to_vec())
1727    }
1728
1729    fn read_external_blob_range(
1730        &self,
1731        hash: &Hash,
1732        start: u64,
1733        end_inclusive: u64,
1734    ) -> Result<Option<Vec<u8>>, StoreError> {
1735        let path = self.external_blob_path(hash).ok_or_else(|| {
1736            StoreError::Other(
1737                "external LMDB blob marker found but external blobs are disabled".into(),
1738            )
1739        })?;
1740        let mut file = File::open(path)?;
1741        let len = file.metadata()?.len();
1742        if len == 0 || start >= len || end_inclusive < start {
1743            return Ok(Some(Vec::new()));
1744        }
1745
1746        let actual_end = end_inclusive.min(len - 1);
1747        let read_len = actual_end.saturating_sub(start).saturating_add(1);
1748        let read_len = usize::try_from(read_len)
1749            .map_err(|_| StoreError::Other("blob range is too large to read".to_string()))?;
1750        let mut data = vec![0; read_len];
1751        file.seek(SeekFrom::Start(start))?;
1752        file.read_exact(&mut data)?;
1753        Ok(Some(data))
1754    }
1755
1756    fn read_external_pack_blob(&self, pack_ref: &ExternalPackRef) -> Result<Vec<u8>, StoreError> {
1757        let read_len = usize::try_from(pack_ref.len).map_err(|_| {
1758            StoreError::Other("external pack blob is too large to read".to_string())
1759        })?;
1760        let path = self.external_pack_path(pack_ref).ok_or_else(|| {
1761            StoreError::Other(
1762                "external LMDB pack marker found but external blobs are disabled".into(),
1763            )
1764        })?;
1765        let mut file = File::open(path)?;
1766        file.seek(SeekFrom::Start(pack_ref.offset))?;
1767        let mut data = vec![0; read_len];
1768        file.read_exact(&mut data)?;
1769        Ok(data)
1770    }
1771
1772    fn read_external_pack_range(
1773        &self,
1774        pack_ref: &ExternalPackRef,
1775        start: u64,
1776        end_inclusive: u64,
1777    ) -> Result<Option<Vec<u8>>, StoreError> {
1778        if pack_ref.len == 0 || start >= pack_ref.len || end_inclusive < start {
1779            return Ok(Some(Vec::new()));
1780        }
1781
1782        let actual_end = end_inclusive.min(pack_ref.len - 1);
1783        let read_len = actual_end.saturating_sub(start).saturating_add(1);
1784        let read_len = usize::try_from(read_len).map_err(|_| {
1785            StoreError::Other("external pack blob range is too large to read".to_string())
1786        })?;
1787        let path = self.external_pack_path(pack_ref).ok_or_else(|| {
1788            StoreError::Other(
1789                "external LMDB pack marker found but external blobs are disabled".into(),
1790            )
1791        })?;
1792        let mut file = File::open(path)?;
1793        file.seek(SeekFrom::Start(pack_ref.offset.saturating_add(start)))?;
1794        let mut data = vec![0; read_len];
1795        file.read_exact(&mut data)?;
1796        Ok(Some(data))
1797    }
1798
1799    fn external_blob_path_for_value(
1800        &self,
1801        hash: &Hash,
1802        value: &[u8],
1803    ) -> Result<Option<PathBuf>, StoreError> {
1804        if self.is_external_blob_ref(hash, value) {
1805            return Ok(self.external_blob_path(hash));
1806        }
1807        if Self::decode_external_pack_ref(value)?.is_some() {
1808            return Ok(None);
1809        }
1810        Ok(None)
1811    }
1812
1813    fn external_blob_path_in_txn(
1814        &self,
1815        txn: &heed::RoTxn,
1816        hash: &Hash,
1817    ) -> Result<Option<PathBuf>, StoreError> {
1818        let Some(value) = self
1819            .blobs
1820            .get(txn, hash)
1821            .map_err(|e| StoreError::Other(e.to_string()))?
1822        else {
1823            return Ok(None);
1824        };
1825        self.external_blob_path_for_value(hash, value)
1826    }
1827
1828    fn remove_external_blob_file(&self, path: Option<PathBuf>) {
1829        if let Some(path) = path {
1830            let _ = fs::remove_file(path);
1831        }
1832    }
1833}
1834
1835fn map_heed_error(error: HeedError) -> StoreError {
1836    match error {
1837        HeedError::Io(io_error) => StoreError::Io(io_error),
1838        other => StoreError::Other(other.to_string()),
1839    }
1840}
1841
1842fn is_map_full_store_error(err: &StoreError) -> bool {
1843    let message = err.to_string();
1844    message.contains("MDB_MAP_FULL") || message.contains("MapFull")
1845}
1846
1847fn env_bool(name: &str) -> Option<bool> {
1848    std::env::var(name).ok().and_then(|value| {
1849        let value = value.trim();
1850        if value == "1" || value.eq_ignore_ascii_case("true") || value.eq_ignore_ascii_case("yes") {
1851            Some(true)
1852        } else if value == "0"
1853            || value.eq_ignore_ascii_case("false")
1854            || value.eq_ignore_ascii_case("no")
1855        {
1856            Some(false)
1857        } else {
1858            None
1859        }
1860    })
1861}
1862
1863fn env_flags_from_env() -> EnvFlags {
1864    env_flags_from_bools(
1865        env_bool(LMDB_NO_READ_AHEAD_ENV).unwrap_or(false),
1866        env_bool(LMDB_NO_SYNC_ENV).unwrap_or(false),
1867        env_bool(LMDB_NO_META_SYNC_ENV).unwrap_or(false),
1868    )
1869}
1870
1871fn env_flags_from_bools(no_read_ahead: bool, no_sync: bool, no_meta_sync: bool) -> EnvFlags {
1872    let mut flags = EnvFlags::empty();
1873    if no_read_ahead {
1874        flags |= EnvFlags::NO_READ_AHEAD;
1875    }
1876    if no_sync {
1877        flags |= EnvFlags::NO_SYNC;
1878    }
1879    if no_meta_sync {
1880        flags |= EnvFlags::NO_META_SYNC;
1881    }
1882    flags
1883}
1884
1885fn unique_temp_path(path: &Path) -> PathBuf {
1886    let file_name = path
1887        .file_name()
1888        .and_then(|name| name.to_str())
1889        .unwrap_or("blob");
1890    let nanos = SystemTime::now()
1891        .duration_since(UNIX_EPOCH)
1892        .unwrap_or_default()
1893        .as_nanos();
1894    path.with_file_name(format!(".{file_name}.tmp.{}.{}", std::process::id(), nanos))
1895}
1896
1897fn unix_timestamp_now() -> u64 {
1898    SystemTime::now()
1899        .duration_since(UNIX_EPOCH)
1900        .unwrap_or_default()
1901        .as_secs()
1902}
1903
1904#[derive(Debug, Clone)]
1905pub struct LmdbStats {
1906    pub count: usize,
1907    pub total_bytes: u64,
1908    pub pinned_count: usize,
1909    pub pinned_bytes: u64,
1910}
1911
1912#[async_trait]
1913impl Store for LmdbBlobStore {
1914    async fn put(&self, hash: Hash, data: Vec<u8>) -> Result<bool, StoreError> {
1915        self.put_sync(hash, &data)
1916    }
1917
1918    async fn put_many(&self, items: Vec<(Hash, Vec<u8>)>) -> Result<usize, StoreError> {
1919        self.put_many_sync(&items)
1920    }
1921
1922    async fn get(&self, hash: &Hash) -> Result<Option<Vec<u8>>, StoreError> {
1923        self.get_sync(hash)
1924    }
1925
1926    async fn has(&self, hash: &Hash) -> Result<bool, StoreError> {
1927        self.exists(hash)
1928    }
1929
1930    async fn delete(&self, hash: &Hash) -> Result<bool, StoreError> {
1931        self.delete_sync(hash)
1932    }
1933
1934    fn set_max_bytes(&self, max: u64) {
1935        self.max_bytes.store(max, Ordering::Relaxed);
1936    }
1937
1938    fn max_bytes(&self) -> Option<u64> {
1939        let max = self.max_bytes.load(Ordering::Relaxed);
1940        if max > 0 {
1941            Some(max)
1942        } else {
1943            None
1944        }
1945    }
1946
1947    async fn stats(&self) -> StoreStats {
1948        match self.stats() {
1949            Ok(stats) => StoreStats {
1950                count: stats.count as u64,
1951                bytes: stats.total_bytes,
1952                pinned_count: stats.pinned_count as u64,
1953                pinned_bytes: stats.pinned_bytes,
1954            },
1955            Err(_) => StoreStats::default(),
1956        }
1957    }
1958
1959    async fn evict_if_needed(&self) -> Result<u64, StoreError> {
1960        let max = self.max_bytes.load(Ordering::Relaxed);
1961        if max == 0 {
1962            return Ok(0);
1963        }
1964
1965        let current = self.current_bytes.load(Ordering::Relaxed);
1966        if current <= max {
1967            return Ok(0);
1968        }
1969
1970        let target = max * 9 / 10;
1971        self.evict_to_target(current, target)
1972    }
1973
1974    async fn pin(&self, hash: &Hash) -> Result<(), StoreError> {
1975        self.pin_sync(hash)
1976    }
1977
1978    async fn unpin(&self, hash: &Hash) -> Result<(), StoreError> {
1979        self.unpin_sync(hash)
1980    }
1981
1982    fn pin_count(&self, hash: &Hash) -> u32 {
1983        let Ok(rtxn) = self.env.read_txn() else {
1984            return 0;
1985        };
1986        self.read_pin_count(&rtxn, hash).unwrap_or(0)
1987    }
1988}
1989
1990#[cfg(test)]
1991mod tests {
1992    use super::*;
1993    use hashtree_core::sha256;
1994    use heed::EnvOpenOptions;
1995    #[cfg(unix)]
1996    use std::path::{Path, PathBuf};
1997    #[cfg(unix)]
1998    use std::process::Command;
1999    use std::sync::{
2000        atomic::{AtomicBool, Ordering},
2001        Arc, Barrier,
2002    };
2003    use std::time::Duration;
2004    use tempfile::TempDir;
2005
2006    #[cfg(unix)]
2007    const STALE_READER_HELPER_ENV: &str = "HASHTREE_LMDB_STALE_READER_HELPER";
2008    #[cfg(unix)]
2009    const STALE_READER_HELPER_MODE_ENV: &str = "HASHTREE_LMDB_STALE_READER_HELPER_MODE";
2010    #[cfg(unix)]
2011    const STALE_READER_DB_PATH_ENV: &str = "HASHTREE_LMDB_STALE_READER_DB_PATH";
2012    #[cfg(unix)]
2013    const STALE_READER_MARKER_PATH_ENV: &str = "HASHTREE_LMDB_STALE_READER_MARKER_PATH";
2014    #[cfg(unix)]
2015    const TEST_MAX_READERS: u32 = 4;
2016
2017    fn count_files_under(path: &std::path::Path) -> std::io::Result<usize> {
2018        if !path.exists() {
2019            return Ok(0);
2020        }
2021        let mut count = 0usize;
2022        for entry in std::fs::read_dir(path)? {
2023            let entry = entry?;
2024            let metadata = entry.metadata()?;
2025            if metadata.is_dir() {
2026                count += count_files_under(&entry.path())?;
2027            } else if metadata.is_file() {
2028                count += 1;
2029            }
2030        }
2031        Ok(count)
2032    }
2033
2034    fn persisted_next_order(store: &LmdbBlobStore) -> Result<u64, StoreError> {
2035        let rtxn = store.env.read_txn().map_err(map_heed_error)?;
2036        LmdbBlobStore::read_next_order(store.stats, &rtxn)?
2037            .ok_or_else(|| StoreError::Other("missing persisted next_order".to_string()))
2038    }
2039
2040    fn metadata_order(store: &LmdbBlobStore, hash: &Hash) -> Result<u64, StoreError> {
2041        let rtxn = store.env.read_txn().map_err(map_heed_error)?;
2042        let meta = store
2043            .metadata
2044            .get(&rtxn, hash)
2045            .map_err(map_heed_error)?
2046            .ok_or_else(|| StoreError::Other("missing blob metadata".to_string()))?;
2047        LmdbBlobStore::decode_blob_meta(meta).map(|meta| meta.order)
2048    }
2049
2050    fn delete_persisted_next_order(store: &LmdbBlobStore) -> Result<(), StoreError> {
2051        let mut wtxn = store.env.write_txn().map_err(map_heed_error)?;
2052        store
2053            .stats
2054            .delete(&mut wtxn, NEXT_ORDER_KEY)
2055            .map_err(map_heed_error)?;
2056        wtxn.commit().map_err(map_heed_error)
2057    }
2058
2059    #[cfg(unix)]
2060    fn run_helper(mode: &str, path: &Path, marker: &Path) {
2061        let output = Command::new(std::env::current_exe().expect("test binary path"))
2062            .arg("--ignored")
2063            .arg("--exact")
2064            .arg("tests::lmdb_stale_reader_helper")
2065            .env(STALE_READER_HELPER_ENV, "1")
2066            .env(STALE_READER_HELPER_MODE_ENV, mode)
2067            .env(STALE_READER_DB_PATH_ENV, path)
2068            .env(STALE_READER_MARKER_PATH_ENV, marker)
2069            .env("RUST_TEST_THREADS", "1")
2070            .output()
2071            .expect("spawn stale-reader helper");
2072
2073        assert!(
2074            output.status.success(),
2075            "stale-reader helper failed: stdout={} stderr={}",
2076            String::from_utf8_lossy(&output.stdout),
2077            String::from_utf8_lossy(&output.stderr)
2078        );
2079        assert!(
2080            marker.exists(),
2081            "stale-reader helper did not run: stdout={} stderr={}",
2082            String::from_utf8_lossy(&output.stdout),
2083            String::from_utf8_lossy(&output.stderr)
2084        );
2085    }
2086
2087    #[test]
2088    fn env_flags_from_bools_enables_bulk_ingest_flags() {
2089        let flags = env_flags_from_bools(true, true, true);
2090
2091        assert!(flags.contains(EnvFlags::NO_READ_AHEAD));
2092        assert!(flags.contains(EnvFlags::NO_SYNC));
2093        assert!(flags.contains(EnvFlags::NO_META_SYNC));
2094    }
2095
2096    #[tokio::test]
2097    async fn test_put_get() -> Result<(), StoreError> {
2098        let temp = TempDir::new().unwrap();
2099        let store = LmdbBlobStore::new(temp.path().join("blobs"))?;
2100
2101        let data = b"hello lmdb";
2102        let hash = sha256(data);
2103        store.put(hash, data.to_vec()).await?;
2104
2105        assert!(store.has(&hash).await?);
2106        assert_eq!(store.get(&hash).await?, Some(data.to_vec()));
2107
2108        Ok(())
2109    }
2110
2111    #[test]
2112    fn external_blob_spill_keeps_large_values_out_of_lmdb() -> Result<(), StoreError> {
2113        let temp = TempDir::new().unwrap();
2114        let external_dir = temp.path().join("external");
2115        let store = LmdbBlobStore::with_map_size_and_settings(
2116            temp.path().join("blobs"),
2117            16 * 1024 * 1024,
2118            EnvFlags::empty(),
2119            |_| {
2120                Some(ExternalBlobConfig {
2121                    base_path: external_dir.clone(),
2122                    min_bytes: 8,
2123                    sync: false,
2124                    pack_target_bytes: None,
2125                })
2126            },
2127        )?;
2128
2129        let small = b"tiny";
2130        let small_hash = sha256(small);
2131        assert!(store.put_sync(small_hash, small)?);
2132
2133        let large = b"large external blob payload";
2134        let large_hash = sha256(large);
2135        assert!(store.put_sync(large_hash, large)?);
2136
2137        assert_eq!(store.get_sync(&small_hash)?, Some(small.to_vec()));
2138        assert_eq!(store.get_sync(&large_hash)?, Some(large.to_vec()));
2139        assert_eq!(
2140            store.get_range_sync(&large_hash, 6, 13)?,
2141            Some(b"external".to_vec())
2142        );
2143
2144        let rtxn = store.env.read_txn().map_err(map_heed_error)?;
2145        let inline_value = store
2146            .blobs
2147            .get(&rtxn, &small_hash)
2148            .map_err(map_heed_error)?
2149            .expect("small inline value");
2150        assert_eq!(inline_value, small);
2151        let external_value = store
2152            .blobs
2153            .get(&rtxn, &large_hash)
2154            .map_err(map_heed_error)?
2155            .expect("large marker value");
2156        assert!(store.is_external_blob_ref(&large_hash, external_value));
2157        drop(rtxn);
2158
2159        let external_path = store.external_blob_path(&large_hash).unwrap();
2160        assert_eq!(std::fs::read(&external_path)?, large);
2161        let stats = store.stats()?;
2162        assert_eq!(stats.count, 2);
2163        assert_eq!(stats.total_bytes, (small.len() + large.len()) as u64);
2164
2165        assert!(store.delete_sync(&large_hash)?);
2166        assert!(!external_path.exists());
2167        Ok(())
2168    }
2169
2170    #[test]
2171    fn external_blob_pack_batches_large_values() -> Result<(), StoreError> {
2172        let temp = TempDir::new().unwrap();
2173        let external_dir = temp.path().join("external");
2174        let store = LmdbBlobStore::with_map_size_and_settings(
2175            temp.path().join("blobs"),
2176            16 * 1024 * 1024,
2177            EnvFlags::empty(),
2178            |_| {
2179                Some(ExternalBlobConfig {
2180                    base_path: external_dir.clone(),
2181                    min_bytes: 8,
2182                    sync: true,
2183                    pack_target_bytes: Some(1024 * 1024),
2184                })
2185            },
2186        )?;
2187
2188        let first = b"first packed blob payload".to_vec();
2189        let second = b"second packed blob payload".to_vec();
2190        let tiny = b"tiny".to_vec();
2191        let first_hash = sha256(&first);
2192        let second_hash = sha256(&second);
2193        let tiny_hash = sha256(&tiny);
2194        let items = vec![
2195            (first_hash, first.clone()),
2196            (tiny_hash, tiny.clone()),
2197            (second_hash, second.clone()),
2198        ];
2199
2200        assert_eq!(store.put_many_sync(&items)?, 3);
2201        assert_eq!(store.get_sync(&first_hash)?, Some(first.clone()));
2202        assert_eq!(store.get_sync(&second_hash)?, Some(second.clone()));
2203        assert_eq!(store.get_sync(&tiny_hash)?, Some(tiny.clone()));
2204        assert_eq!(
2205            store.get_range_sync(&second_hash, 7, 12)?,
2206            Some(b"packed".to_vec())
2207        );
2208
2209        let rtxn = store.env.read_txn().map_err(map_heed_error)?;
2210        let first_value = store
2211            .blobs
2212            .get(&rtxn, &first_hash)
2213            .map_err(map_heed_error)?
2214            .expect("first pack marker");
2215        let second_value = store
2216            .blobs
2217            .get(&rtxn, &second_hash)
2218            .map_err(map_heed_error)?
2219            .expect("second pack marker");
2220        let tiny_value = store
2221            .blobs
2222            .get(&rtxn, &tiny_hash)
2223            .map_err(map_heed_error)?
2224            .expect("tiny inline value");
2225        let first_pack =
2226            LmdbBlobStore::decode_external_pack_ref(first_value)?.expect("first external pack ref");
2227        let second_pack = LmdbBlobStore::decode_external_pack_ref(second_value)?
2228            .expect("second external pack ref");
2229        assert_eq!(first_pack.name, second_pack.name);
2230        assert_eq!(tiny_value, tiny.as_slice());
2231        drop(rtxn);
2232
2233        let pack_path = store.external_pack_path(&first_pack).unwrap();
2234        assert!(pack_path.exists());
2235        assert_eq!(
2236            std::fs::read(&pack_path)?,
2237            [first.as_slice(), second.as_slice()].concat()
2238        );
2239        let pack_count_after_first_write = count_files_under(&external_dir.join("packs"))?;
2240
2241        assert_eq!(store.put_many_sync(&items)?, 0);
2242        assert_eq!(
2243            count_files_under(&external_dir.join("packs"))?,
2244            pack_count_after_first_write,
2245            "rewriting an already-present batch must not create orphan external pack files"
2246        );
2247
2248        assert!(store.delete_sync(&first_hash)?);
2249        assert!(pack_path.exists());
2250        assert_eq!(store.get_sync(&second_hash)?, Some(second));
2251        Ok(())
2252    }
2253
2254    #[tokio::test]
2255    async fn test_delete() -> Result<(), StoreError> {
2256        let temp = TempDir::new().unwrap();
2257        let store = LmdbBlobStore::new(temp.path().join("blobs"))?;
2258
2259        let data = b"delete me";
2260        let hash = sha256(data);
2261        store.put(hash, data.to_vec()).await?;
2262        assert!(store.has(&hash).await?);
2263
2264        assert!(store.delete(&hash).await?);
2265        assert!(!store.has(&hash).await?);
2266        assert!(!store.delete(&hash).await?);
2267
2268        Ok(())
2269    }
2270
2271    #[tokio::test]
2272    async fn test_list() -> Result<(), StoreError> {
2273        let temp = TempDir::new().unwrap();
2274        let store = LmdbBlobStore::new(temp.path().join("blobs"))?;
2275
2276        let d1 = b"one";
2277        let d2 = b"two";
2278        let d3 = b"three";
2279        let h1 = sha256(d1);
2280        let h2 = sha256(d2);
2281        let h3 = sha256(d3);
2282
2283        store.put(h1, d1.to_vec()).await?;
2284        store.put(h2, d2.to_vec()).await?;
2285        store.put(h3, d3.to_vec()).await?;
2286
2287        let hashes = store.list()?;
2288        assert_eq!(hashes.len(), 3);
2289        assert!(hashes.contains(&h1));
2290        assert!(hashes.contains(&h2));
2291        assert!(hashes.contains(&h3));
2292
2293        Ok(())
2294    }
2295
2296    #[test]
2297    fn existing_hashes_in_sorted_candidates_marks_present_hashes() -> Result<(), StoreError> {
2298        let temp = TempDir::new().unwrap();
2299        let store = LmdbBlobStore::new(temp.path().join("blobs"))?;
2300
2301        let h1 = sha256(b"range present one");
2302        let h2 = sha256(b"range present two");
2303        let missing = sha256(b"range missing");
2304        store.put_sync(h1, b"range present one")?;
2305        store.put_sync(h2, b"range present two")?;
2306
2307        let mut candidates = vec![missing, h2, h1, h1];
2308        candidates.sort_unstable();
2309        let existing = store.existing_hashes_in_sorted_candidates(&candidates)?;
2310
2311        assert_eq!(existing.len(), candidates.len());
2312        for (hash, exists) in candidates.iter().zip(existing) {
2313            assert_eq!(exists, *hash == h1 || *hash == h2);
2314        }
2315
2316        Ok(())
2317    }
2318
2319    #[tokio::test]
2320    async fn test_blob_last_accessed_persists_and_updates_eviction_order() -> Result<(), StoreError>
2321    {
2322        let temp = TempDir::new().unwrap();
2323        let store = LmdbBlobStore::new(temp.path().join("blobs"))?;
2324
2325        let first = sha256(b"first");
2326        let second = sha256(b"second");
2327        store.put(first, b"first".to_vec()).await?;
2328        store.put(second, b"second".to_vec()).await?;
2329
2330        assert!(store.last_accessed_at_sync(&first)?.unwrap_or(0) > 0);
2331        let access_time = unix_timestamp_now().saturating_add(1000);
2332        store.touch_accessed_sync(&first, access_time)?;
2333
2334        assert_eq!(store.last_accessed_at_sync(&first)?, Some(access_time));
2335        assert!(store.delete_sync(&first)?);
2336        assert!(store.exists(&second)?);
2337
2338        Ok(())
2339    }
2340
2341    #[test]
2342    fn duplicate_put_is_noop_and_preserves_blob_last_accessed() -> Result<(), StoreError> {
2343        let temp = TempDir::new().unwrap();
2344        let store = LmdbBlobStore::new(temp.path().join("blobs"))?;
2345
2346        let data = b"already here";
2347        let hash = sha256(data);
2348        assert!(store.put_sync(hash, data)?);
2349        let accessed = store.last_accessed_at_sync(&hash)?;
2350        let next_order = persisted_next_order(&store)?;
2351
2352        assert!(!store.put_sync(hash, data)?);
2353        assert_eq!(store.last_accessed_at_sync(&hash)?, accessed);
2354        assert_eq!(persisted_next_order(&store)?, next_order);
2355
2356        let stats = store.stats()?;
2357        assert_eq!(stats.count, 1);
2358        assert_eq!(stats.total_bytes, data.len() as u64);
2359
2360        Ok(())
2361    }
2362
2363    #[test]
2364    fn next_order_persists_across_reopen_and_counts_only_new_blobs() -> Result<(), StoreError> {
2365        let temp = TempDir::new().unwrap();
2366        let path = temp.path().join("blobs");
2367        let first = sha256(b"order first");
2368        let second = sha256(b"order second");
2369        let third = sha256(b"order third");
2370
2371        {
2372            let store = LmdbBlobStore::new(&path)?;
2373            assert!(store.put_sync(first, b"order first")?);
2374            assert!(store.put_sync(second, b"order second")?);
2375            assert_eq!(metadata_order(&store, &first)?, 0);
2376            assert_eq!(metadata_order(&store, &second)?, 1);
2377            assert_eq!(persisted_next_order(&store)?, 2);
2378
2379            assert!(!store.put_sync(first, b"order first")?);
2380            assert_eq!(
2381                store.put_many_sync(&[(second, b"order second".to_vec())])?,
2382                0
2383            );
2384            assert_eq!(persisted_next_order(&store)?, 2);
2385        }
2386
2387        let reopened = LmdbBlobStore::new(&path)?;
2388        assert_eq!(persisted_next_order(&reopened)?, 2);
2389        assert!(reopened.put_sync(third, b"order third")?);
2390        assert_eq!(metadata_order(&reopened, &third)?, 2);
2391        assert_eq!(persisted_next_order(&reopened)?, 3);
2392
2393        Ok(())
2394    }
2395
2396    #[test]
2397    fn legacy_store_without_next_order_seeds_high_counter_without_tail_scan(
2398    ) -> Result<(), StoreError> {
2399        let temp = TempDir::new().unwrap();
2400        let path = temp.path().join("blobs");
2401        let existing = sha256(b"legacy existing");
2402        let new = sha256(b"legacy new");
2403
2404        {
2405            let store = LmdbBlobStore::new(&path)?;
2406            assert!(store.put_sync(existing, b"legacy existing")?);
2407            delete_persisted_next_order(&store)?;
2408        }
2409
2410        let reopened = LmdbBlobStore::new(&path)?;
2411        let migrated_next_order = persisted_next_order(&reopened)?;
2412        assert!(migrated_next_order >= NEXT_ORDER_MIGRATION_FLOOR);
2413        assert!(reopened.put_sync(new, b"legacy new")?);
2414        assert_eq!(metadata_order(&reopened, &new)?, migrated_next_order);
2415        assert_eq!(
2416            persisted_next_order(&reopened)?,
2417            migrated_next_order.saturating_add(1)
2418        );
2419
2420        Ok(())
2421    }
2422
2423    #[test]
2424    fn put_many_report_counts_only_new_hashes_and_bytes() -> Result<(), StoreError> {
2425        let temp = TempDir::new().unwrap();
2426        let store = LmdbBlobStore::new(temp.path().join("blobs"))?;
2427
2428        let existing = b"existing";
2429        let existing_hash = sha256(existing);
2430        assert!(store.put_sync(existing_hash, existing)?);
2431
2432        let new_one = b"new one";
2433        let new_two = b"new two";
2434        let new_one_hash = sha256(new_one);
2435        let new_two_hash = sha256(new_two);
2436        let items = vec![
2437            (existing_hash, existing.to_vec()),
2438            (new_one_hash, new_one.to_vec()),
2439            (new_one_hash, new_one.to_vec()),
2440            (new_two_hash, new_two.to_vec()),
2441        ];
2442
2443        let report = store.put_many_report_sync(&items)?;
2444
2445        assert_eq!(report.total, 4);
2446        assert_eq!(report.inserted, 2);
2447        assert_eq!(
2448            report.inserted_bytes,
2449            (new_one.len() + new_two.len()) as u64
2450        );
2451        assert_eq!(report.inserted_hashes, vec![new_one_hash, new_two_hash]);
2452        assert_eq!(store.put_many_sync(&items)?, 0);
2453        assert_eq!(
2454            store.stats()?.total_bytes,
2455            (existing.len() + new_one.len() + new_two.len()) as u64
2456        );
2457
2458        Ok(())
2459    }
2460
2461    #[test]
2462    fn duplicate_heavy_batch_does_not_evict_by_candidate_bytes() -> Result<(), StoreError> {
2463        let temp = TempDir::new().unwrap();
2464        let store = LmdbBlobStore::new(temp.path().join("blobs"))?;
2465        store.set_max_bytes(35);
2466
2467        let first = [1u8; 10];
2468        let second = [2u8; 10];
2469        let third = [3u8; 10];
2470        let new = [4u8; 5];
2471        let first_hash = sha256(&first);
2472        let second_hash = sha256(&second);
2473        let third_hash = sha256(&third);
2474        let new_hash = sha256(&new);
2475        assert!(store.put_sync(first_hash, &first)?);
2476        assert!(store.put_sync(second_hash, &second)?);
2477        assert!(store.put_sync(third_hash, &third)?);
2478        assert_eq!(store.stats()?.total_bytes, 30);
2479
2480        let report = store.put_many_report_sync(&[
2481            (first_hash, first.to_vec()),
2482            (second_hash, second.to_vec()),
2483            (new_hash, new.to_vec()),
2484        ])?;
2485
2486        assert_eq!(report.inserted, 1);
2487        assert_eq!(report.inserted_bytes, 5);
2488        assert_eq!(store.stats()?.total_bytes, 35);
2489        assert!(store.exists(&first_hash)?);
2490        assert!(store.exists(&second_hash)?);
2491        assert!(store.exists(&third_hash)?);
2492        assert!(store.exists(&new_hash)?);
2493
2494        Ok(())
2495    }
2496
2497    #[test]
2498    fn test_decodes_legacy_blob_metadata_without_access_time() -> Result<(), StoreError> {
2499        let mut encoded = [0u8; LEGACY_BLOB_META_BYTES];
2500        encoded[..8].copy_from_slice(&7u64.to_be_bytes());
2501        encoded[8..].copy_from_slice(&42u64.to_be_bytes());
2502
2503        let meta = LmdbBlobStore::decode_blob_meta(&encoded)?;
2504
2505        assert_eq!(meta.order, 7);
2506        assert_eq!(meta.size, 42);
2507        assert_eq!(meta.last_accessed_at, 0);
2508        Ok(())
2509    }
2510
2511    #[tokio::test]
2512    async fn test_stats() -> Result<(), StoreError> {
2513        let temp = TempDir::new().unwrap();
2514        let store = LmdbBlobStore::new(temp.path().join("blobs"))?;
2515
2516        let d1 = b"hello";
2517        let d2 = b"world";
2518        store.put(sha256(d1), d1.to_vec()).await?;
2519        store.put(sha256(d2), d2.to_vec()).await?;
2520
2521        let stats = store.stats()?;
2522        assert_eq!(stats.count, 2);
2523        assert_eq!(stats.total_bytes, 10);
2524
2525        Ok(())
2526    }
2527
2528    #[tokio::test]
2529    async fn test_stats_persist_across_reopen_and_mutations() -> Result<(), StoreError> {
2530        let temp = TempDir::new().unwrap();
2531        let path = temp.path().join("blobs");
2532        let h1 = sha256(b"hello");
2533        let h2 = sha256(b"world!");
2534        let h3 = sha256(b"prepin");
2535
2536        {
2537            let store = LmdbBlobStore::new(&path)?;
2538            store.put(h1, b"hello".to_vec()).await?;
2539            store.put(h2, b"world!".to_vec()).await?;
2540            store.pin(&h1).await?;
2541            store.pin(&h3).await?;
2542            store.put(h3, b"prepin".to_vec()).await?;
2543
2544            let stats = store.stats()?;
2545            assert_eq!(stats.count, 3);
2546            assert_eq!(stats.total_bytes, 17);
2547            assert_eq!(stats.pinned_count, 2);
2548            assert_eq!(stats.pinned_bytes, 11);
2549        }
2550
2551        {
2552            let reopened = LmdbBlobStore::new(&path)?;
2553            let stats = reopened.stats()?;
2554            assert_eq!(stats.count, 3);
2555            assert_eq!(stats.total_bytes, 17);
2556            assert_eq!(stats.pinned_count, 2);
2557            assert_eq!(stats.pinned_bytes, 11);
2558
2559            assert!(reopened.delete(&h1).await?);
2560            let stats = reopened.stats()?;
2561            assert_eq!(stats.count, 2);
2562            assert_eq!(stats.total_bytes, 12);
2563            assert_eq!(stats.pinned_count, 1);
2564            assert_eq!(stats.pinned_bytes, 6);
2565        }
2566
2567        let reopened = LmdbBlobStore::new(&path)?;
2568        let stats = reopened.stats()?;
2569        assert_eq!(stats.count, 2);
2570        assert_eq!(stats.total_bytes, 12);
2571        assert_eq!(stats.pinned_count, 1);
2572        assert_eq!(stats.pinned_bytes, 6);
2573
2574        Ok(())
2575    }
2576
2577    #[tokio::test]
2578    async fn test_deduplication() -> Result<(), StoreError> {
2579        let temp = TempDir::new().unwrap();
2580        let store = LmdbBlobStore::new(temp.path().join("blobs"))?;
2581
2582        let data = b"same";
2583        let hash = sha256(data);
2584        assert!(store.put(hash, data.to_vec()).await?); // Returns true (newly stored)
2585        assert!(!store.put(hash, data.to_vec()).await?); // Returns false (already existed)
2586
2587        assert_eq!(store.list()?.len(), 1);
2588
2589        Ok(())
2590    }
2591
2592    #[tokio::test]
2593    async fn test_max_bytes() -> Result<(), StoreError> {
2594        let temp = TempDir::new().unwrap();
2595        let store = LmdbBlobStore::new(temp.path().join("blobs"))?;
2596
2597        assert!(store.max_bytes().is_none());
2598
2599        store.set_max_bytes(1000);
2600        assert_eq!(store.max_bytes(), Some(1000));
2601
2602        store.set_max_bytes(0);
2603        assert!(store.max_bytes().is_none());
2604
2605        Ok(())
2606    }
2607
2608    #[test]
2609    fn test_with_max_bytes_expands_lmdb_map_size() -> Result<(), StoreError> {
2610        let temp = TempDir::new().unwrap();
2611        let requested = (DEFAULT_MAP_SIZE as u64) + 64 * 1024 * 1024;
2612        let store = LmdbBlobStore::with_max_bytes(temp.path().join("blobs"), requested)?;
2613
2614        assert!(
2615            store.map_size_bytes() as u64 >= requested,
2616            "expected LMDB map to grow to at least {requested} bytes, got {}",
2617            store.map_size_bytes()
2618        );
2619        assert_eq!(store.max_bytes(), Some(requested));
2620
2621        Ok(())
2622    }
2623
2624    #[tokio::test]
2625    async fn test_eviction_over_limit() -> Result<(), StoreError> {
2626        let temp = TempDir::new().unwrap();
2627        let store = LmdbBlobStore::new(temp.path().join("blobs"))?;
2628        store.set_max_bytes(25);
2629
2630        let h1 = sha256(b"aaaaaaaaaa");
2631        let h2 = sha256(b"bbbbbbbbbb");
2632        let h3 = sha256(b"cccccccccc");
2633
2634        store.put(h1, b"aaaaaaaaaa".to_vec()).await?;
2635        store.put(h2, b"bbbbbbbbbb".to_vec()).await?;
2636        store.put(h3, b"cccccccccc".to_vec()).await?;
2637
2638        let freed = store.evict_if_needed().await?;
2639        assert_eq!(
2640            freed, 0,
2641            "write path should have already evicted stale blobs"
2642        );
2643
2644        assert!(
2645            !store.has(&h1).await?,
2646            "oldest blob should be evicted before the third write"
2647        );
2648        assert!(store.has(&h2).await?);
2649        assert!(store.has(&h3).await?);
2650
2651        let stats = store.stats()?;
2652        assert!(
2653            stats.total_bytes <= 22,
2654            "store should be reduced to 90% target"
2655        );
2656
2657        Ok(())
2658    }
2659
2660    #[tokio::test]
2661    async fn test_eviction_respects_pins() -> Result<(), StoreError> {
2662        let temp = TempDir::new().unwrap();
2663        let store = LmdbBlobStore::new(temp.path().join("blobs"))?;
2664        store.set_max_bytes(25);
2665
2666        let h1 = sha256(b"aaaaaaaaaa");
2667        let h2 = sha256(b"bbbbbbbbbb");
2668        let h3 = sha256(b"cccccccccc");
2669
2670        store.put(h1, b"aaaaaaaaaa".to_vec()).await?;
2671        store.put(h2, b"bbbbbbbbbb".to_vec()).await?;
2672        store.pin(&h1).await?;
2673        store.put(h3, b"cccccccccc".to_vec()).await?;
2674
2675        let freed = store.evict_if_needed().await?;
2676        assert_eq!(
2677            freed, 0,
2678            "write path should have already evicted stale blobs"
2679        );
2680
2681        assert!(store.has(&h1).await?, "pinned blob must not be evicted");
2682        assert!(
2683            !store.has(&h2).await?,
2684            "oldest unpinned blob should be evicted before the third write"
2685        );
2686        assert!(store.has(&h3).await?);
2687
2688        Ok(())
2689    }
2690
2691    #[tokio::test]
2692    async fn test_reopen_with_existing_eviction_order() -> Result<(), StoreError> {
2693        let temp = TempDir::new().unwrap();
2694        let path = temp.path().join("blobs");
2695
2696        {
2697            let store = LmdbBlobStore::new(&path)?;
2698            let h1 = sha256(b"aaaaaaaaaa");
2699            let h2 = sha256(b"bbbbbbbbbb");
2700            store.put(h1, b"aaaaaaaaaa".to_vec()).await?;
2701            store.put(h2, b"bbbbbbbbbb".to_vec()).await?;
2702        }
2703
2704        let reopened = LmdbBlobStore::new(&path)?;
2705        let h3 = sha256(b"cccccccccc");
2706        assert!(reopened.put(h3, b"cccccccccc".to_vec()).await?);
2707        assert!(reopened.has(&h3).await?);
2708
2709        Ok(())
2710    }
2711
2712    #[tokio::test]
2713    async fn test_reopen_with_lower_max_bytes_evicts_existing_blobs() -> Result<(), StoreError> {
2714        let temp = TempDir::new().unwrap();
2715        let path = temp.path().join("blobs");
2716
2717        {
2718            let store = LmdbBlobStore::new(&path)?;
2719            let h1 = sha256(b"aaaaaaaaaa");
2720            let h2 = sha256(b"bbbbbbbbbb");
2721            let h3 = sha256(b"cccccccccc");
2722            store.put(h1, b"aaaaaaaaaa".to_vec()).await?;
2723            store.put(h2, b"bbbbbbbbbb".to_vec()).await?;
2724            store.put(h3, b"cccccccccc".to_vec()).await?;
2725        }
2726
2727        let reopened = LmdbBlobStore::with_max_bytes(&path, 25)?;
2728        let h1 = sha256(b"aaaaaaaaaa");
2729        let h2 = sha256(b"bbbbbbbbbb");
2730        let h3 = sha256(b"cccccccccc");
2731
2732        assert!(
2733            !reopened.has(&h1).await?,
2734            "oldest blob should be evicted when reopening over the new cap"
2735        );
2736        assert!(reopened.has(&h2).await?);
2737        assert!(reopened.has(&h3).await?);
2738
2739        let stats = reopened.stats()?;
2740        assert!(
2741            stats.total_bytes <= 22,
2742            "reopened store should be reduced to the 90% target"
2743        );
2744
2745        Ok(())
2746    }
2747
2748    #[test]
2749    fn test_supports_many_concurrent_readers() -> Result<(), Box<dyn std::error::Error>> {
2750        const READER_THREADS: usize = 160;
2751
2752        let temp = TempDir::new()?;
2753        let store = Arc::new(LmdbBlobStore::new(temp.path().join("blobs"))?);
2754        let hash = sha256(b"many readers");
2755        store.put_sync(hash, b"many readers")?;
2756
2757        let start = Arc::new(Barrier::new(READER_THREADS + 1));
2758        let release = Arc::new(AtomicBool::new(false));
2759        let mut handles = Vec::with_capacity(READER_THREADS);
2760
2761        for _ in 0..READER_THREADS {
2762            let env = store.env.clone();
2763            let start = Arc::clone(&start);
2764            let release = Arc::clone(&release);
2765            handles.push(std::thread::spawn(move || -> Result<(), String> {
2766                start.wait();
2767                let _rtxn = env.read_txn().map_err(|err| err.to_string())?;
2768                while !release.load(Ordering::Relaxed) {
2769                    std::thread::sleep(Duration::from_millis(1));
2770                }
2771                Ok(())
2772            }));
2773        }
2774
2775        start.wait();
2776        std::thread::sleep(Duration::from_millis(50));
2777        release.store(true, Ordering::Relaxed);
2778
2779        let results: Vec<Result<(), String>> = handles
2780            .into_iter()
2781            .map(|handle| handle.join().expect("reader thread panicked"))
2782            .collect();
2783
2784        let failures: Vec<String> = results.into_iter().filter_map(Result::err).collect();
2785        assert!(
2786            failures.is_empty(),
2787            "concurrent reader failures: {}",
2788            failures.join(" | ")
2789        );
2790        assert!(store.exists(&hash)?);
2791
2792        Ok(())
2793    }
2794
2795    #[cfg(unix)]
2796    #[test]
2797    fn test_reclaims_stale_reader_slots() -> Result<(), Box<dyn std::error::Error>> {
2798        let temp = TempDir::new()?;
2799        let path = temp.path().join("blobs");
2800        let data = b"hello stale readers";
2801        let hash = sha256(data);
2802
2803        run_helper("setup", &path, &temp.path().join("setup.marker"));
2804
2805        for index in 0..TEST_MAX_READERS {
2806            let marker = temp.path().join(format!("helper-{index}.marker"));
2807            run_helper("stale", &path, &marker);
2808        }
2809
2810        let store = LmdbBlobStore::with_map_size(&path, 1024 * 1024)?;
2811        assert!(store.exists(&hash)?);
2812
2813        Ok(())
2814    }
2815
2816    #[cfg(unix)]
2817    #[test]
2818    fn test_reopens_existing_env_with_larger_map_size() -> Result<(), Box<dyn std::error::Error>> {
2819        let temp = TempDir::new()?;
2820        let path = temp.path().join("blobs");
2821        run_helper("small-map", &path, &temp.path().join("small-map.marker"));
2822
2823        let reopened = LmdbBlobStore::with_map_size(&path, 8 * 1024 * 1024)?;
2824        assert!(reopened.map_size_bytes() >= 8 * 1024 * 1024);
2825
2826        Ok(())
2827    }
2828
2829    #[cfg(unix)]
2830    #[test]
2831    fn test_reopens_existing_env_with_smaller_requested_map_size(
2832    ) -> Result<(), Box<dyn std::error::Error>> {
2833        let temp = TempDir::new()?;
2834        let path = temp.path().join("blobs");
2835        let hash = sha256(b"large existing blob");
2836        run_helper("large-map", &path, &temp.path().join("large-map.marker"));
2837
2838        let existing_size = std::fs::metadata(path.join("data.mdb"))?.len();
2839        assert!(
2840            existing_size > 1024 * 1024,
2841            "test setup should create an environment larger than the reopen request"
2842        );
2843
2844        let reopened = LmdbBlobStore::with_map_size(&path, 1024 * 1024)?;
2845        assert!(
2846            reopened.map_size_bytes() as u64 >= existing_size,
2847            "expected map size to cover existing data.mdb size {existing_size}, got {}",
2848            reopened.map_size_bytes()
2849        );
2850        assert!(reopened.exists(&hash)?);
2851
2852        Ok(())
2853    }
2854
2855    #[cfg(unix)]
2856    #[test]
2857    #[ignore = "used as a subprocess helper by test_reclaims_stale_reader_slots"]
2858    fn lmdb_stale_reader_helper() {
2859        let Some(db_path) = std::env::var_os(STALE_READER_DB_PATH_ENV) else {
2860            return;
2861        };
2862        let marker_path =
2863            PathBuf::from(std::env::var_os(STALE_READER_MARKER_PATH_ENV).expect("marker path"));
2864        std::fs::write(&marker_path, b"started").expect("write helper marker");
2865
2866        let _env_flag = std::env::var_os(STALE_READER_HELPER_ENV).expect("helper mode enabled");
2867        let mode = std::env::var(STALE_READER_HELPER_MODE_ENV).expect("helper mode");
2868        let db_path = PathBuf::from(db_path);
2869        std::fs::create_dir_all(&db_path).expect("create helper db dir");
2870        let helper_map_size = if mode == "large-map" {
2871            8 * 1024 * 1024
2872        } else {
2873            1024 * 1024
2874        };
2875        let env = unsafe {
2876            EnvOpenOptions::new()
2877                .map_size(helper_map_size)
2878                .max_dbs(DATABASE_COUNT)
2879                .max_readers(TEST_MAX_READERS)
2880                .open(&db_path)
2881                .expect("open lmdb env")
2882        };
2883        match mode.as_str() {
2884            "setup" => {
2885                let mut wtxn = env.write_txn().expect("open write txn");
2886                let blobs: Database<Bytes, Bytes> = env
2887                    .create_database(&mut wtxn, Some("blobs"))
2888                    .expect("create blobs database");
2889                let data = b"hello stale readers";
2890                let hash = sha256(data);
2891                blobs.put(&mut wtxn, &hash, data).expect("seed blob");
2892                wtxn.commit().expect("commit setup txn");
2893                std::process::exit(0);
2894            }
2895            "stale" => {
2896                let _rtxn = env.read_txn().expect("open read txn");
2897                std::process::exit(0);
2898            }
2899            "small-map" => {
2900                let mut wtxn = env.write_txn().expect("open write txn");
2901                let _blobs: Database<Bytes, Bytes> = env
2902                    .create_database(&mut wtxn, Some("blobs"))
2903                    .expect("create blobs database");
2904                let _metadata: Database<Bytes, Bytes> = env
2905                    .create_database(&mut wtxn, Some("metadata"))
2906                    .expect("create metadata database");
2907                let _eviction_order: Database<Bytes, Unit> = env
2908                    .create_database(&mut wtxn, Some("eviction_order"))
2909                    .expect("create eviction_order database");
2910                let _pins: Database<Bytes, Bytes> = env
2911                    .create_database(&mut wtxn, Some("pins"))
2912                    .expect("create pins database");
2913                wtxn.commit().expect("commit small-map setup txn");
2914                std::process::exit(0);
2915            }
2916            "large-map" => {
2917                let mut wtxn = env.write_txn().expect("open write txn");
2918                let blobs: Database<Bytes, Bytes> = env
2919                    .create_database(&mut wtxn, Some("blobs"))
2920                    .expect("create blobs database");
2921                let _metadata: Database<Bytes, Bytes> = env
2922                    .create_database(&mut wtxn, Some("metadata"))
2923                    .expect("create metadata database");
2924                let _eviction_order: Database<Bytes, Unit> = env
2925                    .create_database(&mut wtxn, Some("eviction_order"))
2926                    .expect("create eviction_order database");
2927                let _pins: Database<Bytes, Bytes> = env
2928                    .create_database(&mut wtxn, Some("pins"))
2929                    .expect("create pins database");
2930                let data = vec![42u8; 3 * 1024 * 1024];
2931                let hash = sha256(b"large existing blob");
2932                blobs.put(&mut wtxn, &hash, &data).expect("seed blob");
2933                wtxn.commit().expect("commit large-map setup txn");
2934                std::process::exit(0);
2935            }
2936            other => panic!("unknown helper mode: {other}"),
2937        }
2938    }
2939}