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