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