Skip to main content

hashtree_lmdb/
lib.rs

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