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 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 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 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 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 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 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 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 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 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 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 pub fn catalog_map_size_bytes(&self) -> u64 {
1634 u64::try_from(self.env.info().map_size).unwrap_or(u64::MAX)
1635 }
1636
1637 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 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 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}