Skip to main content

hashtree_lmdb/pool/
mod.rs

1mod adaptive;
2mod catalog;
3mod delete_protection;
4mod gate;
5mod maintenance;
6mod maintenance_batch;
7mod maintenance_move;
8#[cfg(test)]
9mod maintenance_tests;
10mod member;
11mod model;
12mod move_catalog;
13mod read_only;
14mod reader;
15mod temperature;
16mod temperature_balancer;
17mod temperature_catalog;
18mod temperature_worker;
19#[cfg(test)]
20mod tests;
21
22use self::adaptive::AdaptivePoolState;
23pub use self::delete_protection::{PoolDeleteProtectionGuard, POOL_DELETE_PROTECTED};
24use self::gate::ConcurrencyGate;
25use self::member::{open_member_store, prepare_member_paths, validate_member_config};
26use self::model::{LocationRecord, MemberRecord, PoolManifest, MIN_MEMBER_MAP_SIZE_BYTES};
27pub use self::model::{
28    PoolDeleteProtectionChange, PoolDeleteProtectionStatus, PoolMaintenanceReport,
29    PoolMemberConfig, PoolMemberId, PoolMemberRuntimePaths, PoolMemberState, PoolMemberStatus,
30    PoolStalePending, PoolStalePendingCleanupReport, PoolStoreConfig, PoolTemperatureConfig,
31    PoolTemperatureReport,
32};
33use self::move_catalog::{
34    move_cleanup_state_key, move_state_key, rebuild_move_cleanup_member_index_txn,
35    validate_move_cleanup_member_index,
36};
37pub use self::read_only::{
38    ReadOnlyPoolCatalogAudit, ReadOnlyPoolManifestMember, ReadOnlyPoolManifestSnapshot,
39    ReadOnlyPoolStore,
40};
41pub use self::reader::{
42    PoolCatalogLocation, PoolManifestIdentity, PoolPhysicalAudit, PoolReadBatchItem,
43    PoolStoreReader, PoolTerminalAudit,
44};
45use self::temperature::TemperatureRuntime;
46use self::temperature_worker::TemperatureWorker;
47use crate::{managed_env::ManagedEnv, pinned_lmdb_data_len, LmdbBlobStore};
48use async_trait::async_trait;
49use hashtree_core::store::{slice_blob_range, PutManyReport, Store, StoreError, StoreStats};
50use hashtree_core::{sha256, to_hex, types::Hash};
51use heed::types::{Bytes, Unit};
52use heed::{Database, EnvOpenOptions};
53use std::collections::{HashMap, HashSet};
54use std::fs;
55use std::ops::Deref;
56use std::path::{Path, PathBuf};
57use std::sync::atomic::{AtomicU64, Ordering};
58use std::sync::{Arc, Mutex, RwLock};
59use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
60
61const CATALOG_DATABASES: u32 = 6;
62const CATALOG_MAX_READERS: u32 = 1024;
63const MANIFEST_KEY: &[u8] = b"pool-manifest-v1";
64const MEMBER_MARKER_NAME: &str = ".hashtree-pool-member-v1";
65const EXTERNAL_MARKER_NAME: &str = ".hashtree-pool-external-v1";
66const MAX_STALE_PENDING_CLEANUP_ITEMS: usize = 4_096;
67
68fn align_catalog_map_size_bytes(bytes: u64) -> Result<u64, StoreError> {
69    let page_size = (page_size::get() as u64).max(4_096);
70    let remainder = bytes % page_size;
71    if remainder == 0 {
72        Ok(bytes)
73    } else {
74        bytes
75            .checked_add(page_size - remainder)
76            .ok_or_else(|| StoreError::Other("pool catalog map size overflows alignment".into()))
77    }
78}
79
80#[derive(Default)]
81struct RuntimeMembers {
82    generation: Option<u64>,
83    stores: HashMap<PoolMemberId, Arc<LmdbBlobStore>>,
84    read_gates: HashMap<PoolMemberId, Arc<ConcurrencyGate>>,
85    write_gates: HashMap<PoolMemberId, Arc<ConcurrencyGate>>,
86    errors: HashMap<PoolMemberId, String>,
87}
88
89#[derive(Clone)]
90pub struct PoolStore {
91    inner: Arc<PoolStoreInner>,
92}
93
94#[doc(hidden)]
95pub struct PoolStoreInner {
96    env: ManagedEnv,
97    catalog_path: PathBuf,
98    manifest_db: Database<Bytes, Bytes>,
99    locations: Database<Bytes, Bytes>,
100    by_member: Database<Bytes, Unit>,
101    pins: Database<Bytes, Bytes>,
102    last_accessed: Database<Bytes, Bytes>,
103    temperature_state: Database<Bytes, Bytes>,
104    runtime: RwLock<RuntimeMembers>,
105    adaptive: Mutex<AdaptivePoolState>,
106    temperature_config: PoolTemperatureConfig,
107    temperature: Mutex<TemperatureRuntime>,
108    temperature_access_counter: AtomicU64,
109    temperature_owner: PoolMemberId,
110    temperature_cycle: Mutex<()>,
111    temperature_worker: TemperatureWorker,
112    member_runtime_paths: HashMap<PoolMemberId, PoolMemberRuntimePaths>,
113    expected_manifest_sha256: Option<Hash>,
114}
115
116impl Deref for PoolStore {
117    type Target = PoolStoreInner;
118
119    fn deref(&self) -> &Self::Target {
120        &self.inner
121    }
122}
123
124fn validate_new_member_paths(
125    manifest: &PoolManifest,
126    config: &PoolMemberConfig,
127) -> Result<(), StoreError> {
128    if manifest
129        .members
130        .iter()
131        .any(|member| member.config.path == config.path)
132    {
133        return Err(StoreError::Other(format!(
134            "pool member path is already configured: {}",
135            config.path.display()
136        )));
137    }
138    if let Some(external) = config.external_blob_dir.as_ref() {
139        if manifest
140            .members
141            .iter()
142            .any(|member| member.config.external_blob_dir.as_ref() == Some(external))
143        {
144            return Err(StoreError::Other(format!(
145                "pool external blob path is already configured: {}",
146                external.display()
147            )));
148        }
149    }
150    Ok(())
151}
152
153fn open_required_catalog_database<K: 'static, D: 'static>(
154    env: &ManagedEnv,
155    txn: &heed::RoTxn<'_>,
156    name: &str,
157) -> Result<Database<K, D>, StoreError> {
158    env.open_database(txn, Some(name))
159        .map_err(map_heed)?
160        .ok_or_else(|| StoreError::Other(format!("pool {name} database is missing")))
161}
162
163fn validate_controlled_manifest(
164    bytes: &[u8],
165    expected_sha256: Hash,
166    bindings: &HashMap<PoolMemberId, PoolMemberRuntimePaths>,
167) -> Result<PoolManifest, StoreError> {
168    let actual_sha256 = sha256(bytes);
169    if actual_sha256 != expected_sha256 {
170        return Err(StoreError::Other(format!(
171            "live pool manifest SHA-256 differs from controlled authority: expected {}, found {}",
172            to_hex(&expected_sha256),
173            to_hex(&actual_sha256)
174        )));
175    }
176    let manifest = decode_manifest(bytes)?;
177    validate_controlled_manifest_members(&manifest, bindings)?;
178    Ok(manifest)
179}
180
181fn configured_paths_match(pinned: &Path, live: &Path) -> bool {
182    pinned == live
183        || fs::canonicalize(pinned)
184            .ok()
185            .zip(fs::canonicalize(live).ok())
186            .is_some_and(|(pinned, live)| pinned == live)
187}
188
189fn configured_optional_paths_match(pinned: Option<&Path>, live: Option<&Path>) -> bool {
190    match (pinned, live) {
191        (Some(pinned), Some(live)) => configured_paths_match(pinned, live),
192        (None, None) => true,
193        _ => false,
194    }
195}
196
197fn validate_controlled_manifest_members(
198    manifest: &PoolManifest,
199    bindings: &HashMap<PoolMemberId, PoolMemberRuntimePaths>,
200) -> Result<(), StoreError> {
201    if bindings.len() != manifest.members.len() {
202        return Err(StoreError::Other(format!(
203            "pinned pool topology has {} members, live manifest has {}",
204            bindings.len(),
205            manifest.members.len()
206        )));
207    }
208    for member in &manifest.members {
209        let binding = bindings.get(&member.id).ok_or_else(|| {
210            StoreError::Other(format!(
211                "live pool member {} is absent from pinned topology",
212                member.id
213            ))
214        })?;
215        if !configured_paths_match(&binding.configured_path, &member.config.path)
216            || !configured_optional_paths_match(
217                binding.configured_external_path.as_deref(),
218                member.config.external_blob_dir.as_deref(),
219            )
220        {
221            return Err(StoreError::Other(format!(
222                "live pool member {} paths differ from pinned topology",
223                member.id
224            )));
225        }
226        if !member.config.external_blob_sync {
227            return Err(StoreError::Other(format!(
228                "controlled migration requires external_blob_sync=true for pool member {}",
229                member.id
230            )));
231        }
232    }
233    Ok(())
234}
235
236impl PoolStore {
237    pub fn open<P: AsRef<Path>>(path: P, config: PoolStoreConfig) -> Result<Self, StoreError> {
238        config.temperature.validate()?;
239        let mut member_runtime_paths = HashMap::new();
240        for binding in &config.member_runtime_paths {
241            if member_runtime_paths
242                .insert(binding.id, binding.clone())
243                .is_some()
244            {
245                return Err(StoreError::Other(format!(
246                    "duplicate runtime path binding for pool member {}",
247                    binding.id
248                )));
249            }
250            if binding.configured_external_path.is_some() != binding.runtime_external_path.is_some()
251            {
252                return Err(StoreError::Other(format!(
253                    "incomplete external runtime path binding for pool member {}",
254                    binding.id
255                )));
256            }
257        }
258        let controlled = !member_runtime_paths.is_empty()
259            || config.catalog_lmdb_identity.is_some()
260            || config.expected_manifest_sha256.is_some();
261        if controlled
262            && (member_runtime_paths.is_empty()
263                || config.catalog_lmdb_identity.is_none()
264                || config.expected_manifest_sha256.is_none())
265        {
266            return Err(StoreError::Other(
267                "controlled Pool open requires catalog identity, exact manifest SHA-256, and every member runtime binding".into(),
268            ));
269        }
270        let path = path.as_ref();
271        if !controlled {
272            fs::create_dir_all(path).map_err(StoreError::Io)?;
273        }
274        let existing_size = match config.catalog_lmdb_identity {
275            Some(identity) => pinned_lmdb_data_len(path, identity)?,
276            None => fs::metadata(path.join("data.mdb"))
277                .map(|metadata| metadata.len())
278                .unwrap_or(0),
279        };
280        let requested = if controlled {
281            config.catalog_map_size_bytes.max(MIN_MEMBER_MAP_SIZE_BYTES)
282        } else {
283            config
284                .catalog_map_size_bytes
285                .max(existing_size.saturating_add(existing_size / 10))
286                .max(MIN_MEMBER_MAP_SIZE_BYTES)
287        };
288        let requested = align_catalog_map_size_bytes(requested)?;
289        let map_size = usize::try_from(requested)
290            .map_err(|_| StoreError::Other("pool catalog map size exceeds usize".into()))?;
291
292        let mut options = EnvOpenOptions::new();
293        if !controlled {
294            options.map_size(map_size);
295        }
296        options
297            .max_dbs(CATALOG_DATABASES)
298            .max_readers(CATALOG_MAX_READERS);
299        unsafe {
300            options.flags(super::env_flags_from_env());
301        }
302        let env = unsafe {
303            match config.catalog_lmdb_identity {
304                Some(identity) => ManagedEnv::open_pinned(&options, path, identity),
305                None => ManagedEnv::open(&options, path),
306            }
307        }
308        .map_err(|error| {
309            StoreError::Other(format!("open pool catalog {}: {error}", path.display()))
310        })?;
311        if controlled && env.info().map_size < map_size {
312            return Err(StoreError::Other(format!(
313                "controlled Pool catalog map is {} bytes, below its manifest-owned {} bytes; pre-size it before exact migration",
314                env.info().map_size,
315                map_size
316            )));
317        }
318        if !controlled {
319            let _ = env.clear_stale_readers();
320        }
321        let (manifest_db, locations, by_member, pins, last_accessed, temperature_state) =
322            if controlled {
323                let rtxn = env.read_txn().map_err(map_heed)?;
324                let manifest_db: Database<Bytes, Bytes> =
325                    open_required_catalog_database(&env, &rtxn, "manifest")?;
326                let locations: Database<Bytes, Bytes> =
327                    open_required_catalog_database(&env, &rtxn, "locations")?;
328                let by_member: Database<Bytes, Unit> =
329                    open_required_catalog_database(&env, &rtxn, "by_member")?;
330                let pins: Database<Bytes, Bytes> =
331                    open_required_catalog_database(&env, &rtxn, "pins")?;
332                let last_accessed: Database<Bytes, Bytes> =
333                    open_required_catalog_database(&env, &rtxn, "last_accessed")?;
334                let temperature_state: Database<Bytes, Bytes> =
335                    open_required_catalog_database(&env, &rtxn, "temperature_state")?;
336                let manifest_bytes = manifest_db
337                    .get(&rtxn, MANIFEST_KEY)
338                    .map_err(map_heed)?
339                    .ok_or_else(|| StoreError::Other("pool manifest is missing".into()))?;
340                validate_controlled_manifest(
341                    manifest_bytes,
342                    config
343                        .expected_manifest_sha256
344                        .expect("controlled manifest identity checked above"),
345                    &member_runtime_paths,
346                )?;
347                rtxn.commit().map_err(map_heed)?;
348                (
349                    manifest_db,
350                    locations,
351                    by_member,
352                    pins,
353                    last_accessed,
354                    temperature_state,
355                )
356            } else {
357                let mut wtxn = env.write_txn().map_err(map_heed)?;
358                let manifest_db = env
359                    .create_database(&mut wtxn, Some("manifest"))
360                    .map_err(map_heed)?;
361                let locations = env
362                    .create_database(&mut wtxn, Some("locations"))
363                    .map_err(map_heed)?;
364                let by_member = env
365                    .create_database(&mut wtxn, Some("by_member"))
366                    .map_err(map_heed)?;
367                let pins = env
368                    .create_database(&mut wtxn, Some("pins"))
369                    .map_err(map_heed)?;
370                let last_accessed = env
371                    .create_database(&mut wtxn, Some("last_accessed"))
372                    .map_err(map_heed)?;
373                let temperature_state = env
374                    .create_database(&mut wtxn, Some("temperature_state"))
375                    .map_err(map_heed)?;
376                if manifest_db
377                    .get(&wtxn, MANIFEST_KEY)
378                    .map_err(map_heed)?
379                    .is_none()
380                {
381                    let bytes = encode_manifest(&PoolManifest::default())?;
382                    manifest_db
383                        .put(&mut wtxn, MANIFEST_KEY, bytes.as_slice())
384                        .map_err(map_heed)?;
385                }
386                wtxn.commit().map_err(map_heed)?;
387                (
388                    manifest_db,
389                    locations,
390                    by_member,
391                    pins,
392                    last_accessed,
393                    temperature_state,
394                )
395            };
396        if !controlled && env.info().map_size < map_size {
397            unsafe { env.resize(map_size) }.map_err(map_heed)?;
398        }
399
400        if controlled {
401            validate_move_cleanup_member_index(&temperature_state, &by_member, &env)?;
402        } else {
403            let mut wtxn = env.write_txn().map_err(map_heed)?;
404            let cleanup_rebuild_started = Instant::now();
405            let cleanup_rebuild_count =
406                rebuild_move_cleanup_member_index_txn(&temperature_state, &by_member, &mut wtxn)?;
407            let cleanup_rebuild_elapsed = cleanup_rebuild_started.elapsed();
408            wtxn.commit().map_err(map_heed)?;
409            if cleanup_rebuild_count > 0 || cleanup_rebuild_elapsed >= Duration::from_millis(10) {
410                eprintln!(
411                    "Pool cleanup ownership index rebuild: entries {cleanup_rebuild_count}, elapsed {} us",
412                    cleanup_rebuild_elapsed.as_micros()
413                );
414            }
415        }
416
417        let temperature_config = config.temperature.clone();
418        let store = Self {
419            inner: Arc::new(PoolStoreInner {
420                env,
421                catalog_path: path.to_path_buf(),
422                manifest_db,
423                locations,
424                by_member,
425                pins,
426                last_accessed,
427                temperature_state,
428                runtime: RwLock::new(RuntimeMembers::default()),
429                adaptive: Mutex::new(AdaptivePoolState::new(config.member_failure_cooldown)),
430                temperature: Mutex::new(TemperatureRuntime::new(
431                    temperature_config.candidate_capacity,
432                )),
433                temperature_config,
434                temperature_access_counter: AtomicU64::new(0),
435                temperature_owner: PoolMemberId::new(),
436                temperature_cycle: Mutex::new(()),
437                temperature_worker: TemperatureWorker::default(),
438                member_runtime_paths,
439                expected_manifest_sha256: config.expected_manifest_sha256,
440            }),
441        };
442        store.refresh_members()?;
443        store.start_temperature_worker()?;
444        Ok(store)
445    }
446
447    fn start_temperature_worker(&self) -> Result<(), StoreError> {
448        if !self.temperature_config.enabled {
449            return Ok(());
450        }
451        self.temperature_worker.start(
452            Arc::downgrade(&self.inner),
453            self.temperature_config.interval,
454        )
455    }
456
457    /// Stop this process's background temperature balancer.
458    ///
459    /// Exact maintenance and recovery commands use this before taking
460    /// catalog snapshots so the same Pool handle cannot relocate blobs behind
461    /// their authority checks.
462    pub fn stop_temperature_worker(&self) -> Result<(), StoreError> {
463        self.temperature_worker.stop()
464    }
465
466    pub fn add_member(&self, config: PoolMemberConfig) -> Result<PoolMemberId, StoreError> {
467        self.add_member_inner(config, false)?
468            .ok_or_else(|| StoreError::Other("pool member was not added".into()))
469    }
470
471    pub(crate) fn ensure_initial_member(
472        &self,
473        config: PoolMemberConfig,
474    ) -> Result<Option<PoolMemberId>, StoreError> {
475        if !self.read_manifest()?.members.is_empty() {
476            self.refresh_members()?;
477            return Ok(None);
478        }
479        self.add_member_inner(config, true)
480    }
481
482    fn add_member_inner(
483        &self,
484        config: PoolMemberConfig,
485        only_if_empty: bool,
486    ) -> Result<Option<PoolMemberId>, StoreError> {
487        if !self.member_runtime_paths.is_empty() {
488            return Err(StoreError::Other(
489                "cannot change pool membership while exact runtime paths are pinned".into(),
490            ));
491        }
492        validate_member_config(&config)?;
493        let wtxn = self.env.write_txn().map_err(map_heed)?;
494        let manifest = self.manifest_from_txn(&wtxn)?;
495        if only_if_empty && !manifest.members.is_empty() {
496            drop(wtxn);
497            self.refresh_members()?;
498            return Ok(None);
499        }
500        validate_new_member_paths(&manifest, &config)?;
501        let id = prepare_member_paths(&config, PoolMemberId::new())?;
502        drop(wtxn);
503        let store = Arc::new(open_member_store(id, &config, None)?);
504
505        let mut wtxn = self.env.write_txn().map_err(map_heed)?;
506        let mut manifest = self.manifest_from_txn(&wtxn)?;
507        if only_if_empty && !manifest.members.is_empty() {
508            drop(wtxn);
509            self.refresh_members()?;
510            return Ok(None);
511        }
512        validate_new_member_paths(&manifest, &config)?;
513        if manifest.members.iter().any(|member| member.id == id) {
514            return Err(StoreError::Other(format!(
515                "pool member identity is already configured: {id}"
516            )));
517        }
518        manifest.members.push(MemberRecord {
519            id,
520            state: PoolMemberState::Active,
521            config,
522        });
523        manifest.generation = manifest.generation.saturating_add(1);
524        self.put_manifest_txn(&mut wtxn, &manifest)?;
525        wtxn.commit().map_err(map_heed)?;
526
527        let mut runtime = self
528            .runtime
529            .write()
530            .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?;
531        runtime.generation = Some(manifest.generation);
532        runtime.stores.insert(id, store);
533        let member = manifest
534            .members
535            .iter()
536            .find(|member| member.id == id)
537            .expect("new pool member is in committed manifest");
538        runtime.read_gates.insert(
539            id,
540            Arc::new(ConcurrencyGate::new(member.config.max_read_concurrency)),
541        );
542        runtime.write_gates.insert(
543            id,
544            Arc::new(ConcurrencyGate::new(member.config.max_write_concurrency)),
545        );
546        runtime.errors.remove(&id);
547        Ok(Some(id))
548    }
549
550    pub fn begin_drain(&self, id: PoolMemberId) -> Result<(), StoreError> {
551        let mut wtxn = self.env.write_txn().map_err(map_heed)?;
552        let mut manifest = self.manifest_from_txn(&wtxn)?;
553        let has_other_active = manifest
554            .members
555            .iter()
556            .any(|member| member.id != id && member.state == PoolMemberState::Active);
557        let member = manifest
558            .members
559            .iter_mut()
560            .find(|member| member.id == id)
561            .ok_or_else(|| StoreError::Other(format!("unknown pool member {id}")))?;
562        if member.state == PoolMemberState::Draining {
563            return Ok(());
564        }
565        if !has_other_active && self.member_has_locations_txn(&wtxn, id)? {
566            return Err(StoreError::Other(
567                "cannot drain the final member while it still owns blobs".into(),
568            ));
569        }
570        member.state = PoolMemberState::Draining;
571        manifest.generation = manifest.generation.saturating_add(1);
572        self.put_manifest_txn(&mut wtxn, &manifest)?;
573        wtxn.commit().map_err(map_heed)?;
574        self.refresh_members()?;
575        Ok(())
576    }
577
578    pub fn update_member_limits(
579        &self,
580        id: PoolMemberId,
581        capacity_bytes: u64,
582        max_read_concurrency: u32,
583        max_write_concurrency: u32,
584    ) -> Result<(), StoreError> {
585        if capacity_bytes == 0 || max_read_concurrency == 0 || max_write_concurrency == 0 {
586            return Err(StoreError::Other(
587                "pool member capacity and concurrency limits must be non-zero".into(),
588            ));
589        }
590        let mut wtxn = self.env.write_txn().map_err(map_heed)?;
591        let mut manifest = self.manifest_from_txn(&wtxn)?;
592        let member = manifest
593            .members
594            .iter_mut()
595            .find(|member| member.id == id)
596            .ok_or_else(|| StoreError::Other(format!("unknown pool member {id}")))?;
597        member.config.capacity_bytes = capacity_bytes;
598        member.config.max_read_concurrency = max_read_concurrency;
599        member.config.max_write_concurrency = max_write_concurrency;
600        manifest.generation = manifest.generation.saturating_add(1);
601        self.put_manifest_txn(&mut wtxn, &manifest)?;
602        wtxn.commit().map_err(map_heed)?;
603        self.refresh_members()
604    }
605
606    pub fn update_member_temperature_watermarks(
607        &self,
608        id: PoolMemberId,
609        low_percent: u8,
610        high_percent: u8,
611    ) -> Result<(), StoreError> {
612        if low_percent >= high_percent || high_percent > 100 {
613            return Err(StoreError::Other(
614                "pool temperature watermarks must satisfy 0 <= low < high <= 100".into(),
615            ));
616        }
617        let mut wtxn = self.env.write_txn().map_err(map_heed)?;
618        let mut manifest = self.manifest_from_txn(&wtxn)?;
619        let member = manifest
620            .members
621            .iter_mut()
622            .find(|member| member.id == id)
623            .ok_or_else(|| StoreError::Other(format!("unknown pool member {id}")))?;
624        member.config.temperature_low_watermark_percent = low_percent;
625        member.config.temperature_high_watermark_percent = high_percent;
626        manifest.generation = manifest.generation.saturating_add(1);
627        self.put_manifest_txn(&mut wtxn, &manifest)?;
628        wtxn.commit().map_err(map_heed)?;
629        self.refresh_members()
630    }
631
632    pub fn remove_member(&self, id: PoolMemberId) -> Result<(), StoreError> {
633        let mut wtxn = self.env.write_txn().map_err(map_heed)?;
634        let mut manifest = self.manifest_from_txn(&wtxn)?;
635        let index = manifest
636            .members
637            .iter()
638            .position(|member| member.id == id)
639            .ok_or_else(|| StoreError::Other(format!("unknown pool member {id}")))?;
640        if manifest.members[index].state != PoolMemberState::Draining {
641            return Err(StoreError::Other(
642                "pool member must be draining before removal".into(),
643            ));
644        }
645        if self.member_has_locations_txn(&wtxn, id)? {
646            return Err(StoreError::Other(format!(
647                "pool member {id} still owns blob(s)"
648            )));
649        }
650        manifest.members.remove(index);
651        manifest.generation = manifest.generation.saturating_add(1);
652        self.put_manifest_txn(&mut wtxn, &manifest)?;
653        wtxn.commit().map_err(map_heed)?;
654        self.refresh_members()?;
655        Ok(())
656    }
657
658    pub fn member(&self, id: PoolMemberId) -> Result<PoolMemberStatus, StoreError> {
659        self.refresh_members()?;
660        let manifest = self.read_manifest()?;
661        let member = manifest
662            .members
663            .iter()
664            .find(|member| member.id == id)
665            .ok_or_else(|| StoreError::Other(format!("unknown pool member {id}")))?;
666        let located_blobs = self.count_member_locations(id)?;
667        let runtime = self
668            .runtime
669            .read()
670            .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?;
671        let (logical_bytes, available, last_error) = match runtime.stores.get(&id) {
672            Some(store) => match store.stats() {
673                Ok(stats) => (stats.total_bytes, true, None),
674                Err(error) => (0, false, Some(error.to_string())),
675            },
676            None => (0, false, runtime.errors.get(&id).cloned()),
677        };
678        Ok(PoolMemberStatus {
679            id,
680            state: member.state,
681            path: member.config.path.clone(),
682            capacity_bytes: member.config.capacity_bytes,
683            map_size_bytes: member.config.map_size_bytes,
684            external_blob_dir: member.config.external_blob_dir.clone(),
685            external_blob_min_bytes: member.config.external_blob_min_bytes,
686            external_blob_sync: member.config.external_blob_sync,
687            external_pack_target_bytes: member.config.external_pack_target_bytes,
688            max_read_concurrency: member.config.max_read_concurrency,
689            max_write_concurrency: member.config.max_write_concurrency,
690            temperature_low_watermark_percent: member.config.temperature_low_watermark_percent,
691            temperature_high_watermark_percent: member.config.temperature_high_watermark_percent,
692            logical_bytes,
693            located_blobs,
694            available,
695            last_error,
696        })
697    }
698
699    pub fn members(&self) -> Result<Vec<PoolMemberStatus>, StoreError> {
700        let manifest = self.read_manifest()?;
701        manifest
702            .members
703            .iter()
704            .map(|member| self.member(member.id))
705            .collect()
706    }
707
708    pub fn blob_location(&self, hash: &Hash) -> Result<Option<PoolMemberId>, StoreError> {
709        Ok(self
710            .read_location(hash)?
711            .map(LocationRecord::preferred_member))
712    }
713
714    pub fn put_sync(&self, hash: Hash, data: &[u8]) -> Result<bool, StoreError> {
715        if sha256(data) != hash {
716            return Err(StoreError::Other(
717                "pool rejected bytes that do not match their hash".into(),
718            ));
719        }
720
721        if let Some(location) = self.read_location(&hash)? {
722            match self.read_verified_location(&hash, location) {
723                Ok(Some(found)) => {
724                    if matches!(location, LocationRecord::Pending { .. }) {
725                        self.finalize_pending(hash, location)?;
726                    }
727                    debug_assert_eq!(sha256(&found), hash);
728                    return Ok(false);
729                }
730                Ok(None) | Err(_) => return self.repair_location(hash, data, location),
731            }
732        }
733
734        let target = self.choose_write_member(data.len() as u64, None)?;
735        let pending = LocationRecord::Pending {
736            member: target,
737            size: data.len() as u64,
738        };
739        let location = self.reserve_if_absent(hash, pending)?;
740        if location != pending {
741            match self.read_verified_location(&hash, location) {
742                Ok(Some(found)) => {
743                    if matches!(location, LocationRecord::Pending { .. }) {
744                        self.finalize_pending(hash, location)?;
745                    }
746                    debug_assert_eq!(sha256(&found), hash);
747                    return Ok(false);
748                }
749                Ok(None) | Err(_) => return self.repair_location(hash, data, location),
750            }
751        }
752
753        let target = location.preferred_member();
754        let store = self.get_member(target)?;
755        match self.write_verified_member(target, &store, hash, data) {
756            Ok(inserted) => {
757                self.finalize_pending(hash, location)?;
758                Ok(inserted)
759            }
760            Err(_) => {
761                let mut excluded = HashSet::new();
762                excluded.insert(target);
763                self.repair_location_excluding(hash, data, location, excluded)
764            }
765        }
766    }
767
768    pub fn put_many_report_sync(
769        &self,
770        items: &[(Hash, Vec<u8>)],
771    ) -> Result<PutManyReport, StoreError> {
772        self.put_many_report_sync_with_existing_verification(items, true)
773    }
774
775    /// Insert a locally generated content-addressed batch while trusting
776    /// catalogued committed locations. Pending locations still take the
777    /// ordinary repair path; stored and moving locations need no payload read.
778    pub fn put_many_optimistic_report_sync(
779        &self,
780        items: &[(Hash, Vec<u8>)],
781    ) -> Result<PutManyReport, StoreError> {
782        self.put_many_report_sync_with_existing_verification(items, false)
783    }
784
785    fn put_many_report_sync_with_existing_verification(
786        &self,
787        items: &[(Hash, Vec<u8>)],
788        verify_existing: bool,
789    ) -> Result<PutManyReport, StoreError> {
790        let mut seen = HashSet::new();
791        let mut unique = Vec::with_capacity(items.len());
792        let mut ordered = Vec::with_capacity(items.len());
793        for (hash, data) in items {
794            if sha256(data) != *hash {
795                return Err(StoreError::Other(
796                    "pool rejected batch bytes that do not match their hash".into(),
797                ));
798            }
799            if seen.insert(*hash) {
800                unique.push((*hash, data));
801                ordered.push((*hash, data.len() as u64));
802            }
803        }
804
805        let mut inserted = HashSet::new();
806        let mut missing = Vec::new();
807        for (hash, data) in unique {
808            if let Some(location) = self.read_location(&hash)? {
809                if (verify_existing || matches!(location, LocationRecord::Pending { .. }))
810                    && self.put_sync(hash, data)?
811                {
812                    inserted.insert(hash);
813                }
814            } else {
815                missing.push((hash, data));
816            }
817        }
818        if missing.is_empty() {
819            return Ok(put_many_report(items.len(), &ordered, &inserted));
820        }
821
822        let mut reserved_bytes = HashMap::new();
823        let mut assignments = Vec::with_capacity(missing.len());
824        for (hash, data) in missing {
825            let target =
826                self.choose_write_member_with_reserved(data.len() as u64, None, &reserved_bytes)?;
827            let reserved = reserved_bytes.entry(target).or_insert(0u64);
828            *reserved = reserved.saturating_add(data.len() as u64);
829            assignments.push((hash, data, target));
830        }
831
832        let mut wtxn = self.env.write_txn().map_err(map_heed)?;
833        let mut plans = Vec::with_capacity(assignments.len());
834        let mut raced = Vec::new();
835        for (hash, data, target) in assignments {
836            if self
837                .locations
838                .get(&wtxn, &hash)
839                .map_err(map_heed)?
840                .is_some()
841            {
842                raced.push((hash, data));
843                continue;
844            }
845            let pending = LocationRecord::Pending {
846                member: target,
847                size: data.len() as u64,
848            };
849            self.set_location_txn(&mut wtxn, hash, Some(pending))?;
850            plans.push((hash, data, target, pending));
851        }
852        wtxn.commit().map_err(map_heed)?;
853
854        for (hash, data) in raced {
855            if self.put_sync(hash, data)? {
856                inserted.insert(hash);
857            }
858        }
859
860        let mut by_target: HashMap<PoolMemberId, Vec<(Hash, &[u8])>> = HashMap::new();
861        for (hash, data, target, _) in &plans {
862            by_target
863                .entry(*target)
864                .or_default()
865                .push((*hash, data.as_slice()));
866        }
867        for (target, batch) in by_target {
868            let store = self.get_member(target)?;
869            let gate = self.member_gate(target, true)?;
870            let permit = gate.acquire()?;
871            for (hash, _) in &batch {
872                if store
873                    .get_sync(hash)?
874                    .is_some_and(|existing| sha256(&existing) != *hash)
875                {
876                    store.delete_sync(hash)?;
877                }
878            }
879            let started = Instant::now();
880            let result = store.put_many_refs_report_sync(&batch);
881            let success = result.is_ok();
882            let bytes = batch.iter().map(|(_, data)| data.len()).sum::<usize>();
883            self.record_write(target, started.elapsed(), bytes, success);
884            let report = match result {
885                Ok(report) => report,
886                Err(_) => {
887                    drop(permit);
888                    for (hash, data) in batch {
889                        if self.put_sync(hash, data)? {
890                            inserted.insert(hash);
891                        }
892                    }
893                    continue;
894                }
895            };
896            inserted.extend(report.inserted_hashes);
897            for (hash, _) in &batch {
898                self.read_verified_member(target, &store, hash)?
899                    .ok_or_else(|| {
900                        StoreError::Other(format!(
901                            "pool member {target} lost a committed batch write"
902                        ))
903                    })?;
904            }
905        }
906
907        let mut wtxn = self.env.write_txn().map_err(map_heed)?;
908        for (hash, _, target, pending) in plans {
909            let Some(current) = self.locations.get(&wtxn, &hash).map_err(map_heed)? else {
910                continue;
911            };
912            if LocationRecord::decode(current)? == pending {
913                self.set_location_txn(
914                    &mut wtxn,
915                    hash,
916                    Some(LocationRecord::Stored {
917                        member: target,
918                        size: pending.size(),
919                    }),
920                )?;
921            }
922        }
923        wtxn.commit().map_err(map_heed)?;
924        Ok(put_many_report(items.len(), &ordered, &inserted))
925    }
926
927    pub fn put_many_sync(&self, items: &[(Hash, Vec<u8>)]) -> Result<usize, StoreError> {
928        self.put_many_report_sync(items)
929            .map(|report| report.inserted)
930    }
931
932    pub fn put_many_optimistic_sync(&self, items: &[(Hash, Vec<u8>)]) -> Result<usize, StoreError> {
933        self.put_many_optimistic_report_sync(items)
934            .map(|report| report.inserted)
935    }
936
937    pub fn get_sync(&self, hash: &Hash) -> Result<Option<Vec<u8>>, StoreError> {
938        let Some(location) = self.read_location(hash)? else {
939            return Ok(None);
940        };
941        let data = self.read_verified_location(hash, location)?;
942        if data.is_some() && matches!(location, LocationRecord::Pending { .. }) {
943            self.finalize_pending(*hash, location)?;
944        }
945        if data.is_some() {
946            self.sample_temperature_access(*hash, location);
947        }
948        Ok(data)
949    }
950
951    pub fn get_range_sync(
952        &self,
953        hash: &Hash,
954        start: u64,
955        end_inclusive: u64,
956    ) -> Result<Option<Vec<u8>>, StoreError> {
957        let Some(data) = self.get_sync(hash)? else {
958            return Ok(None);
959        };
960        slice_blob_range(&data, start, end_inclusive).map(Some)
961    }
962
963    pub fn blob_size_sync(&self, hash: &Hash) -> Result<Option<u64>, StoreError> {
964        Ok(self.read_location(hash)?.map(LocationRecord::size))
965    }
966
967    pub fn exists(&self, hash: &Hash) -> Result<bool, StoreError> {
968        Ok(self.get_sync(hash)?.is_some())
969    }
970
971    pub fn existing_hashes_in_sorted_candidates(
972        &self,
973        sorted_hashes: &[Hash],
974    ) -> Result<Vec<bool>, StoreError> {
975        let rtxn = self.env.read_txn().map_err(map_heed)?;
976        sorted_hashes
977            .iter()
978            .map(|hash| self.locations.get(&rtxn, hash).map(|value| value.is_some()))
979            .collect::<Result<Vec<_>, _>>()
980            .map_err(map_heed)
981    }
982
983    /// Mark hashes whose pool catalog entries represent committed data.
984    ///
985    /// Pending records deliberately return false so an interrupted migration
986    /// supplies the source bytes again and lets the ordinary repair/finalize
987    /// path complete them. Stored and moving records are already durable and
988    /// can be skipped by an incremental migration without rereading payloads.
989    pub fn committed_hashes_in_sorted_candidates(
990        &self,
991        sorted_hashes: &[Hash],
992    ) -> Result<Vec<bool>, StoreError> {
993        let rtxn = self.env.read_txn().map_err(map_heed)?;
994        sorted_hashes
995            .iter()
996            .map(|hash| {
997                self.locations
998                    .get(&rtxn, hash)
999                    .map_err(map_heed)?
1000                    .map(LocationRecord::decode)
1001                    .transpose()
1002                    .map(|location| {
1003                        location.is_some_and(|location| {
1004                            !matches!(location, LocationRecord::Pending { .. })
1005                        })
1006                    })
1007            })
1008            .collect()
1009    }
1010
1011    /// Return the exact catalog state for one sorted source page using a
1012    /// single read transaction.
1013    ///
1014    /// Migration reconciliation uses this to distinguish a terminal
1015    /// size-matched `Stored` record from `Missing`, crash-left `Pending`, and
1016    /// non-terminal `Moving` without loading target payload bytes.
1017    pub fn catalog_locations_in_sorted_candidates(
1018        &self,
1019        sorted_hashes: &[Hash],
1020    ) -> Result<Vec<PoolCatalogLocation>, StoreError> {
1021        let rtxn = self.env.read_txn().map_err(map_heed)?;
1022        sorted_hashes
1023            .iter()
1024            .map(|hash| {
1025                let record = self
1026                    .locations
1027                    .get(&rtxn, hash)
1028                    .map_err(map_heed)?
1029                    .map(LocationRecord::decode)
1030                    .transpose()?;
1031                Ok(PoolCatalogLocation::from_record(record))
1032            })
1033            .collect()
1034    }
1035
1036    /// Scan one bounded raw catalog page in hash order.
1037    ///
1038    /// Online migration content audits use the raw catalog key as their
1039    /// resumable cursor. Pending and moving rows are returned explicitly:
1040    /// they are not terminal content authorities, but their keys must still
1041    /// advance the scan so a page containing only transient rows makes
1042    /// progress.
1043    pub fn scan_catalog_locations_after(
1044        &self,
1045        after: Option<Hash>,
1046        limit: usize,
1047    ) -> Result<Vec<(Hash, PoolCatalogLocation)>, StoreError> {
1048        if limit == 0 {
1049            return Ok(Vec::new());
1050        }
1051        let rtxn = self.env.read_txn().map_err(map_heed)?;
1052        let mut entries = Vec::with_capacity(limit);
1053        let decode =
1054            |hash: &[u8], encoded: &[u8]| -> Result<(Hash, PoolCatalogLocation), StoreError> {
1055                let hash: Hash = hash
1056                    .try_into()
1057                    .map_err(|_| StoreError::Other("invalid pool catalog hash length".into()))?;
1058                let record = LocationRecord::decode(encoded)?;
1059                Ok((hash, PoolCatalogLocation::from_record(Some(record))))
1060            };
1061        match after {
1062            Some(after) => {
1063                use std::ops::Bound;
1064                let range = (Bound::Excluded(after.as_slice()), Bound::<&[u8]>::Unbounded);
1065                for item in self.locations.range(&rtxn, &range).map_err(map_heed)? {
1066                    let (hash, encoded) = item.map_err(map_heed)?;
1067                    entries.push(decode(hash, encoded)?);
1068                    if entries.len() >= limit {
1069                        break;
1070                    }
1071                }
1072            }
1073            None => {
1074                for item in self.locations.iter(&rtxn).map_err(map_heed)? {
1075                    let (hash, encoded) = item.map_err(map_heed)?;
1076                    entries.push(decode(hash, encoded)?);
1077                    if entries.len() >= limit {
1078                        break;
1079                    }
1080                }
1081            }
1082        }
1083        Ok(entries)
1084    }
1085
1086    /// Largest map size among members that are available in this process.
1087    ///
1088    /// The pool catalog is a separate LMDB environment and is intentionally
1089    /// not reported as blob capacity.
1090    pub fn largest_member_map_size_bytes(&self) -> Result<Option<usize>, StoreError> {
1091        self.refresh_members()?;
1092        let runtime = self
1093            .runtime
1094            .read()
1095            .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?;
1096        Ok(runtime
1097            .stores
1098            .values()
1099            .map(|store| store.map_size_bytes())
1100            .max())
1101    }
1102
1103    pub fn delete_sync(&self, hash: &Hash) -> Result<bool, StoreError> {
1104        let _coordination = self.acquire_delete_coordination_lock(false)?;
1105        self.require_deletes_unprotected()?;
1106        let Some(location) = self.read_location(hash)? else {
1107            return Ok(false);
1108        };
1109        let (members, len) = location.members();
1110        let mut deleted = false;
1111        for member in members.into_iter().take(len) {
1112            if let Ok(store) = self.get_member(member) {
1113                let gate = self.member_gate(member, true)?;
1114                let _permit = gate.acquire()?;
1115                deleted |= store.delete_sync(hash)?;
1116            }
1117        }
1118        let mut wtxn = self.env.write_txn().map_err(map_heed)?;
1119        self.set_location_txn(&mut wtxn, *hash, None)?;
1120        self.pins.delete(&mut wtxn, hash).map_err(map_heed)?;
1121        wtxn.commit().map_err(map_heed)?;
1122        Ok(deleted)
1123    }
1124
1125    /// Delete a batch with one transaction per affected member and one pool
1126    /// catalog transaction.
1127    pub fn delete_many_sync(&self, hashes: &[Hash]) -> Result<usize, StoreError> {
1128        if hashes.is_empty() {
1129            return Ok(0);
1130        }
1131        let _coordination = self.acquire_delete_coordination_lock(false)?;
1132        self.require_deletes_unprotected()?;
1133        let rtxn = self.env.read_txn().map_err(map_heed)?;
1134        let mut seen = HashSet::with_capacity(hashes.len());
1135        let mut located = Vec::new();
1136        let mut by_member = HashMap::<PoolMemberId, Vec<Hash>>::new();
1137        for hash in hashes {
1138            if !seen.insert(*hash) {
1139                continue;
1140            }
1141            let Some(encoded) = self.locations.get(&rtxn, hash).map_err(map_heed)? else {
1142                continue;
1143            };
1144            let location = LocationRecord::decode(encoded)?;
1145            let (members, len) = location.members();
1146            for member in members.into_iter().take(len) {
1147                by_member.entry(member).or_default().push(*hash);
1148            }
1149            located.push(*hash);
1150        }
1151        drop(rtxn);
1152
1153        for (member, member_hashes) in by_member {
1154            if let Ok(store) = self.get_member(member) {
1155                let gate = self.member_gate(member, true)?;
1156                let _permit = gate.acquire()?;
1157                store.delete_many_sync(&member_hashes)?;
1158            }
1159        }
1160
1161        let mut wtxn = self.env.write_txn().map_err(map_heed)?;
1162        for hash in &located {
1163            self.set_location_txn(&mut wtxn, *hash, None)?;
1164            self.pins.delete(&mut wtxn, hash).map_err(map_heed)?;
1165        }
1166        wtxn.commit().map_err(map_heed)?;
1167        Ok(located.len())
1168    }
1169
1170    /// Clear one exact, bounded set of crash-abandoned `Pending` records.
1171    ///
1172    /// This is deliberately narrower than [`Self::delete_many_sync`]: every
1173    /// live location must still equal the authorized `(hash, member, size)`,
1174    /// no hash may be pinned or owned by a move, and no member may contain a
1175    /// physical record for the hash. The catalog and all secondary indexes
1176    /// change in one transaction and are force-synced before success.
1177    ///
1178    /// A caller must stop and fence **every** Pool writer before its strict
1179    /// read-only audit and keep that fence through this call. This includes
1180    /// clones and handles in this process: a writer can commit `Pending`, pause
1181    /// before entering its member gate, and otherwise resume after cleanup.
1182    /// Catalog and member LMDB environments do not provide one atomic
1183    /// cross-environment transaction, so this method intentionally makes no
1184    /// online-safety claim.
1185    pub fn cleanup_stale_pending_exact_offline_sync(
1186        &self,
1187        expected: &[PoolStalePending],
1188    ) -> Result<PoolStalePendingCleanupReport, StoreError> {
1189        if expected.is_empty() || expected.len() > MAX_STALE_PENDING_CLEANUP_ITEMS {
1190            return Err(StoreError::Other(format!(
1191                "stale Pending cleanup requires 1..={MAX_STALE_PENDING_CLEANUP_ITEMS} exact records"
1192            )));
1193        }
1194        if self.expected_manifest_sha256.is_none() || self.member_runtime_paths.is_empty() {
1195            return Err(StoreError::Other(
1196                "stale Pending cleanup requires a controlled Pool open".into(),
1197            ));
1198        }
1199        if self.temperature_config.enabled {
1200            return Err(StoreError::Other(
1201                "stale Pending cleanup requires temperature tracking to be disabled".into(),
1202            ));
1203        }
1204        if expected.windows(2).any(|pair| pair[0].hash >= pair[1].hash) {
1205            return Err(StoreError::Other(
1206                "stale Pending cleanup records must be strictly hash-sorted and unique".into(),
1207            ));
1208        }
1209        let declared_bytes = expected.iter().try_fold(0u64, |total, item| {
1210            total.checked_add(item.size).ok_or_else(|| {
1211                StoreError::Other("stale Pending declared byte total overflow".into())
1212            })
1213        })?;
1214
1215        self.validate_controlled_authority_and_sync()?;
1216
1217        let mut member_ids = expected.iter().map(|item| item.member).collect::<Vec<_>>();
1218        member_ids.sort_unstable();
1219        member_ids.dedup();
1220        let member_resources = member_ids
1221            .iter()
1222            .map(|member| {
1223                Ok((
1224                    *member,
1225                    self.get_member(*member)?,
1226                    self.member_gate(*member, true)?,
1227                ))
1228            })
1229            .collect::<Result<Vec<_>, StoreError>>()?;
1230        let _exclusive_permits = member_resources
1231            .iter()
1232            .map(|(_, _, gate)| gate.acquire_exclusive())
1233            .collect::<Result<Vec<_>, StoreError>>()?;
1234
1235        let ensure_physically_absent = || -> Result<(), StoreError> {
1236            for (member, store, _) in &member_resources {
1237                let hashes = expected
1238                    .iter()
1239                    .filter(|item| item.member == *member)
1240                    .map(|item| item.hash)
1241                    .collect::<Vec<_>>();
1242                let present = store.existing_hashes_in_sorted_candidates(&hashes)?;
1243                if let Some(hash) = hashes
1244                    .iter()
1245                    .zip(present)
1246                    .find_map(|(hash, present)| present.then_some(hash))
1247                {
1248                    return Err(StoreError::Other(format!(
1249                        "stale Pending cleanup rejected physically present member record {} on {member}",
1250                        to_hex(hash)
1251                    )));
1252                }
1253            }
1254            Ok(())
1255        };
1256        ensure_physically_absent()?;
1257
1258        let mut wtxn = self.env.write_txn().map_err(map_heed)?;
1259        let mut exact_pending = 0usize;
1260        let mut already_missing = 0usize;
1261        for item in expected {
1262            let current = self
1263                .locations
1264                .get(&wtxn, &item.hash)
1265                .map_err(map_heed)?
1266                .map(LocationRecord::decode)
1267                .transpose()?;
1268            let authorized = LocationRecord::Pending {
1269                member: item.member,
1270                size: item.size,
1271            };
1272            match current {
1273                Some(current) if current == authorized => {
1274                    exact_pending = exact_pending.saturating_add(1);
1275                    if self
1276                        .by_member
1277                        .get(&wtxn, &member_hash_key(item.member, item.hash))
1278                        .map_err(map_heed)?
1279                        .is_none()
1280                    {
1281                        return Err(StoreError::Other(format!(
1282                            "stale Pending cleanup found a missing member index for {}",
1283                            to_hex(&item.hash)
1284                        )));
1285                    }
1286                }
1287                None => {
1288                    already_missing = already_missing.saturating_add(1);
1289                    if self
1290                        .by_member
1291                        .get(&wtxn, &member_hash_key(item.member, item.hash))
1292                        .map_err(map_heed)?
1293                        .is_some()
1294                        || self
1295                            .last_accessed
1296                            .get(&wtxn, &item.hash)
1297                            .map_err(map_heed)?
1298                            .is_some()
1299                    {
1300                        return Err(StoreError::Other(format!(
1301                            "already-cleared stale Pending record {} retains catalog indexes",
1302                            to_hex(&item.hash)
1303                        )));
1304                    }
1305                }
1306                current => {
1307                    return Err(StoreError::Other(format!(
1308                        "stale Pending authority no longer matches {}: expected {authorized:?}, found {current:?}",
1309                        to_hex(&item.hash)
1310                    )));
1311                }
1312            }
1313            let pin_count = self
1314                .pins
1315                .get(&wtxn, &item.hash)
1316                .map_err(map_heed)?
1317                .map(decode_pin_count)
1318                .transpose()?
1319                .unwrap_or(0);
1320            if pin_count != 0 {
1321                return Err(StoreError::Other(format!(
1322                    "stale Pending cleanup rejected pinned hash {}",
1323                    to_hex(&item.hash)
1324                )));
1325            }
1326            if self
1327                .temperature_state
1328                .get(&wtxn, &move_state_key(item.hash))
1329                .map_err(map_heed)?
1330                .is_some()
1331                || self
1332                    .temperature_state
1333                    .get(&wtxn, &move_cleanup_state_key(item.hash))
1334                    .map_err(map_heed)?
1335                    .is_some()
1336            {
1337                return Err(StoreError::Other(format!(
1338                    "stale Pending cleanup rejected move-owned hash {}",
1339                    to_hex(&item.hash)
1340                )));
1341            }
1342        }
1343        if exact_pending != 0 && already_missing != 0 {
1344            return Err(StoreError::Other(
1345                "stale Pending cleanup refuses a partially changed authority set".into(),
1346            ));
1347        }
1348
1349        // Repeat the cross-environment absence proof as late as possible. The
1350        // operational writer fence remains mandatory for other Pool handles.
1351        ensure_physically_absent()?;
1352        if exact_pending != 0 {
1353            for item in expected {
1354                self.set_location_txn(&mut wtxn, item.hash, None)?;
1355                self.pins.delete(&mut wtxn, &item.hash).map_err(map_heed)?;
1356            }
1357        }
1358        wtxn.commit().map_err(map_heed)?;
1359        self.validate_controlled_authority_and_sync()?;
1360
1361        let rtxn = self.env.read_txn().map_err(map_heed)?;
1362        for item in expected {
1363            if self
1364                .locations
1365                .get(&rtxn, &item.hash)
1366                .map_err(map_heed)?
1367                .is_some()
1368                || self
1369                    .by_member
1370                    .get(&rtxn, &member_hash_key(item.member, item.hash))
1371                    .map_err(map_heed)?
1372                    .is_some()
1373                || self
1374                    .pins
1375                    .get(&rtxn, &item.hash)
1376                    .map_err(map_heed)?
1377                    .is_some()
1378                || self
1379                    .last_accessed
1380                    .get(&rtxn, &item.hash)
1381                    .map_err(map_heed)?
1382                    .is_some()
1383            {
1384                return Err(StoreError::Other(format!(
1385                    "stale Pending cleanup did not durably clear {}",
1386                    to_hex(&item.hash)
1387                )));
1388            }
1389        }
1390        rtxn.commit().map_err(map_heed)?;
1391        Ok(PoolStalePendingCleanupReport {
1392            requested: expected.len(),
1393            declared_bytes,
1394            already_cleaned: exact_pending == 0,
1395        })
1396    }
1397
1398    pub fn pin_sync(&self, hash: &Hash) -> Result<(), StoreError> {
1399        let mut wtxn = self.env.write_txn().map_err(map_heed)?;
1400        let previous = self
1401            .pins
1402            .get(&wtxn, hash)
1403            .map_err(map_heed)?
1404            .map(decode_pin_count)
1405            .transpose()?
1406            .unwrap_or(0);
1407        self.pins
1408            .put(&mut wtxn, hash, &previous.saturating_add(1).to_be_bytes())
1409            .map_err(map_heed)?;
1410        wtxn.commit().map_err(map_heed)
1411    }
1412
1413    pub fn unpin_sync(&self, hash: &Hash) -> Result<(), StoreError> {
1414        let mut wtxn = self.env.write_txn().map_err(map_heed)?;
1415        let count = self
1416            .pins
1417            .get(&wtxn, hash)
1418            .map_err(map_heed)?
1419            .map(decode_pin_count)
1420            .transpose()?
1421            .unwrap_or(0);
1422        if count <= 1 {
1423            self.pins.delete(&mut wtxn, hash).map_err(map_heed)?;
1424        } else {
1425            self.pins
1426                .put(&mut wtxn, hash, &(count - 1).to_be_bytes())
1427                .map_err(map_heed)?;
1428        }
1429        wtxn.commit().map_err(map_heed)
1430    }
1431
1432    pub fn pin_count_sync(&self, hash: &Hash) -> Result<u32, StoreError> {
1433        let rtxn = self.env.read_txn().map_err(map_heed)?;
1434        self.pins
1435            .get(&rtxn, hash)
1436            .map_err(map_heed)?
1437            .map(decode_pin_count)
1438            .transpose()
1439            .map(|count| count.unwrap_or(0))
1440    }
1441
1442    pub fn touch_accessed_sync(&self, hash: &Hash, now: u64) -> Result<bool, StoreError> {
1443        let mut wtxn = self.env.write_txn().map_err(map_heed)?;
1444        if self.locations.get(&wtxn, hash).map_err(map_heed)?.is_none() {
1445            return Ok(false);
1446        }
1447        let previous = self
1448            .last_accessed
1449            .get(&wtxn, hash)
1450            .map_err(map_heed)?
1451            .and_then(self::temperature::AccessRecord::decode)
1452            .unwrap_or_else(|| self::temperature::AccessRecord::new(now));
1453        let access = self::temperature::AccessRecord {
1454            last_accessed_at: now,
1455            ..previous
1456        }
1457        .encode();
1458        self.last_accessed
1459            .put(&mut wtxn, hash, &access)
1460            .map_err(map_heed)?;
1461        wtxn.commit().map_err(map_heed)?;
1462        Ok(true)
1463    }
1464
1465    pub fn touch_many_accessed_sync(&self, hashes: &[Hash], now: u64) -> Result<usize, StoreError> {
1466        let mut wtxn = self.env.write_txn().map_err(map_heed)?;
1467        let mut updated = 0usize;
1468        let mut seen = HashSet::new();
1469        for hash in hashes {
1470            if !seen.insert(*hash) || self.locations.get(&wtxn, hash).map_err(map_heed)?.is_none() {
1471                continue;
1472            }
1473            let previous = self
1474                .last_accessed
1475                .get(&wtxn, hash)
1476                .map_err(map_heed)?
1477                .and_then(self::temperature::AccessRecord::decode)
1478                .unwrap_or_else(|| self::temperature::AccessRecord::new(now));
1479            let access = self::temperature::AccessRecord {
1480                last_accessed_at: now,
1481                ..previous
1482            }
1483            .encode();
1484            self.last_accessed
1485                .put(&mut wtxn, hash, &access)
1486                .map_err(map_heed)?;
1487            updated += 1;
1488        }
1489        wtxn.commit().map_err(map_heed)?;
1490        Ok(updated)
1491    }
1492
1493    pub fn last_accessed_at_sync(&self, hash: &Hash) -> Result<Option<u64>, StoreError> {
1494        let rtxn = self.env.read_txn().map_err(map_heed)?;
1495        self.last_accessed
1496            .get(&rtxn, hash)
1497            .map_err(map_heed)?
1498            .map(|bytes| {
1499                self::temperature::AccessRecord::decode(bytes)
1500                    .map(|access| access.last_accessed_at)
1501                    .ok_or_else(|| StoreError::Other("invalid pool access record".into()))
1502            })
1503            .transpose()
1504    }
1505
1506    pub fn many_last_accessed_at_sync(
1507        &self,
1508        hashes: &[Hash],
1509    ) -> Result<Vec<(Hash, u64)>, StoreError> {
1510        let rtxn = self.env.read_txn().map_err(map_heed)?;
1511        let mut values = Vec::new();
1512        for hash in hashes {
1513            if let Some(value) = self.last_accessed.get(&rtxn, hash).map_err(map_heed)? {
1514                let access = self::temperature::AccessRecord::decode(value)
1515                    .ok_or_else(|| StoreError::Other("invalid pool access record".into()))?;
1516                values.push((*hash, access.last_accessed_at));
1517            }
1518        }
1519        Ok(values)
1520    }
1521
1522    pub fn list(&self) -> Result<Vec<Hash>, StoreError> {
1523        let rtxn = self.env.read_txn().map_err(map_heed)?;
1524        let mut hashes = Vec::new();
1525        for item in self.locations.iter(&rtxn).map_err(map_heed)? {
1526            let (hash, _) = item.map_err(map_heed)?;
1527            let hash: Hash = hash
1528                .try_into()
1529                .map_err(|_| StoreError::Other("invalid pool hash key".into()))?;
1530            hashes.push(hash);
1531        }
1532        Ok(hashes)
1533    }
1534
1535    pub fn stats(&self) -> Result<StoreStats, StoreError> {
1536        let rtxn = self.env.read_txn().map_err(map_heed)?;
1537        let mut stats = StoreStats::default();
1538        for item in self.locations.iter(&rtxn).map_err(map_heed)? {
1539            let (_, location) = item.map_err(map_heed)?;
1540            let location = LocationRecord::decode(location)?;
1541            stats.count = stats.count.saturating_add(1);
1542            stats.bytes = stats.bytes.saturating_add(location.size());
1543        }
1544        for item in self.pins.iter(&rtxn).map_err(map_heed)? {
1545            let (hash, count) = item.map_err(map_heed)?;
1546            if decode_pin_count(count)? == 0 {
1547                continue;
1548            }
1549            let Some(location) = self.locations.get(&rtxn, hash).map_err(map_heed)? else {
1550                continue;
1551            };
1552            stats.pinned_count = stats.pinned_count.saturating_add(1);
1553            stats.pinned_bytes = stats
1554                .pinned_bytes
1555                .saturating_add(LocationRecord::decode(location)?.size());
1556        }
1557        Ok(stats)
1558    }
1559
1560    /// Return conservative physical usage for quota admission without scanning
1561    /// the logical Pool catalog.
1562    ///
1563    /// Every member LMDB maintains its own totals transactionally, so summing
1564    /// those counters is constant in the member count and remains accurate
1565    /// across processes and mixed binary versions. Blobs temporarily duplicated
1566    /// during a move are counted in both members, which is intentionally
1567    /// conservative for writable-space admission. Missing member statistics
1568    /// fail closed rather than undercounting storage.
1569    pub fn writable_physical_stats(&self) -> Result<StoreStats, StoreError> {
1570        self.refresh_members()?;
1571        let manifest = self.read_manifest()?;
1572        let runtime = self
1573            .runtime
1574            .read()
1575            .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?;
1576        let mut total = StoreStats::default();
1577        for member in manifest.members {
1578            let store = runtime.stores.get(&member.id).ok_or_else(|| {
1579                let detail = runtime
1580                    .errors
1581                    .get(&member.id)
1582                    .map(String::as_str)
1583                    .unwrap_or("member store is not open");
1584                StoreError::Other(format!(
1585                    "pool member {} unavailable for physical quota accounting: {detail}",
1586                    member.id
1587                ))
1588            })?;
1589            let stats = store.stats().map_err(|error| {
1590                StoreError::Other(format!(
1591                    "pool member {} unavailable for physical quota accounting: {error}",
1592                    member.id
1593                ))
1594            })?;
1595            let count = u64::try_from(stats.count)
1596                .map_err(|_| StoreError::Other("pool physical blob count exceeds u64".into()))?;
1597            let pinned_count = u64::try_from(stats.pinned_count)
1598                .map_err(|_| StoreError::Other("pool physical pin count exceeds u64".into()))?;
1599            total.count = total
1600                .count
1601                .checked_add(count)
1602                .ok_or_else(|| StoreError::Other("pool physical blob count overflow".into()))?;
1603            total.bytes = total
1604                .bytes
1605                .checked_add(stats.total_bytes)
1606                .ok_or_else(|| StoreError::Other("pool physical byte count overflow".into()))?;
1607            total.pinned_count = total
1608                .pinned_count
1609                .checked_add(pinned_count)
1610                .ok_or_else(|| StoreError::Other("pool physical pin count overflow".into()))?;
1611            total.pinned_bytes = total
1612                .pinned_bytes
1613                .checked_add(stats.pinned_bytes)
1614                .ok_or_else(|| StoreError::Other("pool physical pinned bytes overflow".into()))?;
1615        }
1616        Ok(total)
1617    }
1618
1619    pub fn force_sync(&self) -> Result<(), StoreError> {
1620        self.env.force_sync().map_err(map_heed)?;
1621        self.refresh_members()?;
1622        let runtime = self
1623            .runtime
1624            .read()
1625            .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?;
1626        for store in runtime.stores.values() {
1627            store.force_sync()?;
1628        }
1629        Ok(())
1630    }
1631
1632    /// Current LMDB virtual map size for the Pool catalog.
1633    pub fn catalog_map_size_bytes(&self) -> u64 {
1634        u64::try_from(self.env.info().map_size).unwrap_or(u64::MAX)
1635    }
1636
1637    /// Exhaustively prove every controlled member's physical indexes and
1638    /// persisted aggregate counters before the first migration mutation.
1639    ///
1640    /// Controlled reopen performs only constant-time structural checks so
1641    /// bounded mapping epochs do not rescan the complete target. The caller
1642    /// must hold the external writer fence while this one full proof runs and
1643    /// must keep that fence held until migration completion.
1644    pub fn validate_controlled_member_state_exact(&self) -> Result<(), StoreError> {
1645        if self.expected_manifest_sha256.is_none() {
1646            return Err(StoreError::Other(
1647                "exact controlled member validation requires an exact manifest SHA-256".into(),
1648            ));
1649        }
1650        self.refresh_members()?;
1651        let manifest = self.read_manifest()?;
1652        let runtime = self
1653            .runtime
1654            .read()
1655            .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?;
1656        if runtime.stores.len() != manifest.members.len() {
1657            return Err(StoreError::Other(
1658                "not every authority-pinned Pool member is open for exact validation".into(),
1659            ));
1660        }
1661        for member in manifest.members {
1662            let store = runtime.stores.get(&member.id).ok_or_else(|| {
1663                StoreError::Other(format!(
1664                    "controlled Pool member {} is unavailable for exact validation",
1665                    member.id
1666                ))
1667            })?;
1668            store.validate_exact_member_state().map_err(|error| {
1669                StoreError::Other(format!(
1670                    "controlled Pool member {} failed exact state validation: {error}",
1671                    member.id
1672                ))
1673            })?;
1674        }
1675        Ok(())
1676    }
1677
1678    /// Revalidate every controlled catalog/member authority without forcing
1679    /// unrelated writer commits to storage.
1680    ///
1681    /// This is sufficient before advancing a migration cursor when the
1682    /// migration observed exact `Stored` rows and performed no target writes.
1683    /// Any batch that attempted target writes must use
1684    /// [`Self::validate_controlled_authority_and_sync`] instead.
1685    pub fn validate_controlled_authority(&self) -> Result<(), StoreError> {
1686        if self.expected_manifest_sha256.is_none() {
1687            return Err(StoreError::Other(
1688                "controlled authority validation requires an exact manifest SHA-256".into(),
1689            ));
1690        }
1691        self.refresh_members()?;
1692        {
1693            let runtime = self
1694                .runtime
1695                .read()
1696                .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?;
1697            if runtime.stores.len() != self.member_runtime_paths.len() {
1698                return Err(StoreError::Other(
1699                    "not every authority-pinned Pool member is open".into(),
1700                ));
1701            }
1702        }
1703        self.refresh_members()
1704    }
1705
1706    /// Revalidate every controlled authority and force the catalog/member
1707    /// commits durable before an external migration cursor may advance.
1708    pub fn validate_controlled_authority_and_sync(&self) -> Result<(), StoreError> {
1709        self.validate_controlled_authority()?;
1710        self.env.force_sync().map_err(map_heed)?;
1711        {
1712            let runtime = self
1713                .runtime
1714                .read()
1715                .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?;
1716            for store in runtime.stores.values() {
1717                store.force_sync()?;
1718            }
1719        }
1720        self.validate_controlled_authority()
1721    }
1722
1723    fn refresh_members(&self) -> Result<(), StoreError> {
1724        let (manifest, manifest_sha256) = self.read_manifest_with_identity()?;
1725        if let Some(expected) = self.expected_manifest_sha256 {
1726            if manifest_sha256 != expected {
1727                return Err(StoreError::Other(format!(
1728                    "live pool manifest SHA-256 differs from controlled authority: expected {}, found {}",
1729                    to_hex(&expected),
1730                    to_hex(&manifest_sha256)
1731                )));
1732            }
1733            validate_controlled_manifest_members(&manifest, &self.member_runtime_paths)?;
1734        }
1735        let mut runtime = self
1736            .runtime
1737            .write()
1738            .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?;
1739        if runtime.generation == Some(manifest.generation) {
1740            return Ok(());
1741        }
1742
1743        let configured = manifest
1744            .members
1745            .iter()
1746            .map(|member| member.id)
1747            .collect::<HashSet<_>>();
1748        runtime.stores.retain(|id, _| configured.contains(id));
1749        runtime.read_gates.retain(|id, _| configured.contains(id));
1750        runtime.write_gates.retain(|id, _| configured.contains(id));
1751        runtime.errors.retain(|id, _| configured.contains(id));
1752        for member in &manifest.members {
1753            let read_gate = runtime
1754                .read_gates
1755                .entry(member.id)
1756                .or_insert_with(|| {
1757                    Arc::new(ConcurrencyGate::new(member.config.max_read_concurrency))
1758                })
1759                .clone();
1760            read_gate.set_limit(member.config.max_read_concurrency)?;
1761            let write_gate = runtime
1762                .write_gates
1763                .entry(member.id)
1764                .or_insert_with(|| {
1765                    Arc::new(ConcurrencyGate::new(member.config.max_write_concurrency))
1766                })
1767                .clone();
1768            write_gate.set_limit(member.config.max_write_concurrency)?;
1769            if runtime.stores.contains_key(&member.id) {
1770                continue;
1771            }
1772            let config = self.runtime_member_config(member)?;
1773            let pinned_identity = self
1774                .member_runtime_paths
1775                .get(&member.id)
1776                .map(|binding| binding.lmdb_identity);
1777            match open_member_store(member.id, &config, pinned_identity) {
1778                Ok(store) => {
1779                    runtime.stores.insert(member.id, Arc::new(store));
1780                    runtime.errors.remove(&member.id);
1781                }
1782                Err(error) => {
1783                    if self.expected_manifest_sha256.is_some() {
1784                        return Err(StoreError::Other(format!(
1785                            "controlled Pool member {} failed its exact open: {error}",
1786                            member.id
1787                        )));
1788                    }
1789                    runtime.errors.insert(member.id, error.to_string());
1790                }
1791            }
1792        }
1793        runtime.generation = Some(manifest.generation);
1794        drop(runtime);
1795        self.adaptive
1796            .lock()
1797            .map_err(|_| StoreError::Other("pool adaptive lock poisoned".into()))?
1798            .retain(&configured);
1799        Ok(())
1800    }
1801
1802    fn runtime_member_config(&self, member: &MemberRecord) -> Result<PoolMemberConfig, StoreError> {
1803        let Some(binding) = self.member_runtime_paths.get(&member.id) else {
1804            if self.member_runtime_paths.is_empty() {
1805                return Ok(member.config.clone());
1806            }
1807            return Err(StoreError::Other(format!(
1808                "pool member {} has no pinned runtime path",
1809                member.id
1810            )));
1811        };
1812        if !configured_paths_match(&binding.configured_path, &member.config.path)
1813            || !configured_optional_paths_match(
1814                binding.configured_external_path.as_deref(),
1815                member.config.external_blob_dir.as_deref(),
1816            )
1817        {
1818            return Err(StoreError::Other(format!(
1819                "pool member {} paths differ from pinned runtime authority",
1820                member.id
1821            )));
1822        }
1823        let mut config = member.config.clone();
1824        config.path = binding.runtime_path.clone();
1825        config.external_blob_dir = binding.runtime_external_path.clone();
1826        Ok(config)
1827    }
1828
1829    fn get_member(&self, id: PoolMemberId) -> Result<Arc<LmdbBlobStore>, StoreError> {
1830        self.refresh_members()?;
1831        if let Some(store) = self
1832            .runtime
1833            .read()
1834            .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?
1835            .stores
1836            .get(&id)
1837            .cloned()
1838        {
1839            return Ok(store);
1840        }
1841
1842        let member = self
1843            .read_manifest()?
1844            .members
1845            .into_iter()
1846            .find(|member| member.id == id)
1847            .ok_or_else(|| StoreError::Other(format!("unknown pool member {id}")))?;
1848        let config = self.runtime_member_config(&member)?;
1849        let pinned_identity = self
1850            .member_runtime_paths
1851            .get(&id)
1852            .map(|binding| binding.lmdb_identity);
1853        match open_member_store(id, &config, pinned_identity) {
1854            Ok(store) => {
1855                let store = Arc::new(store);
1856                let mut runtime = self
1857                    .runtime
1858                    .write()
1859                    .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?;
1860                runtime.stores.insert(id, Arc::clone(&store));
1861                runtime.errors.remove(&id);
1862                Ok(store)
1863            }
1864            Err(error) => {
1865                self.record_member_failure(id, false);
1866                let mut runtime = self
1867                    .runtime
1868                    .write()
1869                    .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?;
1870                runtime.errors.insert(id, error.to_string());
1871                Err(error)
1872            }
1873        }
1874    }
1875
1876    fn member_state(&self, id: PoolMemberId) -> Result<Option<PoolMemberState>, StoreError> {
1877        Ok(self
1878            .read_manifest()?
1879            .members
1880            .into_iter()
1881            .find(|member| member.id == id)
1882            .map(|member| member.state))
1883    }
1884
1885    fn member_gate(
1886        &self,
1887        id: PoolMemberId,
1888        write: bool,
1889    ) -> Result<Arc<ConcurrencyGate>, StoreError> {
1890        self.refresh_members()?;
1891        let runtime = self
1892            .runtime
1893            .read()
1894            .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?;
1895        let gates = if write {
1896            &runtime.write_gates
1897        } else {
1898            &runtime.read_gates
1899        };
1900        gates
1901            .get(&id)
1902            .cloned()
1903            .ok_or_else(|| StoreError::Other(format!("unknown pool member {id}")))
1904    }
1905
1906    fn choose_write_member(
1907        &self,
1908        incoming_bytes: u64,
1909        exclude: Option<PoolMemberId>,
1910    ) -> Result<PoolMemberId, StoreError> {
1911        self.choose_write_member_with_reserved(incoming_bytes, exclude, &HashMap::new())
1912    }
1913
1914    fn choose_write_member_with_reserved(
1915        &self,
1916        incoming_bytes: u64,
1917        exclude: Option<PoolMemberId>,
1918        reserved_bytes: &HashMap<PoolMemberId, u64>,
1919    ) -> Result<PoolMemberId, StoreError> {
1920        let excluded = exclude.into_iter().collect::<HashSet<_>>();
1921        self.choose_write_member_excluding(incoming_bytes, &excluded, reserved_bytes)
1922    }
1923
1924    fn choose_write_member_excluding(
1925        &self,
1926        incoming_bytes: u64,
1927        excluded: &HashSet<PoolMemberId>,
1928        reserved_bytes: &HashMap<PoolMemberId, u64>,
1929    ) -> Result<PoolMemberId, StoreError> {
1930        self.refresh_members()?;
1931        let manifest = self.read_manifest()?;
1932        let runtime = self
1933            .runtime
1934            .read()
1935            .map_err(|_| StoreError::Other("pool runtime lock poisoned".into()))?;
1936        let mut candidates = Vec::new();
1937        let mut below_high_watermark = Vec::new();
1938        for member in manifest.members.iter().filter(|member| {
1939            member.state == PoolMemberState::Active && !excluded.contains(&member.id)
1940        }) {
1941            let Some(store) = runtime.stores.get(&member.id) else {
1942                continue;
1943            };
1944            let stats = match store.stats() {
1945                Ok(stats) => stats,
1946                Err(_) => continue,
1947            };
1948            let effective_bytes = stats
1949                .total_bytes
1950                .saturating_add(reserved_bytes.get(&member.id).copied().unwrap_or(0));
1951            if member.config.capacity_bytes > 0
1952                && effective_bytes.saturating_add(incoming_bytes) > member.config.capacity_bytes
1953            {
1954                continue;
1955            }
1956            candidates.push((member.id, effective_bytes, member.config.capacity_bytes));
1957            let projected_fill = effective_bytes
1958                .saturating_add(incoming_bytes)
1959                .saturating_mul(100)
1960                .saturating_div(member.config.capacity_bytes)
1961                .min(100);
1962            if projected_fill <= u64::from(member.config.temperature_high_watermark_percent) {
1963                below_high_watermark.push((
1964                    member.id,
1965                    effective_bytes,
1966                    member.config.capacity_bytes,
1967                ));
1968            }
1969        }
1970        drop(runtime);
1971        if !below_high_watermark.is_empty() {
1972            candidates = below_high_watermark;
1973        }
1974        self.adaptive
1975            .lock()
1976            .map_err(|_| StoreError::Other("pool adaptive lock poisoned".into()))?
1977            .choose_write(&candidates)
1978            .ok_or_else(|| StoreError::Other("no writable pool member has capacity".into()))
1979    }
1980
1981    fn repair_location(
1982        &self,
1983        hash: Hash,
1984        data: &[u8],
1985        expected: LocationRecord,
1986    ) -> Result<bool, StoreError> {
1987        self.repair_location_excluding(hash, data, expected, HashSet::new())
1988    }
1989
1990    fn repair_location_excluding(
1991        &self,
1992        hash: Hash,
1993        data: &[u8],
1994        expected: LocationRecord,
1995        mut excluded: HashSet<PoolMemberId>,
1996    ) -> Result<bool, StoreError> {
1997        let preferred = expected.preferred_member();
1998        let mut next = (!excluded.contains(&preferred)
1999            && self.member_state(preferred)? == Some(PoolMemberState::Active))
2000        .then_some(preferred);
2001        let mut last_error = None;
2002        let (target, inserted) = loop {
2003            let target = match next.take() {
2004                Some(target) => target,
2005                None => match self.choose_write_member_excluding(
2006                    data.len() as u64,
2007                    &excluded,
2008                    &HashMap::new(),
2009                ) {
2010                    Ok(target) => target,
2011                    Err(error) => {
2012                        return Err(last_error.unwrap_or(error));
2013                    }
2014                },
2015            };
2016            let result = self
2017                .get_member(target)
2018                .and_then(|store| self.write_verified_member(target, &store, hash, data));
2019            match result {
2020                Ok(inserted) => break (target, inserted),
2021                Err(error) => {
2022                    excluded.insert(target);
2023                    last_error = Some(error);
2024                }
2025            }
2026        };
2027        let stored = LocationRecord::Stored {
2028            member: target,
2029            size: data.len() as u64,
2030        };
2031
2032        let mut wtxn = self.env.write_txn().map_err(map_heed)?;
2033        let current = self
2034            .locations
2035            .get(&wtxn, &hash)
2036            .map_err(map_heed)?
2037            .map(LocationRecord::decode)
2038            .transpose()?;
2039        if current == Some(expected) {
2040            self.set_location_txn(&mut wtxn, hash, Some(stored))?;
2041            wtxn.commit().map_err(map_heed)?;
2042            return Ok(inserted);
2043        }
2044        drop(wtxn);
2045
2046        if let Some(current) = current {
2047            if self.read_verified_location(&hash, current)?.is_some() {
2048                return Ok(false);
2049            }
2050        }
2051        Err(StoreError::Other(format!(
2052            "pool location changed while repairing {hash:?}"
2053        )))
2054    }
2055
2056    fn read_verified_location(
2057        &self,
2058        hash: &Hash,
2059        location: LocationRecord,
2060    ) -> Result<Option<Vec<u8>>, StoreError> {
2061        let mut ids = match location {
2062            LocationRecord::Pending { member, .. } | LocationRecord::Stored { member, .. } => {
2063                vec![member]
2064            }
2065            LocationRecord::Moving { source, target, .. } => vec![target, source],
2066        };
2067        self.adaptive
2068            .lock()
2069            .map_err(|_| StoreError::Other("pool adaptive lock poisoned".into()))?
2070            .order_reads(&mut ids);
2071
2072        let mut first_error = None;
2073        for id in ids {
2074            let store = match self.get_member(id) {
2075                Ok(store) => store,
2076                Err(error) => {
2077                    first_error.get_or_insert(error);
2078                    continue;
2079                }
2080            };
2081            match self.read_verified_member(id, &store, hash) {
2082                Ok(Some(data)) => return Ok(Some(data)),
2083                Ok(None) => {}
2084                Err(error) => {
2085                    first_error.get_or_insert(error);
2086                }
2087            }
2088        }
2089        match first_error {
2090            Some(error) => Err(error),
2091            None => Ok(None),
2092        }
2093    }
2094
2095    fn read_verified_member(
2096        &self,
2097        id: PoolMemberId,
2098        store: &LmdbBlobStore,
2099        hash: &Hash,
2100    ) -> Result<Option<Vec<u8>>, StoreError> {
2101        let gate = self.member_gate(id, false)?;
2102        let _permit = gate.acquire()?;
2103        let started = Instant::now();
2104        let result = store.get_sync(hash);
2105        match result {
2106            Ok(Some(data)) if sha256(&data) == *hash => {
2107                self.record_read(id, started.elapsed(), true);
2108                Ok(Some(data))
2109            }
2110            Ok(Some(_)) => {
2111                self.record_read(id, started.elapsed(), false);
2112                Err(StoreError::Other(format!(
2113                    "pool member {id} returned corrupt bytes"
2114                )))
2115            }
2116            Ok(None) => {
2117                self.record_read(id, started.elapsed(), true);
2118                Ok(None)
2119            }
2120            Err(error) => {
2121                self.record_read(id, started.elapsed(), false);
2122                Err(error)
2123            }
2124        }
2125    }
2126
2127    fn write_verified_member(
2128        &self,
2129        id: PoolMemberId,
2130        store: &LmdbBlobStore,
2131        hash: Hash,
2132        data: &[u8],
2133    ) -> Result<bool, StoreError> {
2134        let gate = self.member_gate(id, true)?;
2135        let _permit = gate.acquire()?;
2136        if let Some(existing) = store.get_sync(&hash)? {
2137            if sha256(&existing) == hash {
2138                return Ok(false);
2139            }
2140            store.delete_sync(&hash)?;
2141        }
2142        let started = Instant::now();
2143        let result = store.put_sync(hash, data);
2144        let success = result.is_ok();
2145        self.record_write(id, started.elapsed(), data.len(), success);
2146        let inserted = result?;
2147        let written = store
2148            .get_sync(&hash)?
2149            .ok_or_else(|| StoreError::Other(format!("pool member {id} lost a committed write")))?;
2150        if sha256(&written) != hash {
2151            self.record_member_failure(id, true);
2152            return Err(StoreError::Other(format!(
2153                "pool member {id} committed corrupt bytes"
2154            )));
2155        }
2156        Ok(inserted)
2157    }
2158
2159    fn delete_member_blob(
2160        &self,
2161        id: PoolMemberId,
2162        store: &LmdbBlobStore,
2163        hash: &Hash,
2164    ) -> Result<bool, StoreError> {
2165        let gate = self.member_gate(id, true)?;
2166        let _permit = gate.acquire()?;
2167        store.delete_sync(hash)
2168    }
2169
2170    fn record_read(&self, id: PoolMemberId, elapsed: Duration, success: bool) {
2171        if let Ok(mut adaptive) = self.adaptive.lock() {
2172            adaptive.record_read(id, elapsed, success);
2173        }
2174    }
2175
2176    fn record_write(&self, id: PoolMemberId, elapsed: Duration, bytes: usize, success: bool) {
2177        if let Ok(mut adaptive) = self.adaptive.lock() {
2178            adaptive.record_write(id, elapsed, bytes, success);
2179        }
2180    }
2181
2182    fn record_member_failure(&self, id: PoolMemberId, write: bool) {
2183        if write {
2184            self.record_write(id, Duration::ZERO, 0, false);
2185        } else {
2186            self.record_read(id, Duration::ZERO, false);
2187        }
2188    }
2189
2190    fn sample_temperature_access(&self, hash: Hash, location: LocationRecord) {
2191        if !self.temperature_config.enabled {
2192            return;
2193        }
2194        let sample_rate = u64::from(self.temperature_config.read_sample_rate.max(1));
2195        let access = self
2196            .temperature_access_counter
2197            .fetch_add(1, Ordering::Relaxed)
2198            .saturating_add(1);
2199        if !access.is_multiple_of(sample_rate) {
2200            return;
2201        }
2202        if let Ok(mut temperature) = self.temperature.lock() {
2203            temperature
2204                .samples
2205                .observe(hash, location, unix_timestamp_now());
2206        }
2207    }
2208}
2209
2210#[async_trait]
2211impl Store for PoolStore {
2212    async fn put(&self, hash: Hash, data: Vec<u8>) -> Result<bool, StoreError> {
2213        self.put_sync(hash, &data)
2214    }
2215
2216    async fn put_many(&self, items: Vec<(Hash, Vec<u8>)>) -> Result<usize, StoreError> {
2217        self.put_many_sync(&items)
2218    }
2219
2220    async fn put_many_optimistic(&self, items: Vec<(Hash, Vec<u8>)>) -> Result<usize, StoreError> {
2221        self.put_many_optimistic_sync(&items)
2222    }
2223
2224    async fn get(&self, hash: &Hash) -> Result<Option<Vec<u8>>, StoreError> {
2225        self.get_sync(hash)
2226    }
2227
2228    async fn get_range(
2229        &self,
2230        hash: &Hash,
2231        start: u64,
2232        end_inclusive: u64,
2233    ) -> Result<Option<Vec<u8>>, StoreError> {
2234        self.get_range_sync(hash, start, end_inclusive)
2235    }
2236
2237    async fn blob_size(&self, hash: &Hash) -> Result<Option<u64>, StoreError> {
2238        self.blob_size_sync(hash)
2239    }
2240
2241    async fn has(&self, hash: &Hash) -> Result<bool, StoreError> {
2242        self.exists(hash)
2243    }
2244
2245    async fn delete(&self, hash: &Hash) -> Result<bool, StoreError> {
2246        self.delete_sync(hash)
2247    }
2248
2249    async fn delete_many(&self, hashes: Vec<Hash>) -> Result<usize, StoreError> {
2250        self.delete_many_sync(&hashes)
2251    }
2252
2253    async fn stats(&self) -> StoreStats {
2254        PoolStore::stats(self).unwrap_or_default()
2255    }
2256
2257    async fn pin(&self, hash: &Hash) -> Result<(), StoreError> {
2258        self.pin_sync(hash)
2259    }
2260
2261    async fn unpin(&self, hash: &Hash) -> Result<(), StoreError> {
2262        self.unpin_sync(hash)
2263    }
2264
2265    fn pin_count(&self, hash: &Hash) -> u32 {
2266        self.pin_count_sync(hash).unwrap_or(0)
2267    }
2268}
2269
2270fn encode_manifest(manifest: &PoolManifest) -> Result<Vec<u8>, StoreError> {
2271    rmp_serde::to_vec_named(manifest)
2272        .map_err(|error| StoreError::Other(format!("encode pool manifest: {error}")))
2273}
2274
2275fn decode_manifest(bytes: &[u8]) -> Result<PoolManifest, StoreError> {
2276    let manifest: PoolManifest = rmp_serde::from_slice(bytes)
2277        .map_err(|error| StoreError::Other(format!("decode pool manifest: {error}")))?;
2278    if manifest.version != 1 {
2279        return Err(StoreError::Other(format!(
2280            "unsupported pool manifest version {}",
2281            manifest.version
2282        )));
2283    }
2284    Ok(manifest)
2285}
2286
2287fn member_hash_key(member: PoolMemberId, hash: Hash) -> [u8; 48] {
2288    let mut key = [0u8; 48];
2289    key[..16].copy_from_slice(member.as_bytes());
2290    key[16..].copy_from_slice(&hash);
2291    key
2292}
2293
2294fn decode_pin_count(bytes: &[u8]) -> Result<u32, StoreError> {
2295    Ok(u32::from_be_bytes(bytes.try_into().map_err(|_| {
2296        StoreError::Other("invalid pool pin count".into())
2297    })?))
2298}
2299
2300fn unix_timestamp_now() -> u64 {
2301    SystemTime::now()
2302        .duration_since(UNIX_EPOCH)
2303        .unwrap_or_default()
2304        .as_secs()
2305}
2306
2307fn put_many_report(
2308    total: usize,
2309    ordered: &[(Hash, u64)],
2310    inserted: &HashSet<Hash>,
2311) -> PutManyReport {
2312    let mut report = PutManyReport {
2313        total,
2314        ..PutManyReport::default()
2315    };
2316    for (hash, bytes) in ordered {
2317        if inserted.contains(hash) {
2318            report.inserted = report.inserted.saturating_add(1);
2319            report.inserted_bytes = report.inserted_bytes.saturating_add(*bytes);
2320            report.inserted_hashes.push(*hash);
2321        }
2322    }
2323    report
2324}
2325
2326fn map_heed(error: heed::Error) -> StoreError {
2327    StoreError::Other(error.to_string())
2328}