1mod adaptive;
2mod catalog;
3mod gate;
4mod maintenance;
5mod member;
6mod model;
7mod temperature;
8mod temperature_balancer;
9mod temperature_catalog;
10mod temperature_worker;
11#[cfg(test)]
12mod tests;
13
14use self::adaptive::AdaptivePoolState;
15use self::gate::ConcurrencyGate;
16use self::member::{open_member_store, prepare_member_paths, validate_member_config};
17use self::model::{LocationRecord, MemberRecord, PoolManifest, MIN_MEMBER_MAP_SIZE_BYTES};
18pub use self::model::{
19 PoolMaintenanceReport, PoolMemberConfig, PoolMemberId, PoolMemberState, PoolMemberStatus,
20 PoolStoreConfig, PoolTemperatureConfig, PoolTemperatureReport,
21};
22use self::temperature::TemperatureRuntime;
23use self::temperature_worker::TemperatureWorker;
24use crate::{managed_env::ManagedEnv, LmdbBlobStore};
25use async_trait::async_trait;
26use hashtree_core::store::{slice_blob_range, PutManyReport, Store, StoreError, StoreStats};
27use hashtree_core::{sha256, types::Hash};
28use heed::types::{Bytes, Unit};
29use heed::{Database, EnvOpenOptions};
30use std::collections::{HashMap, HashSet};
31use std::fs;
32use std::ops::Deref;
33use std::path::Path;
34use std::sync::atomic::{AtomicU64, Ordering};
35use std::sync::{Arc, Mutex, RwLock};
36use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
37
38const CATALOG_DATABASES: u32 = 6;
39const CATALOG_MAX_READERS: u32 = 1024;
40const MANIFEST_KEY: &[u8] = b"pool-manifest-v1";
41const MEMBER_MARKER_NAME: &str = ".hashtree-pool-member-v1";
42const EXTERNAL_MARKER_NAME: &str = ".hashtree-pool-external-v1";
43
44#[derive(Default)]
45struct RuntimeMembers {
46 generation: Option<u64>,
47 stores: HashMap<PoolMemberId, Arc<LmdbBlobStore>>,
48 read_gates: HashMap<PoolMemberId, Arc<ConcurrencyGate>>,
49 write_gates: HashMap<PoolMemberId, Arc<ConcurrencyGate>>,
50 errors: HashMap<PoolMemberId, String>,
51}
52
53#[derive(Clone)]
54pub struct PoolStore {
55 inner: Arc<PoolStoreInner>,
56}
57
58#[doc(hidden)]
59pub struct PoolStoreInner {
60 env: ManagedEnv,
61 manifest_db: Database<Bytes, Bytes>,
62 locations: Database<Bytes, Bytes>,
63 by_member: Database<Bytes, Unit>,
64 pins: Database<Bytes, Bytes>,
65 last_accessed: Database<Bytes, Bytes>,
66 temperature_state: Database<Bytes, Bytes>,
67 runtime: RwLock<RuntimeMembers>,
68 adaptive: Mutex<AdaptivePoolState>,
69 temperature_config: PoolTemperatureConfig,
70 temperature: Mutex<TemperatureRuntime>,
71 temperature_access_counter: AtomicU64,
72 temperature_owner: PoolMemberId,
73 temperature_cycle: Mutex<()>,
74 temperature_worker: TemperatureWorker,
75}
76
77impl Deref for PoolStore {
78 type Target = PoolStoreInner;
79
80 fn deref(&self) -> &Self::Target {
81 &self.inner
82 }
83}
84
85fn validate_new_member_paths(
86 manifest: &PoolManifest,
87 config: &PoolMemberConfig,
88) -> Result<(), StoreError> {
89 if manifest
90 .members
91 .iter()
92 .any(|member| member.config.path == config.path)
93 {
94 return Err(StoreError::Other(format!(
95 "pool member path is already configured: {}",
96 config.path.display()
97 )));
98 }
99 if let Some(external) = config.external_blob_dir.as_ref() {
100 if manifest
101 .members
102 .iter()
103 .any(|member| member.config.external_blob_dir.as_ref() == Some(external))
104 {
105 return Err(StoreError::Other(format!(
106 "pool external blob path is already configured: {}",
107 external.display()
108 )));
109 }
110 }
111 Ok(())
112}
113
114impl PoolStore {
115 pub fn open<P: AsRef<Path>>(path: P, config: PoolStoreConfig) -> Result<Self, StoreError> {
116 config.temperature.validate()?;
117 let path = path.as_ref();
118 fs::create_dir_all(path).map_err(StoreError::Io)?;
119 let existing_size = fs::metadata(path.join("data.mdb"))
120 .map(|metadata| metadata.len())
121 .unwrap_or(0);
122 let requested = config
123 .catalog_map_size_bytes
124 .max(existing_size.saturating_add(existing_size / 10))
125 .max(MIN_MEMBER_MAP_SIZE_BYTES);
126 let map_size = usize::try_from(requested)
127 .map_err(|_| StoreError::Other("pool catalog map size exceeds usize".into()))?;
128
129 let mut options = EnvOpenOptions::new();
130 options
131 .map_size(map_size)
132 .max_dbs(CATALOG_DATABASES)
133 .max_readers(CATALOG_MAX_READERS);
134 unsafe {
135 options.flags(super::env_flags_from_env());
136 }
137 let env = unsafe { ManagedEnv::open(&options, path) }.map_err(|error| {
138 StoreError::Other(format!("open pool catalog {}: {error}", path.display()))
139 })?;
140 let _ = env.clear_stale_readers();
141 if env.info().map_size < map_size {
142 unsafe { env.resize(map_size) }.map_err(map_heed)?;
143 }
144
145 let mut wtxn = env.write_txn().map_err(map_heed)?;
146 let manifest_db = env
147 .create_database(&mut wtxn, Some("manifest"))
148 .map_err(map_heed)?;
149 let locations = env
150 .create_database(&mut wtxn, Some("locations"))
151 .map_err(map_heed)?;
152 let by_member = env
153 .create_database(&mut wtxn, Some("by_member"))
154 .map_err(map_heed)?;
155 let pins = env
156 .create_database(&mut wtxn, Some("pins"))
157 .map_err(map_heed)?;
158 let last_accessed = env
159 .create_database(&mut wtxn, Some("last_accessed"))
160 .map_err(map_heed)?;
161 let temperature_state = env
162 .create_database(&mut wtxn, Some("temperature_state"))
163 .map_err(map_heed)?;
164 if manifest_db
165 .get(&wtxn, MANIFEST_KEY)
166 .map_err(map_heed)?
167 .is_none()
168 {
169 let bytes = encode_manifest(&PoolManifest::default())?;
170 manifest_db
171 .put(&mut wtxn, MANIFEST_KEY, bytes.as_slice())
172 .map_err(map_heed)?;
173 }
174 wtxn.commit().map_err(map_heed)?;
175
176 let temperature_config = config.temperature.clone();
177 let store = Self {
178 inner: Arc::new(PoolStoreInner {
179 env,
180 manifest_db,
181 locations,
182 by_member,
183 pins,
184 last_accessed,
185 temperature_state,
186 runtime: RwLock::new(RuntimeMembers::default()),
187 adaptive: Mutex::new(AdaptivePoolState::new(config.member_failure_cooldown)),
188 temperature: Mutex::new(TemperatureRuntime::new(
189 temperature_config.candidate_capacity,
190 )),
191 temperature_config,
192 temperature_access_counter: AtomicU64::new(0),
193 temperature_owner: PoolMemberId::new(),
194 temperature_cycle: Mutex::new(()),
195 temperature_worker: TemperatureWorker::default(),
196 }),
197 };
198 store.refresh_members()?;
199 store.start_temperature_worker()?;
200 Ok(store)
201 }
202
203 fn start_temperature_worker(&self) -> Result<(), StoreError> {
204 if !self.temperature_config.enabled {
205 return Ok(());
206 }
207 self.temperature_worker.start(
208 Arc::downgrade(&self.inner),
209 self.temperature_config.interval,
210 )
211 }
212
213 pub fn add_member(&self, config: PoolMemberConfig) -> Result<PoolMemberId, StoreError> {
214 self.add_member_inner(config, false)?
215 .ok_or_else(|| StoreError::Other("pool member was not added".into()))
216 }
217
218 pub(crate) fn ensure_initial_member(
219 &self,
220 config: PoolMemberConfig,
221 ) -> Result<Option<PoolMemberId>, StoreError> {
222 if !self.read_manifest()?.members.is_empty() {
223 self.refresh_members()?;
224 return Ok(None);
225 }
226 self.add_member_inner(config, true)
227 }
228
229 fn add_member_inner(
230 &self,
231 config: PoolMemberConfig,
232 only_if_empty: bool,
233 ) -> Result<Option<PoolMemberId>, StoreError> {
234 validate_member_config(&config)?;
235 let wtxn = self.env.write_txn().map_err(map_heed)?;
236 let manifest = self.manifest_from_txn(&wtxn)?;
237 if only_if_empty && !manifest.members.is_empty() {
238 drop(wtxn);
239 self.refresh_members()?;
240 return Ok(None);
241 }
242 validate_new_member_paths(&manifest, &config)?;
243 let id = prepare_member_paths(&config, PoolMemberId::new())?;
244 drop(wtxn);
245 let store = Arc::new(open_member_store(id, &config)?);
246
247 let mut wtxn = self.env.write_txn().map_err(map_heed)?;
248 let mut manifest = self.manifest_from_txn(&wtxn)?;
249 if only_if_empty && !manifest.members.is_empty() {
250 drop(wtxn);
251 self.refresh_members()?;
252 return Ok(None);
253 }
254 validate_new_member_paths(&manifest, &config)?;
255 if manifest.members.iter().any(|member| member.id == id) {
256 return Err(StoreError::Other(format!(
257 "pool member identity is already configured: {id}"
258 )));
259 }
260 manifest.members.push(MemberRecord {
261 id,
262 state: PoolMemberState::Active,
263 config,
264 });
265 manifest.generation = manifest.generation.saturating_add(1);
266 self.put_manifest_txn(&mut wtxn, &manifest)?;
267 wtxn.commit().map_err(map_heed)?;
268
269 let mut runtime = self
270 .runtime
271 .write()
272 .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?;
273 runtime.generation = Some(manifest.generation);
274 runtime.stores.insert(id, store);
275 let member = manifest
276 .members
277 .iter()
278 .find(|member| member.id == id)
279 .expect("new pool member is in committed manifest");
280 runtime.read_gates.insert(
281 id,
282 Arc::new(ConcurrencyGate::new(member.config.max_read_concurrency)),
283 );
284 runtime.write_gates.insert(
285 id,
286 Arc::new(ConcurrencyGate::new(member.config.max_write_concurrency)),
287 );
288 runtime.errors.remove(&id);
289 Ok(Some(id))
290 }
291
292 pub fn begin_drain(&self, id: PoolMemberId) -> Result<(), StoreError> {
293 let mut wtxn = self.env.write_txn().map_err(map_heed)?;
294 let mut manifest = self.manifest_from_txn(&wtxn)?;
295 let has_other_active = manifest
296 .members
297 .iter()
298 .any(|member| member.id != id && member.state == PoolMemberState::Active);
299 let member = manifest
300 .members
301 .iter_mut()
302 .find(|member| member.id == id)
303 .ok_or_else(|| StoreError::Other(format!("unknown pool member {id}")))?;
304 if member.state == PoolMemberState::Draining {
305 return Ok(());
306 }
307 if !has_other_active && self.count_member_locations_txn(&wtxn, id)? > 0 {
308 return Err(StoreError::Other(
309 "cannot drain the final member while it still owns blobs".into(),
310 ));
311 }
312 member.state = PoolMemberState::Draining;
313 manifest.generation = manifest.generation.saturating_add(1);
314 self.put_manifest_txn(&mut wtxn, &manifest)?;
315 wtxn.commit().map_err(map_heed)?;
316 self.refresh_members()?;
317 Ok(())
318 }
319
320 pub fn update_member_limits(
321 &self,
322 id: PoolMemberId,
323 capacity_bytes: u64,
324 max_read_concurrency: u32,
325 max_write_concurrency: u32,
326 ) -> Result<(), StoreError> {
327 if capacity_bytes == 0 || max_read_concurrency == 0 || max_write_concurrency == 0 {
328 return Err(StoreError::Other(
329 "pool member capacity and concurrency limits must be non-zero".into(),
330 ));
331 }
332 let mut wtxn = self.env.write_txn().map_err(map_heed)?;
333 let mut manifest = self.manifest_from_txn(&wtxn)?;
334 let member = manifest
335 .members
336 .iter_mut()
337 .find(|member| member.id == id)
338 .ok_or_else(|| StoreError::Other(format!("unknown pool member {id}")))?;
339 member.config.capacity_bytes = capacity_bytes;
340 member.config.max_read_concurrency = max_read_concurrency;
341 member.config.max_write_concurrency = max_write_concurrency;
342 manifest.generation = manifest.generation.saturating_add(1);
343 self.put_manifest_txn(&mut wtxn, &manifest)?;
344 wtxn.commit().map_err(map_heed)?;
345 self.refresh_members()
346 }
347
348 pub fn update_member_temperature_watermarks(
349 &self,
350 id: PoolMemberId,
351 low_percent: u8,
352 high_percent: u8,
353 ) -> Result<(), StoreError> {
354 if low_percent >= high_percent || high_percent > 100 {
355 return Err(StoreError::Other(
356 "pool temperature watermarks must satisfy 0 <= low < high <= 100".into(),
357 ));
358 }
359 let mut wtxn = self.env.write_txn().map_err(map_heed)?;
360 let mut manifest = self.manifest_from_txn(&wtxn)?;
361 let member = manifest
362 .members
363 .iter_mut()
364 .find(|member| member.id == id)
365 .ok_or_else(|| StoreError::Other(format!("unknown pool member {id}")))?;
366 member.config.temperature_low_watermark_percent = low_percent;
367 member.config.temperature_high_watermark_percent = high_percent;
368 manifest.generation = manifest.generation.saturating_add(1);
369 self.put_manifest_txn(&mut wtxn, &manifest)?;
370 wtxn.commit().map_err(map_heed)?;
371 self.refresh_members()
372 }
373
374 pub fn remove_member(&self, id: PoolMemberId) -> Result<(), StoreError> {
375 let mut wtxn = self.env.write_txn().map_err(map_heed)?;
376 let mut manifest = self.manifest_from_txn(&wtxn)?;
377 let index = manifest
378 .members
379 .iter()
380 .position(|member| member.id == id)
381 .ok_or_else(|| StoreError::Other(format!("unknown pool member {id}")))?;
382 if manifest.members[index].state != PoolMemberState::Draining {
383 return Err(StoreError::Other(
384 "pool member must be draining before removal".into(),
385 ));
386 }
387 let located = self.count_member_locations_txn(&wtxn, id)?;
388 if located != 0 {
389 return Err(StoreError::Other(format!(
390 "pool member {id} still owns {located} blob(s)"
391 )));
392 }
393 manifest.members.remove(index);
394 manifest.generation = manifest.generation.saturating_add(1);
395 self.put_manifest_txn(&mut wtxn, &manifest)?;
396 wtxn.commit().map_err(map_heed)?;
397 self.refresh_members()?;
398 Ok(())
399 }
400
401 pub fn member(&self, id: PoolMemberId) -> Result<PoolMemberStatus, StoreError> {
402 self.refresh_members()?;
403 let manifest = self.read_manifest()?;
404 let member = manifest
405 .members
406 .iter()
407 .find(|member| member.id == id)
408 .ok_or_else(|| StoreError::Other(format!("unknown pool member {id}")))?;
409 let located_blobs = self.count_member_locations(id)?;
410 let runtime = self
411 .runtime
412 .read()
413 .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?;
414 let (logical_bytes, available, last_error) = match runtime.stores.get(&id) {
415 Some(store) => match store.stats() {
416 Ok(stats) => (stats.total_bytes, true, None),
417 Err(error) => (0, false, Some(error.to_string())),
418 },
419 None => (0, false, runtime.errors.get(&id).cloned()),
420 };
421 Ok(PoolMemberStatus {
422 id,
423 state: member.state,
424 path: member.config.path.clone(),
425 capacity_bytes: member.config.capacity_bytes,
426 map_size_bytes: member.config.map_size_bytes,
427 external_blob_dir: member.config.external_blob_dir.clone(),
428 external_blob_min_bytes: member.config.external_blob_min_bytes,
429 external_blob_sync: member.config.external_blob_sync,
430 external_pack_target_bytes: member.config.external_pack_target_bytes,
431 max_read_concurrency: member.config.max_read_concurrency,
432 max_write_concurrency: member.config.max_write_concurrency,
433 temperature_low_watermark_percent: member.config.temperature_low_watermark_percent,
434 temperature_high_watermark_percent: member.config.temperature_high_watermark_percent,
435 logical_bytes,
436 located_blobs,
437 available,
438 last_error,
439 })
440 }
441
442 pub fn members(&self) -> Result<Vec<PoolMemberStatus>, StoreError> {
443 let manifest = self.read_manifest()?;
444 manifest
445 .members
446 .iter()
447 .map(|member| self.member(member.id))
448 .collect()
449 }
450
451 pub fn blob_location(&self, hash: &Hash) -> Result<Option<PoolMemberId>, StoreError> {
452 Ok(self
453 .read_location(hash)?
454 .map(LocationRecord::preferred_member))
455 }
456
457 pub fn put_sync(&self, hash: Hash, data: &[u8]) -> Result<bool, StoreError> {
458 if sha256(data) != hash {
459 return Err(StoreError::Other(
460 "pool rejected bytes that do not match their hash".into(),
461 ));
462 }
463
464 if let Some(location) = self.read_location(&hash)? {
465 match self.read_verified_location(&hash, location) {
466 Ok(Some(found)) => {
467 if matches!(location, LocationRecord::Pending { .. }) {
468 self.finalize_pending(hash, location)?;
469 }
470 debug_assert_eq!(sha256(&found), hash);
471 return Ok(false);
472 }
473 Ok(None) | Err(_) => return self.repair_location(hash, data, location),
474 }
475 }
476
477 let target = self.choose_write_member(data.len() as u64, None)?;
478 let pending = LocationRecord::Pending {
479 member: target,
480 size: data.len() as u64,
481 };
482 let location = self.reserve_if_absent(hash, pending)?;
483 if location != pending {
484 match self.read_verified_location(&hash, location) {
485 Ok(Some(found)) => {
486 if matches!(location, LocationRecord::Pending { .. }) {
487 self.finalize_pending(hash, location)?;
488 }
489 debug_assert_eq!(sha256(&found), hash);
490 return Ok(false);
491 }
492 Ok(None) | Err(_) => return self.repair_location(hash, data, location),
493 }
494 }
495
496 let target = location.preferred_member();
497 let store = self.get_member(target)?;
498 match self.write_verified_member(target, &store, hash, data) {
499 Ok(inserted) => {
500 self.finalize_pending(hash, location)?;
501 Ok(inserted)
502 }
503 Err(_) => {
504 let mut excluded = HashSet::new();
505 excluded.insert(target);
506 self.repair_location_excluding(hash, data, location, excluded)
507 }
508 }
509 }
510
511 pub fn put_many_report_sync(
512 &self,
513 items: &[(Hash, Vec<u8>)],
514 ) -> Result<PutManyReport, StoreError> {
515 self.put_many_report_sync_with_existing_verification(items, true)
516 }
517
518 pub fn put_many_optimistic_report_sync(
522 &self,
523 items: &[(Hash, Vec<u8>)],
524 ) -> Result<PutManyReport, StoreError> {
525 self.put_many_report_sync_with_existing_verification(items, false)
526 }
527
528 fn put_many_report_sync_with_existing_verification(
529 &self,
530 items: &[(Hash, Vec<u8>)],
531 verify_existing: bool,
532 ) -> Result<PutManyReport, StoreError> {
533 let mut seen = HashSet::new();
534 let mut unique = Vec::with_capacity(items.len());
535 let mut ordered = Vec::with_capacity(items.len());
536 for (hash, data) in items {
537 if sha256(data) != *hash {
538 return Err(StoreError::Other(
539 "pool rejected batch bytes that do not match their hash".into(),
540 ));
541 }
542 if seen.insert(*hash) {
543 unique.push((*hash, data));
544 ordered.push((*hash, data.len() as u64));
545 }
546 }
547
548 let mut inserted = HashSet::new();
549 let mut missing = Vec::new();
550 for (hash, data) in unique {
551 if let Some(location) = self.read_location(&hash)? {
552 if verify_existing || matches!(location, LocationRecord::Pending { .. }) {
553 if self.put_sync(hash, data)? {
554 inserted.insert(hash);
555 }
556 }
557 } else {
558 missing.push((hash, data));
559 }
560 }
561 if missing.is_empty() {
562 return Ok(put_many_report(items.len(), &ordered, &inserted));
563 }
564
565 let mut reserved_bytes = HashMap::new();
566 let mut assignments = Vec::with_capacity(missing.len());
567 for (hash, data) in missing {
568 let target =
569 self.choose_write_member_with_reserved(data.len() as u64, None, &reserved_bytes)?;
570 let reserved = reserved_bytes.entry(target).or_insert(0u64);
571 *reserved = reserved.saturating_add(data.len() as u64);
572 assignments.push((hash, data, target));
573 }
574
575 let mut wtxn = self.env.write_txn().map_err(map_heed)?;
576 let mut plans = Vec::with_capacity(assignments.len());
577 let mut raced = Vec::new();
578 for (hash, data, target) in assignments {
579 if self
580 .locations
581 .get(&wtxn, &hash)
582 .map_err(map_heed)?
583 .is_some()
584 {
585 raced.push((hash, data));
586 continue;
587 }
588 let pending = LocationRecord::Pending {
589 member: target,
590 size: data.len() as u64,
591 };
592 self.set_location_txn(&mut wtxn, hash, Some(pending))?;
593 plans.push((hash, data, target, pending));
594 }
595 wtxn.commit().map_err(map_heed)?;
596
597 for (hash, data) in raced {
598 if self.put_sync(hash, data)? {
599 inserted.insert(hash);
600 }
601 }
602
603 let mut by_target: HashMap<PoolMemberId, Vec<(Hash, &[u8])>> = HashMap::new();
604 for (hash, data, target, _) in &plans {
605 by_target
606 .entry(*target)
607 .or_default()
608 .push((*hash, data.as_slice()));
609 }
610 for (target, batch) in by_target {
611 let store = self.get_member(target)?;
612 let gate = self.member_gate(target, true)?;
613 let permit = gate.acquire()?;
614 for (hash, _) in &batch {
615 if store
616 .get_sync(hash)?
617 .is_some_and(|existing| sha256(&existing) != *hash)
618 {
619 store.delete_sync(hash)?;
620 }
621 }
622 let started = Instant::now();
623 let result = store.put_many_refs_report_sync(&batch);
624 let success = result.is_ok();
625 let bytes = batch.iter().map(|(_, data)| data.len()).sum::<usize>();
626 self.record_write(target, started.elapsed(), bytes, success);
627 let report = match result {
628 Ok(report) => report,
629 Err(_) => {
630 drop(permit);
631 for (hash, data) in batch {
632 if self.put_sync(hash, data)? {
633 inserted.insert(hash);
634 }
635 }
636 continue;
637 }
638 };
639 inserted.extend(report.inserted_hashes);
640 for (hash, _) in &batch {
641 self.read_verified_member(target, &store, hash)?
642 .ok_or_else(|| {
643 StoreError::Other(format!(
644 "pool member {target} lost a committed batch write"
645 ))
646 })?;
647 }
648 }
649
650 let mut wtxn = self.env.write_txn().map_err(map_heed)?;
651 for (hash, _, target, pending) in plans {
652 let Some(current) = self.locations.get(&wtxn, &hash).map_err(map_heed)? else {
653 continue;
654 };
655 if LocationRecord::decode(current)? == pending {
656 self.set_location_txn(
657 &mut wtxn,
658 hash,
659 Some(LocationRecord::Stored {
660 member: target,
661 size: pending.size(),
662 }),
663 )?;
664 }
665 }
666 wtxn.commit().map_err(map_heed)?;
667 Ok(put_many_report(items.len(), &ordered, &inserted))
668 }
669
670 pub fn put_many_sync(&self, items: &[(Hash, Vec<u8>)]) -> Result<usize, StoreError> {
671 self.put_many_report_sync(items)
672 .map(|report| report.inserted)
673 }
674
675 pub fn put_many_optimistic_sync(&self, items: &[(Hash, Vec<u8>)]) -> Result<usize, StoreError> {
676 self.put_many_optimistic_report_sync(items)
677 .map(|report| report.inserted)
678 }
679
680 pub fn get_sync(&self, hash: &Hash) -> Result<Option<Vec<u8>>, StoreError> {
681 let Some(location) = self.read_location(hash)? else {
682 return Ok(None);
683 };
684 let data = self.read_verified_location(hash, location)?;
685 if data.is_some() && matches!(location, LocationRecord::Pending { .. }) {
686 self.finalize_pending(*hash, location)?;
687 }
688 if data.is_some() {
689 self.sample_temperature_access(*hash, location);
690 }
691 Ok(data)
692 }
693
694 pub fn get_range_sync(
695 &self,
696 hash: &Hash,
697 start: u64,
698 end_inclusive: u64,
699 ) -> Result<Option<Vec<u8>>, StoreError> {
700 let Some(data) = self.get_sync(hash)? else {
701 return Ok(None);
702 };
703 slice_blob_range(&data, start, end_inclusive).map(Some)
704 }
705
706 pub fn blob_size_sync(&self, hash: &Hash) -> Result<Option<u64>, StoreError> {
707 Ok(self.read_location(hash)?.map(LocationRecord::size))
708 }
709
710 pub fn exists(&self, hash: &Hash) -> Result<bool, StoreError> {
711 Ok(self.get_sync(hash)?.is_some())
712 }
713
714 pub fn existing_hashes_in_sorted_candidates(
715 &self,
716 sorted_hashes: &[Hash],
717 ) -> Result<Vec<bool>, StoreError> {
718 let rtxn = self.env.read_txn().map_err(map_heed)?;
719 sorted_hashes
720 .iter()
721 .map(|hash| self.locations.get(&rtxn, hash).map(|value| value.is_some()))
722 .collect::<Result<Vec<_>, _>>()
723 .map_err(map_heed)
724 }
725
726 pub fn largest_member_map_size_bytes(&self) -> Result<Option<usize>, StoreError> {
731 self.refresh_members()?;
732 let runtime = self
733 .runtime
734 .read()
735 .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?;
736 Ok(runtime
737 .stores
738 .values()
739 .map(|store| store.map_size_bytes())
740 .max())
741 }
742
743 pub fn delete_sync(&self, hash: &Hash) -> Result<bool, StoreError> {
744 let Some(location) = self.read_location(hash)? else {
745 return Ok(false);
746 };
747 let (members, len) = location.members();
748 let mut deleted = false;
749 for member in members.into_iter().take(len) {
750 if let Ok(store) = self.get_member(member) {
751 let gate = self.member_gate(member, true)?;
752 let _permit = gate.acquire()?;
753 deleted |= store.delete_sync(hash)?;
754 }
755 }
756 let mut wtxn = self.env.write_txn().map_err(map_heed)?;
757 self.set_location_txn(&mut wtxn, *hash, None)?;
758 self.pins.delete(&mut wtxn, hash).map_err(map_heed)?;
759 wtxn.commit().map_err(map_heed)?;
760 Ok(deleted)
761 }
762
763 pub fn pin_sync(&self, hash: &Hash) -> Result<(), StoreError> {
764 let mut wtxn = self.env.write_txn().map_err(map_heed)?;
765 let previous = self
766 .pins
767 .get(&wtxn, hash)
768 .map_err(map_heed)?
769 .map(decode_pin_count)
770 .transpose()?
771 .unwrap_or(0);
772 self.pins
773 .put(&mut wtxn, hash, &previous.saturating_add(1).to_be_bytes())
774 .map_err(map_heed)?;
775 wtxn.commit().map_err(map_heed)
776 }
777
778 pub fn unpin_sync(&self, hash: &Hash) -> Result<(), StoreError> {
779 let mut wtxn = self.env.write_txn().map_err(map_heed)?;
780 let count = self
781 .pins
782 .get(&wtxn, hash)
783 .map_err(map_heed)?
784 .map(decode_pin_count)
785 .transpose()?
786 .unwrap_or(0);
787 if count <= 1 {
788 self.pins.delete(&mut wtxn, hash).map_err(map_heed)?;
789 } else {
790 self.pins
791 .put(&mut wtxn, hash, &(count - 1).to_be_bytes())
792 .map_err(map_heed)?;
793 }
794 wtxn.commit().map_err(map_heed)
795 }
796
797 pub fn pin_count_sync(&self, hash: &Hash) -> Result<u32, StoreError> {
798 let rtxn = self.env.read_txn().map_err(map_heed)?;
799 self.pins
800 .get(&rtxn, hash)
801 .map_err(map_heed)?
802 .map(decode_pin_count)
803 .transpose()
804 .map(|count| count.unwrap_or(0))
805 }
806
807 pub fn touch_accessed_sync(&self, hash: &Hash, now: u64) -> Result<bool, StoreError> {
808 let mut wtxn = self.env.write_txn().map_err(map_heed)?;
809 if self.locations.get(&wtxn, hash).map_err(map_heed)?.is_none() {
810 return Ok(false);
811 }
812 let previous = self
813 .last_accessed
814 .get(&wtxn, hash)
815 .map_err(map_heed)?
816 .and_then(self::temperature::AccessRecord::decode)
817 .unwrap_or_else(|| self::temperature::AccessRecord::new(now));
818 let access = self::temperature::AccessRecord {
819 last_accessed_at: now,
820 ..previous
821 }
822 .encode();
823 self.last_accessed
824 .put(&mut wtxn, hash, &access)
825 .map_err(map_heed)?;
826 wtxn.commit().map_err(map_heed)?;
827 Ok(true)
828 }
829
830 pub fn touch_many_accessed_sync(&self, hashes: &[Hash], now: u64) -> Result<usize, StoreError> {
831 let mut wtxn = self.env.write_txn().map_err(map_heed)?;
832 let mut updated = 0usize;
833 let mut seen = HashSet::new();
834 for hash in hashes {
835 if !seen.insert(*hash) || self.locations.get(&wtxn, hash).map_err(map_heed)?.is_none() {
836 continue;
837 }
838 let previous = self
839 .last_accessed
840 .get(&wtxn, hash)
841 .map_err(map_heed)?
842 .and_then(self::temperature::AccessRecord::decode)
843 .unwrap_or_else(|| self::temperature::AccessRecord::new(now));
844 let access = self::temperature::AccessRecord {
845 last_accessed_at: now,
846 ..previous
847 }
848 .encode();
849 self.last_accessed
850 .put(&mut wtxn, hash, &access)
851 .map_err(map_heed)?;
852 updated += 1;
853 }
854 wtxn.commit().map_err(map_heed)?;
855 Ok(updated)
856 }
857
858 pub fn last_accessed_at_sync(&self, hash: &Hash) -> Result<Option<u64>, StoreError> {
859 let rtxn = self.env.read_txn().map_err(map_heed)?;
860 self.last_accessed
861 .get(&rtxn, hash)
862 .map_err(map_heed)?
863 .map(|bytes| {
864 self::temperature::AccessRecord::decode(bytes)
865 .map(|access| access.last_accessed_at)
866 .ok_or_else(|| StoreError::Other("invalid pool access record".into()))
867 })
868 .transpose()
869 }
870
871 pub fn many_last_accessed_at_sync(
872 &self,
873 hashes: &[Hash],
874 ) -> Result<Vec<(Hash, u64)>, StoreError> {
875 let rtxn = self.env.read_txn().map_err(map_heed)?;
876 let mut values = Vec::new();
877 for hash in hashes {
878 if let Some(value) = self.last_accessed.get(&rtxn, hash).map_err(map_heed)? {
879 let access = self::temperature::AccessRecord::decode(value)
880 .ok_or_else(|| StoreError::Other("invalid pool access record".into()))?;
881 values.push((*hash, access.last_accessed_at));
882 }
883 }
884 Ok(values)
885 }
886
887 pub fn list(&self) -> Result<Vec<Hash>, StoreError> {
888 let rtxn = self.env.read_txn().map_err(map_heed)?;
889 let mut hashes = Vec::new();
890 for item in self.locations.iter(&rtxn).map_err(map_heed)? {
891 let (hash, _) = item.map_err(map_heed)?;
892 let hash: Hash = hash
893 .try_into()
894 .map_err(|_| StoreError::Other("invalid pool hash key".into()))?;
895 hashes.push(hash);
896 }
897 Ok(hashes)
898 }
899
900 pub fn stats(&self) -> Result<StoreStats, StoreError> {
901 let rtxn = self.env.read_txn().map_err(map_heed)?;
902 let mut stats = StoreStats::default();
903 for item in self.locations.iter(&rtxn).map_err(map_heed)? {
904 let (_, location) = item.map_err(map_heed)?;
905 let location = LocationRecord::decode(location)?;
906 stats.count = stats.count.saturating_add(1);
907 stats.bytes = stats.bytes.saturating_add(location.size());
908 }
909 for item in self.pins.iter(&rtxn).map_err(map_heed)? {
910 let (hash, count) = item.map_err(map_heed)?;
911 if decode_pin_count(count)? == 0 {
912 continue;
913 }
914 let Some(location) = self.locations.get(&rtxn, hash).map_err(map_heed)? else {
915 continue;
916 };
917 stats.pinned_count = stats.pinned_count.saturating_add(1);
918 stats.pinned_bytes = stats
919 .pinned_bytes
920 .saturating_add(LocationRecord::decode(location)?.size());
921 }
922 Ok(stats)
923 }
924
925 pub fn force_sync(&self) -> Result<(), StoreError> {
926 self.env.force_sync().map_err(map_heed)?;
927 self.refresh_members()?;
928 let runtime = self
929 .runtime
930 .read()
931 .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?;
932 for store in runtime.stores.values() {
933 store.force_sync()?;
934 }
935 Ok(())
936 }
937
938 fn refresh_members(&self) -> Result<(), StoreError> {
939 let manifest = self.read_manifest()?;
940 let mut runtime = self
941 .runtime
942 .write()
943 .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?;
944 if runtime.generation == Some(manifest.generation) {
945 return Ok(());
946 }
947
948 let configured = manifest
949 .members
950 .iter()
951 .map(|member| member.id)
952 .collect::<HashSet<_>>();
953 runtime.stores.retain(|id, _| configured.contains(id));
954 runtime.read_gates.retain(|id, _| configured.contains(id));
955 runtime.write_gates.retain(|id, _| configured.contains(id));
956 runtime.errors.retain(|id, _| configured.contains(id));
957 for member in &manifest.members {
958 let read_gate = runtime
959 .read_gates
960 .entry(member.id)
961 .or_insert_with(|| {
962 Arc::new(ConcurrencyGate::new(member.config.max_read_concurrency))
963 })
964 .clone();
965 read_gate.set_limit(member.config.max_read_concurrency)?;
966 let write_gate = runtime
967 .write_gates
968 .entry(member.id)
969 .or_insert_with(|| {
970 Arc::new(ConcurrencyGate::new(member.config.max_write_concurrency))
971 })
972 .clone();
973 write_gate.set_limit(member.config.max_write_concurrency)?;
974 if runtime.stores.contains_key(&member.id) {
975 continue;
976 }
977 match open_member_store(member.id, &member.config) {
978 Ok(store) => {
979 runtime.stores.insert(member.id, Arc::new(store));
980 runtime.errors.remove(&member.id);
981 }
982 Err(error) => {
983 runtime.errors.insert(member.id, error.to_string());
984 }
985 }
986 }
987 runtime.generation = Some(manifest.generation);
988 drop(runtime);
989 self.adaptive
990 .lock()
991 .map_err(|_| StoreError::Other("pool adaptive lock poisoned".into()))?
992 .retain(&configured);
993 Ok(())
994 }
995
996 fn get_member(&self, id: PoolMemberId) -> Result<Arc<LmdbBlobStore>, StoreError> {
997 self.refresh_members()?;
998 if let Some(store) = self
999 .runtime
1000 .read()
1001 .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?
1002 .stores
1003 .get(&id)
1004 .cloned()
1005 {
1006 return Ok(store);
1007 }
1008
1009 let member = self
1010 .read_manifest()?
1011 .members
1012 .into_iter()
1013 .find(|member| member.id == id)
1014 .ok_or_else(|| StoreError::Other(format!("unknown pool member {id}")))?;
1015 match open_member_store(id, &member.config) {
1016 Ok(store) => {
1017 let store = Arc::new(store);
1018 let mut runtime = self
1019 .runtime
1020 .write()
1021 .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?;
1022 runtime.stores.insert(id, Arc::clone(&store));
1023 runtime.errors.remove(&id);
1024 Ok(store)
1025 }
1026 Err(error) => {
1027 self.record_member_failure(id, false);
1028 let mut runtime = self
1029 .runtime
1030 .write()
1031 .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?;
1032 runtime.errors.insert(id, error.to_string());
1033 Err(error)
1034 }
1035 }
1036 }
1037
1038 fn member_state(&self, id: PoolMemberId) -> Result<Option<PoolMemberState>, StoreError> {
1039 Ok(self
1040 .read_manifest()?
1041 .members
1042 .into_iter()
1043 .find(|member| member.id == id)
1044 .map(|member| member.state))
1045 }
1046
1047 fn member_gate(
1048 &self,
1049 id: PoolMemberId,
1050 write: bool,
1051 ) -> Result<Arc<ConcurrencyGate>, StoreError> {
1052 self.refresh_members()?;
1053 let runtime = self
1054 .runtime
1055 .read()
1056 .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?;
1057 let gates = if write {
1058 &runtime.write_gates
1059 } else {
1060 &runtime.read_gates
1061 };
1062 gates
1063 .get(&id)
1064 .cloned()
1065 .ok_or_else(|| StoreError::Other(format!("unknown pool member {id}")))
1066 }
1067
1068 fn choose_write_member(
1069 &self,
1070 incoming_bytes: u64,
1071 exclude: Option<PoolMemberId>,
1072 ) -> Result<PoolMemberId, StoreError> {
1073 self.choose_write_member_with_reserved(incoming_bytes, exclude, &HashMap::new())
1074 }
1075
1076 fn choose_write_member_with_reserved(
1077 &self,
1078 incoming_bytes: u64,
1079 exclude: Option<PoolMemberId>,
1080 reserved_bytes: &HashMap<PoolMemberId, u64>,
1081 ) -> Result<PoolMemberId, StoreError> {
1082 let excluded = exclude.into_iter().collect::<HashSet<_>>();
1083 self.choose_write_member_excluding(incoming_bytes, &excluded, reserved_bytes)
1084 }
1085
1086 fn choose_write_member_excluding(
1087 &self,
1088 incoming_bytes: u64,
1089 excluded: &HashSet<PoolMemberId>,
1090 reserved_bytes: &HashMap<PoolMemberId, u64>,
1091 ) -> Result<PoolMemberId, StoreError> {
1092 self.refresh_members()?;
1093 let manifest = self.read_manifest()?;
1094 let runtime = self
1095 .runtime
1096 .read()
1097 .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?;
1098 let mut candidates = Vec::new();
1099 let mut below_high_watermark = Vec::new();
1100 for member in manifest.members.iter().filter(|member| {
1101 member.state == PoolMemberState::Active && !excluded.contains(&member.id)
1102 }) {
1103 let Some(store) = runtime.stores.get(&member.id) else {
1104 continue;
1105 };
1106 let stats = match store.stats() {
1107 Ok(stats) => stats,
1108 Err(_) => continue,
1109 };
1110 let effective_bytes = stats
1111 .total_bytes
1112 .saturating_add(reserved_bytes.get(&member.id).copied().unwrap_or(0));
1113 if member.config.capacity_bytes > 0
1114 && effective_bytes.saturating_add(incoming_bytes) > member.config.capacity_bytes
1115 {
1116 continue;
1117 }
1118 candidates.push((member.id, effective_bytes, member.config.capacity_bytes));
1119 let projected_fill = effective_bytes
1120 .saturating_add(incoming_bytes)
1121 .saturating_mul(100)
1122 .saturating_div(member.config.capacity_bytes)
1123 .min(100);
1124 if projected_fill <= u64::from(member.config.temperature_high_watermark_percent) {
1125 below_high_watermark.push((
1126 member.id,
1127 effective_bytes,
1128 member.config.capacity_bytes,
1129 ));
1130 }
1131 }
1132 drop(runtime);
1133 if !below_high_watermark.is_empty() {
1134 candidates = below_high_watermark;
1135 }
1136 self.adaptive
1137 .lock()
1138 .map_err(|_| StoreError::Other("pool adaptive lock poisoned".into()))?
1139 .choose_write(&candidates)
1140 .ok_or_else(|| StoreError::Other("no writable pool member has capacity".into()))
1141 }
1142
1143 fn repair_location(
1144 &self,
1145 hash: Hash,
1146 data: &[u8],
1147 expected: LocationRecord,
1148 ) -> Result<bool, StoreError> {
1149 self.repair_location_excluding(hash, data, expected, HashSet::new())
1150 }
1151
1152 fn repair_location_excluding(
1153 &self,
1154 hash: Hash,
1155 data: &[u8],
1156 expected: LocationRecord,
1157 mut excluded: HashSet<PoolMemberId>,
1158 ) -> Result<bool, StoreError> {
1159 let preferred = expected.preferred_member();
1160 let mut next = (!excluded.contains(&preferred)
1161 && self.member_state(preferred)? == Some(PoolMemberState::Active))
1162 .then_some(preferred);
1163 let mut last_error = None;
1164 let (target, inserted) = loop {
1165 let target = match next.take() {
1166 Some(target) => target,
1167 None => match self.choose_write_member_excluding(
1168 data.len() as u64,
1169 &excluded,
1170 &HashMap::new(),
1171 ) {
1172 Ok(target) => target,
1173 Err(error) => {
1174 return Err(last_error.unwrap_or(error));
1175 }
1176 },
1177 };
1178 let result = self
1179 .get_member(target)
1180 .and_then(|store| self.write_verified_member(target, &store, hash, data));
1181 match result {
1182 Ok(inserted) => break (target, inserted),
1183 Err(error) => {
1184 excluded.insert(target);
1185 last_error = Some(error);
1186 }
1187 }
1188 };
1189 let stored = LocationRecord::Stored {
1190 member: target,
1191 size: data.len() as u64,
1192 };
1193
1194 let mut wtxn = self.env.write_txn().map_err(map_heed)?;
1195 let current = self
1196 .locations
1197 .get(&wtxn, &hash)
1198 .map_err(map_heed)?
1199 .map(LocationRecord::decode)
1200 .transpose()?;
1201 if current == Some(expected) {
1202 self.set_location_txn(&mut wtxn, hash, Some(stored))?;
1203 wtxn.commit().map_err(map_heed)?;
1204 return Ok(inserted);
1205 }
1206 drop(wtxn);
1207
1208 if let Some(current) = current {
1209 if self.read_verified_location(&hash, current)?.is_some() {
1210 return Ok(false);
1211 }
1212 }
1213 Err(StoreError::Other(format!(
1214 "pool location changed while repairing {hash:?}"
1215 )))
1216 }
1217
1218 fn read_verified_location(
1219 &self,
1220 hash: &Hash,
1221 location: LocationRecord,
1222 ) -> Result<Option<Vec<u8>>, StoreError> {
1223 let mut ids = match location {
1224 LocationRecord::Pending { member, .. } | LocationRecord::Stored { member, .. } => {
1225 vec![member]
1226 }
1227 LocationRecord::Moving { source, target, .. } => vec![target, source],
1228 };
1229 self.adaptive
1230 .lock()
1231 .map_err(|_| StoreError::Other("pool adaptive lock poisoned".into()))?
1232 .order_reads(&mut ids);
1233
1234 let mut first_error = None;
1235 for id in ids {
1236 let store = match self.get_member(id) {
1237 Ok(store) => store,
1238 Err(error) => {
1239 first_error.get_or_insert(error);
1240 continue;
1241 }
1242 };
1243 match self.read_verified_member(id, &store, hash) {
1244 Ok(Some(data)) => return Ok(Some(data)),
1245 Ok(None) => {}
1246 Err(error) => {
1247 first_error.get_or_insert(error);
1248 }
1249 }
1250 }
1251 match first_error {
1252 Some(error) => Err(error),
1253 None => Ok(None),
1254 }
1255 }
1256
1257 fn read_verified_member(
1258 &self,
1259 id: PoolMemberId,
1260 store: &LmdbBlobStore,
1261 hash: &Hash,
1262 ) -> Result<Option<Vec<u8>>, StoreError> {
1263 let gate = self.member_gate(id, false)?;
1264 let _permit = gate.acquire()?;
1265 let started = Instant::now();
1266 let result = store.get_sync(hash);
1267 match result {
1268 Ok(Some(data)) if sha256(&data) == *hash => {
1269 self.record_read(id, started.elapsed(), true);
1270 Ok(Some(data))
1271 }
1272 Ok(Some(_)) => {
1273 self.record_read(id, started.elapsed(), false);
1274 Err(StoreError::Other(format!(
1275 "pool member {id} returned corrupt bytes"
1276 )))
1277 }
1278 Ok(None) => {
1279 self.record_read(id, started.elapsed(), true);
1280 Ok(None)
1281 }
1282 Err(error) => {
1283 self.record_read(id, started.elapsed(), false);
1284 Err(error)
1285 }
1286 }
1287 }
1288
1289 fn write_verified_member(
1290 &self,
1291 id: PoolMemberId,
1292 store: &LmdbBlobStore,
1293 hash: Hash,
1294 data: &[u8],
1295 ) -> Result<bool, StoreError> {
1296 let gate = self.member_gate(id, true)?;
1297 let _permit = gate.acquire()?;
1298 if let Some(existing) = store.get_sync(&hash)? {
1299 if sha256(&existing) == hash {
1300 return Ok(false);
1301 }
1302 store.delete_sync(&hash)?;
1303 }
1304 let started = Instant::now();
1305 let result = store.put_sync(hash, data);
1306 let success = result.is_ok();
1307 self.record_write(id, started.elapsed(), data.len(), success);
1308 let inserted = result?;
1309 let written = store
1310 .get_sync(&hash)?
1311 .ok_or_else(|| StoreError::Other(format!("pool member {id} lost a committed write")))?;
1312 if sha256(&written) != hash {
1313 self.record_member_failure(id, true);
1314 return Err(StoreError::Other(format!(
1315 "pool member {id} committed corrupt bytes"
1316 )));
1317 }
1318 Ok(inserted)
1319 }
1320
1321 fn delete_member_blob(
1322 &self,
1323 id: PoolMemberId,
1324 store: &LmdbBlobStore,
1325 hash: &Hash,
1326 ) -> Result<bool, StoreError> {
1327 let gate = self.member_gate(id, true)?;
1328 let _permit = gate.acquire()?;
1329 store.delete_sync(hash)
1330 }
1331
1332 fn record_read(&self, id: PoolMemberId, elapsed: Duration, success: bool) {
1333 if let Ok(mut adaptive) = self.adaptive.lock() {
1334 adaptive.record_read(id, elapsed, success);
1335 }
1336 }
1337
1338 fn record_write(&self, id: PoolMemberId, elapsed: Duration, bytes: usize, success: bool) {
1339 if let Ok(mut adaptive) = self.adaptive.lock() {
1340 adaptive.record_write(id, elapsed, bytes, success);
1341 }
1342 }
1343
1344 fn record_member_failure(&self, id: PoolMemberId, write: bool) {
1345 if write {
1346 self.record_write(id, Duration::ZERO, 0, false);
1347 } else {
1348 self.record_read(id, Duration::ZERO, false);
1349 }
1350 }
1351
1352 fn sample_temperature_access(&self, hash: Hash, location: LocationRecord) {
1353 if !self.temperature_config.enabled {
1354 return;
1355 }
1356 let sample_rate = u64::from(self.temperature_config.read_sample_rate.max(1));
1357 let access = self
1358 .temperature_access_counter
1359 .fetch_add(1, Ordering::Relaxed)
1360 .saturating_add(1);
1361 if !access.is_multiple_of(sample_rate) {
1362 return;
1363 }
1364 if let Ok(mut temperature) = self.temperature.lock() {
1365 temperature
1366 .samples
1367 .observe(hash, location, unix_timestamp_now());
1368 }
1369 }
1370}
1371
1372#[async_trait]
1373impl Store for PoolStore {
1374 async fn put(&self, hash: Hash, data: Vec<u8>) -> Result<bool, StoreError> {
1375 self.put_sync(hash, &data)
1376 }
1377
1378 async fn put_many(&self, items: Vec<(Hash, Vec<u8>)>) -> Result<usize, StoreError> {
1379 self.put_many_sync(&items)
1380 }
1381
1382 async fn put_many_optimistic(&self, items: Vec<(Hash, Vec<u8>)>) -> Result<usize, StoreError> {
1383 self.put_many_optimistic_sync(&items)
1384 }
1385
1386 async fn get(&self, hash: &Hash) -> Result<Option<Vec<u8>>, StoreError> {
1387 self.get_sync(hash)
1388 }
1389
1390 async fn get_range(
1391 &self,
1392 hash: &Hash,
1393 start: u64,
1394 end_inclusive: u64,
1395 ) -> Result<Option<Vec<u8>>, StoreError> {
1396 self.get_range_sync(hash, start, end_inclusive)
1397 }
1398
1399 async fn blob_size(&self, hash: &Hash) -> Result<Option<u64>, StoreError> {
1400 self.blob_size_sync(hash)
1401 }
1402
1403 async fn has(&self, hash: &Hash) -> Result<bool, StoreError> {
1404 self.exists(hash)
1405 }
1406
1407 async fn delete(&self, hash: &Hash) -> Result<bool, StoreError> {
1408 self.delete_sync(hash)
1409 }
1410
1411 async fn stats(&self) -> StoreStats {
1412 PoolStore::stats(self).unwrap_or_default()
1413 }
1414
1415 async fn pin(&self, hash: &Hash) -> Result<(), StoreError> {
1416 self.pin_sync(hash)
1417 }
1418
1419 async fn unpin(&self, hash: &Hash) -> Result<(), StoreError> {
1420 self.unpin_sync(hash)
1421 }
1422
1423 fn pin_count(&self, hash: &Hash) -> u32 {
1424 self.pin_count_sync(hash).unwrap_or(0)
1425 }
1426}
1427
1428fn encode_manifest(manifest: &PoolManifest) -> Result<Vec<u8>, StoreError> {
1429 rmp_serde::to_vec_named(manifest)
1430 .map_err(|error| StoreError::Other(format!("encode pool manifest: {error}")))
1431}
1432
1433fn decode_manifest(bytes: &[u8]) -> Result<PoolManifest, StoreError> {
1434 let manifest: PoolManifest = rmp_serde::from_slice(bytes)
1435 .map_err(|error| StoreError::Other(format!("decode pool manifest: {error}")))?;
1436 if manifest.version != 1 {
1437 return Err(StoreError::Other(format!(
1438 "unsupported pool manifest version {}",
1439 manifest.version
1440 )));
1441 }
1442 Ok(manifest)
1443}
1444
1445fn member_hash_key(member: PoolMemberId, hash: Hash) -> [u8; 48] {
1446 let mut key = [0u8; 48];
1447 key[..16].copy_from_slice(member.as_bytes());
1448 key[16..].copy_from_slice(&hash);
1449 key
1450}
1451
1452fn decode_pin_count(bytes: &[u8]) -> Result<u32, StoreError> {
1453 Ok(u32::from_be_bytes(bytes.try_into().map_err(|_| {
1454 StoreError::Other("invalid pool pin count".into())
1455 })?))
1456}
1457
1458fn unix_timestamp_now() -> u64 {
1459 SystemTime::now()
1460 .duration_since(UNIX_EPOCH)
1461 .unwrap_or_default()
1462 .as_secs()
1463}
1464
1465fn put_many_report(
1466 total: usize,
1467 ordered: &[(Hash, u64)],
1468 inserted: &HashSet<Hash>,
1469) -> PutManyReport {
1470 let mut report = PutManyReport {
1471 total,
1472 ..PutManyReport::default()
1473 };
1474 for (hash, bytes) in ordered {
1475 if inserted.contains(hash) {
1476 report.inserted = report.inserted.saturating_add(1);
1477 report.inserted_bytes = report.inserted_bytes.saturating_add(*bytes);
1478 report.inserted_hashes.push(*hash);
1479 }
1480 }
1481 report
1482}
1483
1484fn map_heed(error: heed::Error) -> StoreError {
1485 StoreError::Other(error.to_string())
1486}