Skip to main content

hashtree_lmdb/pool/
mod.rs

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    /// Insert a locally generated content-addressed batch while trusting
519    /// catalogued committed locations. Pending locations still take the
520    /// ordinary repair path; stored and moving locations need no payload read.
521    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    /// Largest map size among members that are available in this process.
727    ///
728    /// The pool catalog is a separate LMDB environment and is intentionally
729    /// not reported as blob capacity.
730    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}