1mod 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
32pub 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#[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 pub fn from_env(env_path: &Path) -> Option<Self> {
100 ExternalBlobConfig::from_env(env_path).map(Into::into)
101 }
102
103 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
164pub struct LmdbBlobStore {
166 env: ManagedEnv,
167 blobs: Database<Bytes, Bytes>,
169 metadata: Database<Bytes, Bytes>,
171 eviction_order: Database<Bytes, Unit>,
173 pins: Database<Bytes, Bytes>,
175 stats: Database<Bytes, Bytes>,
177 max_bytes: AtomicU64,
178 next_order: AtomicU64,
179 external_blobs: Option<ExternalBlobConfig>,
180}
181
182pub 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 pub fn new<P: AsRef<Path>>(path: P) -> Result<Self, StoreError> {
256 Self::with_map_size(path, DEFAULT_MAP_SIZE)
257 }
258
259 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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?); assert!(!store.put(hash, data.to_vec()).await?); 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}