Skip to main content

hashtree_lmdb/pool/
mod.rs

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