1use crate::warm_start::key::Fingerprint;
11use serde::{Deserialize, Serialize};
12use sha2::{Digest, Sha256};
13use std::collections::HashMap;
14use std::fs;
15use std::io::{self, Write as _};
16use std::path::{Path, PathBuf};
17use std::sync::atomic::{AtomicU64, Ordering};
18use std::sync::{Arc, Mutex, OnceLock};
19use std::time::{Duration, SystemTime, UNIX_EPOCH};
20
21pub(crate) const SCHEMA_VERSION: u32 = 1;
24
25pub(crate) const DEFAULT_SIZE_BUDGET_BYTES: u64 = 1024 * 1024 * 1024;
27
28pub(crate) const DEFAULT_TTL_SECS: u64 = 60 * 60 * 24 * 30;
30
31#[derive(Debug, thiserror::Error)]
32pub enum StoreError {
33 #[error("io: {0}")]
34 Io(#[from] io::Error),
35 #[error("json: {0}")]
36 Json(#[from] serde_json::Error),
37}
38
39#[derive(Debug, Clone)]
41pub struct WarmStartEntry {
42 pub payload: Vec<u8>,
43 pub objective: Option<f64>,
44 pub iteration: Option<u64>,
45 pub written_unix_secs: u64,
46 pub kind: EntryKind,
47}
48
49#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
50pub enum EntryKind {
51 Checkpoint,
53 Final,
55}
56
57#[derive(Debug, Clone, Serialize, Deserialize)]
58struct OnDiskMeta {
59 schema_version: u32,
60 written_unix_secs: u64,
61 #[serde(default)]
65 written_nanos: u32,
66 objective: Option<f64>,
67 iteration: Option<u64>,
68 kind: EntryKind,
69 checksum_hex: String,
70 payload_bytes: u64,
71 #[serde(default)]
74 accessed: bool,
75 #[serde(default)]
86 accessed_unix_secs: u64,
87 #[serde(default)]
88 accessed_nanos: u32,
89}
90
91fn meta_activity_nanos(meta: &OnDiskMeta) -> u128 {
96 let written = (meta.written_unix_secs as u128) * 1_000_000_000u128 + meta.written_nanos as u128;
97 let accessed =
98 (meta.accessed_unix_secs as u128) * 1_000_000_000u128 + meta.accessed_nanos as u128;
99 written.max(accessed)
100}
101
102#[derive(Debug, Clone)]
103pub struct StoreOptions {
104 pub size_budget_bytes: u64,
105 pub ttl: Duration,
106}
107
108impl Default for StoreOptions {
109 fn default() -> Self {
110 Self {
111 size_budget_bytes: DEFAULT_SIZE_BUDGET_BYTES,
112 ttl: Duration::from_secs(DEFAULT_TTL_SECS),
113 }
114 }
115}
116
117#[derive(Debug)]
118pub struct WarmStartStore {
119 root: PathBuf,
120 opts: StoreOptions,
121 index: Arc<Mutex<MetadataIndex>>,
124 byte_total: Arc<AtomicU64>,
133 save_counter: Arc<AtomicU64>,
136 last_evict_root_mtime: Arc<Mutex<Option<SystemTime>>>,
151 test_time_offset_ns: AtomicU64,
160}
161
162impl Clone for WarmStartStore {
163 fn clone(&self) -> Self {
164 Self {
165 root: self.root.clone(),
166 opts: self.opts.clone(),
167 index: Arc::clone(&self.index),
168 byte_total: Arc::clone(&self.byte_total),
173 save_counter: Arc::clone(&self.save_counter),
174 last_evict_root_mtime: Arc::clone(&self.last_evict_root_mtime),
175 test_time_offset_ns: AtomicU64::new(self.test_time_offset_ns.load(Ordering::Relaxed)),
176 }
177 }
178}
179
180impl WarmStartStore {
181 pub fn open(root: PathBuf, opts: StoreOptions) -> Result<Self, StoreError> {
183 fs::create_dir_all(&root)?;
184 Ok(Self {
185 root,
186 opts,
187 index: Arc::new(Mutex::new(MetadataIndex::default())),
188 byte_total: Arc::new(AtomicU64::new(0)),
189 save_counter: Arc::new(AtomicU64::new(0)),
190 last_evict_root_mtime: Arc::new(Mutex::new(None)),
191 test_time_offset_ns: AtomicU64::new(0),
192 })
193 }
194
195 pub fn root(&self) -> &Path {
196 &self.root
197 }
198
199 pub fn options(&self) -> &StoreOptions {
200 &self.opts
201 }
202
203 fn key_dir(&self, key: &Fingerprint) -> PathBuf {
204 self.root.join(key.to_hex())
205 }
206
207 pub fn lookup(&self, key: &Fingerprint) -> Result<Option<WarmStartEntry>, StoreError> {
215 self.lookup_with(key, LookupMode::Best)
216 }
217
218 pub fn lookup_latest(&self, key: &Fingerprint) -> Result<Option<WarmStartEntry>, StoreError> {
226 self.lookup_with(key, LookupMode::Latest)
227 }
228
229 fn lookup_with(
230 &self,
231 key: &Fingerprint,
232 mode: LookupMode,
233 ) -> Result<Option<WarmStartEntry>, StoreError> {
234 let dir = self.key_dir(key);
235 if !dir.exists() {
236 lookup_cache_invalidate(&LookupCacheKey { fp: *key, mode });
240 self.metadata_index_remove_key(key);
241 return Ok(None);
242 }
243 let cache_key = LookupCacheKey { fp: *key, mode };
252 let now_nanos = self.nanos_now();
253 if let Some(hit) = lookup_cache_get(&cache_key) {
254 if let Ok(md) = fs::metadata(&hit.meta_path)
255 && md.modified().ok() == Some(hit.meta_mtime)
256 {
257 let expired = self.opts.ttl.as_nanos() > 0
258 && now_nanos.saturating_sub(hit.write_nanos) >= self.opts.ttl.as_nanos();
259 if !expired {
260 let entry = self.touch_lookup_hit(&hit.meta_path, hit.entry)?;
261 return Ok(Some(entry));
262 }
263 lookup_cache_invalidate(&cache_key);
264 let bin = hit.meta_path.with_extension("bin");
265 fs::remove_file(&hit.meta_path).ok();
266 fs::remove_file(&bin).ok();
267 self.metadata_index_remove(&hit.meta_path);
269 return Ok(None);
270 }
271 lookup_cache_invalidate(&cache_key);
272 }
273 let mut best: Option<(OnDiskMeta, PathBuf)> = None;
279 for scanned in self.scan_key_dir(&dir, now_nanos) {
280 let take = match best {
281 None => true,
282 Some((ref cur, _)) => mode.better(&scanned.meta, cur),
283 };
284 if take {
285 best = Some((scanned.meta, scanned.meta_path));
286 }
287 }
288 let (meta, meta_path) = match best {
289 Some(b) => b,
290 None => {
291 lookup_cache_invalidate(&cache_key);
292 return Ok(None);
293 }
294 };
295 let bin_path = meta_path.with_extension("bin");
296 let payload = match fs::read(&bin_path) {
297 Ok(v) => v,
298 Err(_) => return Ok(None),
299 };
300 if checksum_hex(&payload) != meta.checksum_hex {
302 fs::remove_file(&meta_path).ok();
303 fs::remove_file(&bin_path).ok();
304 lookup_cache_invalidate(&cache_key);
305 self.metadata_index_remove(&meta_path);
306 return Ok(None);
307 }
308 let entry = WarmStartEntry {
309 payload,
310 objective: meta.objective,
311 iteration: meta.iteration,
312 written_unix_secs: meta.written_unix_secs,
313 kind: meta.kind,
314 };
315 let (meta, entry) = self.touch_lookup_meta(&meta_path, meta, entry)?;
316 if let Ok(md) = fs::metadata(&meta_path)
322 && let Ok(mtime) = md.modified()
323 {
324 let write_nanos = meta_activity_nanos(&meta);
325 lookup_cache_insert(
326 cache_key,
327 CachedLookup {
328 meta_path: meta_path.clone(),
329 meta_mtime: mtime,
330 write_nanos,
331 entry: entry.clone(),
332 },
333 );
334 }
335 Ok(Some(entry))
336 }
337
338 pub fn save(
341 &self,
342 key: &Fingerprint,
343 payload: &[u8],
344 objective: Option<f64>,
345 iteration: Option<u64>,
346 kind: EntryKind,
347 ) -> Result<String, StoreError> {
348 let run_id = self.fresh_run_id();
349 self.save_overwrite(key, &run_id, payload, objective, iteration, kind)?;
350 Ok(run_id)
351 }
352
353 pub fn save_overwrite(
356 &self,
357 key: &Fingerprint,
358 run_id: &str,
359 payload: &[u8],
360 objective: Option<f64>,
361 iteration: Option<u64>,
362 kind: EntryKind,
363 ) -> Result<(), StoreError> {
364 lookup_cache_invalidate(&LookupCacheKey {
371 fp: *key,
372 mode: LookupMode::Best,
373 });
374 lookup_cache_invalidate(&LookupCacheKey {
375 fp: *key,
376 mode: LookupMode::Latest,
377 });
378 let dir = self.key_dir(key);
379 let pid = std::process::id();
380 let checksum = checksum_hex(payload);
382 let objective_finite = objective.filter(|o| o.is_finite());
383 let nonce = self.nanos_now();
414 let bin_final = dir.join(format!("{run_id}.bin"));
415 let meta_final = dir.join(format!("{run_id}.json"));
416 let mut attempt = 0u8;
417 let build_meta_json = |secs: u64, subsec_nanos: u32| -> Result<Vec<u8>, StoreError> {
418 let meta = OnDiskMeta {
419 schema_version: SCHEMA_VERSION,
420 written_unix_secs: secs,
421 written_nanos: subsec_nanos,
422 objective: objective_finite,
423 iteration,
424 kind,
425 checksum_hex: checksum.clone(),
426 payload_bytes: payload.len() as u64,
427 accessed: false,
428 accessed_unix_secs: 0,
429 accessed_nanos: 0,
430 };
431 Ok(serde_json::to_vec_pretty(&meta)?)
432 };
433 loop {
434 let bin_tmp = dir.join(format!("{run_id}.bin.tmp.{pid}.{nonce}.{attempt}"));
435 let meta_tmp = dir.join(format!("{run_id}.json.tmp.{pid}.{nonce}.{attempt}"));
436 let stamp_fn = || self.unix_now_parts();
437 let build_meta_for_io = |secs: u64, subsec_nanos: u32| -> io::Result<Vec<u8>> {
438 build_meta_json(secs, subsec_nanos)
439 .map_err(|e| io::Error::other(format!("meta build: {e:?}")))
440 };
441 match write_and_promote_entry(&EntryWrite {
442 dir: &dir,
443 bin_tmp: &bin_tmp,
444 meta_tmp: &meta_tmp,
445 payload,
446 bin_final: &bin_final,
447 meta_final: &meta_final,
448 stamp_fn: &stamp_fn,
449 build_meta_json: &build_meta_for_io,
450 }) {
451 Ok(()) => break,
452 Err(e) if e.kind() == io::ErrorKind::NotFound && attempt == 0 => {
453 fs::remove_file(&bin_tmp).ok();
457 fs::remove_file(&meta_tmp).ok();
458 attempt += 1;
459 continue;
460 }
461 Err(e) => {
462 fs::remove_file(&bin_tmp).ok();
463 fs::remove_file(&meta_tmp).ok();
464 fs::remove_file(&bin_final).ok();
465 return Err(StoreError::Io(e));
466 }
467 }
468 }
469 if let Ok(d) = fs::File::open(&dir) {
476 d.sync_all().ok();
477 }
478 self.metadata_index_upsert(&meta_final, &bin_final).ok();
479 let approx_added = payload.len() as u64 + APPROX_META_BYTES;
500 let new_total = self.byte_total.fetch_add(approx_added, Ordering::Relaxed) + approx_added;
501 let n = self.save_counter.fetch_add(1, Ordering::Relaxed);
502 if n == 0
503 || n.is_multiple_of(EVICT_EVERY_N_SAVES)
504 || new_total > self.opts.size_budget_bytes
505 {
506 self.evict_overflow().ok();
507 }
508 Ok(())
509 }
510
511 pub fn evict_overflow(&self) -> Result<(), StoreError> {
521 let current_root_mtime = fs::metadata(&self.root)
534 .ok()
535 .and_then(|m| m.modified().ok());
536 if self.byte_total.load(Ordering::Relaxed) <= self.opts.size_budget_bytes
537 && let Some(now_mtime) = current_root_mtime
538 && let Ok(last) = self.last_evict_root_mtime.lock()
539 && *last == Some(now_mtime)
540 {
541 return Ok(());
542 }
543 let read_dir = match fs::read_dir(&self.root) {
544 Ok(rd) => rd,
545 Err(_) => return Ok(()),
546 };
547 let mut all: Vec<(PathBuf, PathBuf, u64, u128, bool)> = Vec::new();
549 let now_nanos = self.nanos_now();
550 for key_dir_entry in read_dir {
551 let key_dir = match key_dir_entry {
552 Ok(e) => e.path(),
553 Err(_) => continue,
554 };
555 if !key_dir.is_dir() {
556 continue;
557 }
558 let scanned = self.scan_key_dir(&key_dir, now_nanos);
564 for entry in &scanned {
565 let write_nanos = (entry.meta.written_unix_secs as u128) * 1_000_000_000u128
566 + entry.meta.written_nanos as u128;
567 let total_bytes = entry.meta_len + entry.bin_len;
568 all.push((
569 entry.meta_path.clone(),
570 entry.bin_path.clone(),
571 total_bytes,
572 write_nanos,
573 entry.meta.accessed,
574 ));
575 }
576 if scanned.is_empty()
578 && fs::read_dir(&key_dir)
579 .map(|mut it| it.next().is_none())
580 .unwrap_or(false)
581 {
582 fs::remove_dir(&key_dir).ok();
583 if let Ok(mut index) = self.index.lock() {
584 index.by_key_dir.remove(&key_dir);
585 }
586 }
587 }
588 let total: u64 = all.iter().map(|e| e.2).sum();
589 if total <= self.opts.size_budget_bytes {
590 self.byte_total.store(total, Ordering::Relaxed);
597 if let (Ok(mut last), Some(m)) = (
604 self.last_evict_root_mtime.lock(),
605 fs::metadata(&self.root)
606 .ok()
607 .and_then(|m| m.modified().ok()),
608 ) {
609 *last = Some(m);
610 }
611 return Ok(());
612 }
613 all.sort_by(|a, b| {
614 a.4.cmp(&b.4)
615 .then_with(|| a.3.cmp(&b.3))
616 .then_with(|| a.0.cmp(&b.0))
617 });
618 let mut remaining = total;
619 for (meta, bin, bytes, _, _) in all.into_iter() {
620 if remaining <= self.opts.size_budget_bytes {
621 break;
622 }
623 fs::remove_file(&meta).ok();
624 fs::remove_file(&bin).ok();
625 self.metadata_index_remove(&meta);
626 remaining = remaining.saturating_sub(bytes);
627 }
628 self.byte_total.store(remaining, Ordering::Relaxed);
631 Ok(())
632 }
633}
634
635struct EntryWrite<'a> {
647 dir: &'a Path,
648 bin_tmp: &'a Path,
649 meta_tmp: &'a Path,
650 payload: &'a [u8],
651 bin_final: &'a Path,
652 meta_final: &'a Path,
653 stamp_fn: &'a dyn Fn() -> (u64, u32),
660 build_meta_json: &'a dyn Fn(u64, u32) -> io::Result<Vec<u8>>,
663}
664
665fn write_and_promote_entry(w: &EntryWrite<'_>) -> io::Result<()> {
666 fs::create_dir_all(w.dir)?;
670 {
671 let mut f = fs::File::create(w.bin_tmp)?;
672 f.write_all(w.payload)?;
673 f.sync_all().ok();
674 }
675 fs::rename(w.bin_tmp, w.bin_final)?;
679 let (secs, subsec_nanos) = (w.stamp_fn)();
686 let meta_json = (w.build_meta_json)(secs, subsec_nanos)?;
687 {
688 let mut f = fs::File::create(w.meta_tmp)?;
689 f.write_all(&meta_json)?;
690 f.sync_all().ok();
691 }
692 if let Err(e) = fs::rename(w.meta_tmp, w.meta_final) {
693 fs::remove_file(w.bin_final).ok();
696 return Err(e);
697 }
698 Ok(())
699}
700
701const APPROX_META_BYTES: u64 = 512;
705
706#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
708enum LookupMode {
709 Best,
711 Latest,
713}
714
715impl LookupMode {
716 fn better(&self, candidate: &OnDiskMeta, current: &OnDiskMeta) -> bool {
717 match self {
718 LookupMode::Best => entry_better(candidate, current),
719 LookupMode::Latest => entry_newer(candidate, current),
720 }
721 }
722}
723
724#[derive(Clone, Copy, PartialEq, Eq, Hash)]
725struct LookupCacheKey {
726 fp: Fingerprint,
727 mode: LookupMode,
728}
729
730#[derive(Clone)]
731struct CachedLookup {
732 meta_path: PathBuf,
733 meta_mtime: SystemTime,
734 write_nanos: u128,
738 entry: WarmStartEntry,
739}
740
741#[derive(Debug, Default)]
742struct MetadataIndex {
743 by_meta_path: HashMap<PathBuf, IndexedMeta>,
744 by_key_dir: HashMap<PathBuf, ScannedDir>,
755}
756
757#[derive(Debug, Clone)]
758struct IndexedMeta {
759 meta_mtime: SystemTime,
760 meta_len: u64,
761 bin_len: u64,
762 meta: OnDiskMeta,
763}
764
765impl IndexedMeta {
766 fn matches(&self, meta_md: &fs::Metadata, bin_md: &fs::Metadata) -> bool {
767 meta_md.modified().ok() == Some(self.meta_mtime)
768 && meta_md.len() == self.meta_len
769 && bin_md.len() == self.bin_len
770 }
771}
772
773#[derive(Debug, Clone)]
776struct ScannedDir {
777 dir_mtime: SystemTime,
778 entries: Vec<ScannedEntry>,
779}
780
781#[derive(Debug, Clone)]
785struct ScannedEntry {
786 meta_path: PathBuf,
787 bin_path: PathBuf,
788 meta_len: u64,
789 bin_len: u64,
790 meta_mtime: Option<SystemTime>,
791 bin_mtime: Option<SystemTime>,
792 meta: OnDiskMeta,
793}
794
795impl ScannedEntry {
796 fn matches_files(&self, meta_md: &fs::Metadata, bin_md: &fs::Metadata) -> bool {
797 meta_md.len() == self.meta_len
798 && bin_md.len() == self.bin_len
799 && meta_md.modified().ok() == self.meta_mtime
800 && bin_md.modified().ok() == self.bin_mtime
801 }
802}
803
804const fn meta_expired(activity_nanos: u128, ttl: Duration, now_nanos: u128) -> bool {
812 let ttl_nanos = ttl.as_nanos();
813 if ttl_nanos == 0 {
814 return false;
815 }
816 now_nanos.saturating_sub(activity_nanos) >= ttl_nanos
817}
818
819fn lookup_cache() -> &'static Mutex<HashMap<LookupCacheKey, CachedLookup>> {
826 static CACHE: OnceLock<Mutex<HashMap<LookupCacheKey, CachedLookup>>> = OnceLock::new();
827 CACHE.get_or_init(|| Mutex::new(HashMap::new()))
828}
829
830const LOOKUP_CACHE_MAX_ENTRIES: usize = 128;
831const LOOKUP_CACHE_MAX_BYTES: usize = 256 * 1024 * 1024;
832
833const fn cached_lookup_resident_bytes(value: &CachedLookup) -> usize {
834 std::mem::size_of::<CachedLookup>().saturating_add(value.entry.payload.capacity())
835}
836
837fn lookup_cache_get(key: &LookupCacheKey) -> Option<CachedLookup> {
838 let guard = lookup_cache().lock().ok()?;
839 guard.get(key).cloned()
840}
841
842fn lookup_cache_insert(key: LookupCacheKey, val: CachedLookup) {
843 if let Ok(mut guard) = lookup_cache().lock() {
844 let new_bytes = cached_lookup_resident_bytes(&val);
845 if new_bytes > LOOKUP_CACHE_MAX_BYTES {
846 return;
847 }
848 let mut resident_bytes: usize = guard.values().map(cached_lookup_resident_bytes).sum();
849 if let Some(old) = guard.remove(&key) {
850 resident_bytes = resident_bytes.saturating_sub(cached_lookup_resident_bytes(&old));
851 }
852 while guard.len() >= LOOKUP_CACHE_MAX_ENTRIES
853 || resident_bytes.saturating_add(new_bytes) > LOOKUP_CACHE_MAX_BYTES
854 {
855 let oldest = guard
856 .iter()
857 .min_by_key(|(_, cached)| cached.write_nanos)
858 .map(|(old_key, _)| *old_key);
859 let Some(oldest) = oldest else {
860 break;
861 };
862 if let Some(old) = guard.remove(&oldest) {
863 resident_bytes = resident_bytes.saturating_sub(cached_lookup_resident_bytes(&old));
864 }
865 }
866 guard.insert(key, val);
867 }
868}
869
870fn lookup_cache_invalidate(key: &LookupCacheKey) {
871 if let Ok(mut guard) = lookup_cache().lock() {
872 guard.remove(key);
873 }
874}
875
876const EVICT_EVERY_N_SAVES: u64 = 32;
881
882fn parse_tmp_pid(name: &str) -> Option<u32> {
883 let tail = name.split(".tmp.").nth(1)?;
887 let pid_str = tail.split('.').next()?;
888 pid_str.parse::<u32>().ok()
889}
890
891fn read_meta(path: &Path) -> Result<OnDiskMeta, StoreError> {
892 let bytes = fs::read(path)?;
893 let parsed: OnDiskMeta = serde_json::from_slice(&bytes)?;
894 Ok(parsed)
895}
896
897fn entry_better(candidate: &OnDiskMeta, current: &OnDiskMeta) -> bool {
898 match (candidate.objective, current.objective) {
899 (Some(c), Some(d)) => {
900 if (c - d).abs() < 1e-12 {
901 match (candidate.kind, current.kind) {
902 (EntryKind::Final, EntryKind::Checkpoint) => true,
903 (EntryKind::Checkpoint, EntryKind::Final) => false,
904 _ => entry_newer(candidate, current),
905 }
906 } else {
907 c < d
908 }
909 }
910 (Some(_), None) => true,
911 (None, Some(_)) => false,
912 (None, None) => entry_newer(candidate, current),
913 }
914}
915
916fn entry_newer(candidate: &OnDiskMeta, current: &OnDiskMeta) -> bool {
917 let candidate_stamp = (
918 candidate.written_unix_secs,
919 candidate.written_nanos,
920 candidate_kind_rank(candidate.kind),
921 );
922 let current_stamp = (
923 current.written_unix_secs,
924 current.written_nanos,
925 candidate_kind_rank(current.kind),
926 );
927 candidate_stamp > current_stamp
928}
929
930const fn candidate_kind_rank(kind: EntryKind) -> u8 {
931 match kind {
932 EntryKind::Checkpoint => 0,
933 EntryKind::Final => 1,
934 }
935}
936
937fn checksum_hex(payload: &[u8]) -> String {
938 let mut h = Sha256::new();
939 h.update(payload);
940 let out = h.finalize();
941 let mut s = String::with_capacity(out.len() * 2);
942 for b in out.iter() {
943 use std::fmt::Write;
944 write!(&mut s, "{:02x}", b).expect("writing to String is infallible");
945 }
946 s
947}
948
949impl WarmStartStore {
950 fn touch_lookup_hit(
951 &self,
952 meta_path: &Path,
953 entry: WarmStartEntry,
954 ) -> Result<WarmStartEntry, StoreError> {
955 let meta = read_meta(meta_path)?;
956 let (_meta, entry) = self.touch_lookup_meta(meta_path, meta, entry)?;
957 Ok(entry)
958 }
959
960 fn touch_lookup_meta(
961 &self,
962 meta_path: &Path,
963 mut meta: OnDiskMeta,
964 entry: WarmStartEntry,
965 ) -> Result<(OnDiskMeta, WarmStartEntry), StoreError> {
966 let now = self.nanos_now();
967 let old_access =
973 (meta.accessed_unix_secs as u128) * 1_000_000_000u128 + meta.accessed_nanos as u128;
974 let touched = now.max(old_access.saturating_add(1));
975 meta.accessed_unix_secs = (touched / 1_000_000_000u128) as u64;
976 meta.accessed_nanos = (touched % 1_000_000_000u128) as u32;
977 meta.accessed = true;
978 let json = serde_json::to_vec_pretty(&meta)?;
979 let tmp = meta_path.with_extension(format!(
980 "json.touch.tmp.{}.{}",
981 std::process::id(),
982 self.nanos_now()
983 ));
984 {
985 let mut f = fs::File::create(&tmp)?;
986 f.write_all(&json)?;
987 f.sync_all()?;
988 }
989 fs::rename(&tmp, meta_path)?;
990 if let Some(dir) = meta_path.parent()
991 && let Ok(d) = fs::File::open(dir)
992 {
993 d.sync_all().ok();
994 }
995 self.metadata_index_remove(meta_path);
996 Ok((meta, entry))
999 }
1000
1001 fn read_meta_indexed(
1002 &self,
1003 path: &Path,
1004 meta_md: &fs::Metadata,
1005 bin_md: &fs::Metadata,
1006 ) -> Result<OnDiskMeta, StoreError> {
1007 if let Ok(index) = self.index.lock()
1008 && let Some(cached) = index.by_meta_path.get(path)
1009 && cached.matches(meta_md, bin_md)
1010 {
1011 return Ok(cached.meta.clone());
1012 }
1013
1014 let meta = read_meta(path)?;
1015 let Some(meta_mtime) = meta_md.modified().ok() else {
1016 return Ok(meta);
1017 };
1018 if let Ok(mut index) = self.index.lock() {
1019 index.by_meta_path.insert(
1020 path.to_path_buf(),
1021 IndexedMeta {
1022 meta_mtime,
1023 meta_len: meta_md.len(),
1024 bin_len: bin_md.len(),
1025 meta: meta.clone(),
1026 },
1027 );
1028 }
1029 Ok(meta)
1030 }
1031
1032 fn metadata_index_upsert(&self, meta_path: &Path, bin_path: &Path) -> Result<(), StoreError> {
1033 if let Ok(mut index) = self.index.lock() {
1041 index.by_meta_path.remove(meta_path);
1042 if let Some(parent) = meta_path.parent() {
1043 index.by_key_dir.remove(parent);
1044 }
1045 }
1046 let meta_md = fs::metadata(meta_path)?;
1047 let bin_md = fs::metadata(bin_path)?;
1048 self.read_meta_indexed(meta_path, &meta_md, &bin_md)?;
1049 Ok(())
1050 }
1051
1052 fn metadata_index_remove(&self, meta_path: &Path) {
1053 if let Ok(mut index) = self.index.lock() {
1054 index.by_meta_path.remove(meta_path);
1055 if let Some(parent) = meta_path.parent() {
1056 index.by_key_dir.remove(parent);
1057 }
1058 }
1059 }
1060
1061 fn metadata_index_remove_key(&self, key: &Fingerprint) {
1062 let dir = self.key_dir(key);
1063 if let Ok(mut index) = self.index.lock() {
1064 index.by_meta_path.retain(|path, _| !path.starts_with(&dir));
1065 index.by_key_dir.remove(&dir);
1066 }
1067 }
1068
1069 fn cached_dir_scan(&self, dir: &Path, dir_md: &fs::Metadata) -> Option<Vec<ScannedEntry>> {
1080 let dir_mtime = dir_md.modified().ok()?;
1081 let index = self.index.lock().ok()?;
1082 let cached = index.by_key_dir.get(dir)?;
1083 if cached.dir_mtime != dir_mtime {
1084 return None;
1085 }
1086 for entry in &cached.entries {
1087 let meta_md = fs::metadata(&entry.meta_path).ok()?;
1088 let bin_md = fs::metadata(&entry.bin_path).ok()?;
1089 if !entry.matches_files(&meta_md, &bin_md) {
1090 return None;
1091 }
1092 }
1093 Some(cached.entries.clone())
1094 }
1095
1096 fn store_dir_scan(&self, dir: &Path, dir_mtime: SystemTime, entries: &[ScannedEntry]) {
1097 if let Ok(mut index) = self.index.lock() {
1098 index.by_key_dir.insert(
1099 dir.to_path_buf(),
1100 ScannedDir {
1101 dir_mtime,
1102 entries: entries.to_vec(),
1103 },
1104 );
1105 }
1106 }
1107
1108 fn scan_key_dir(&self, dir: &Path, now_nanos: u128) -> Vec<ScannedEntry> {
1122 let dir_md = match fs::metadata(dir) {
1123 Ok(m) => m,
1124 Err(_) => return Vec::new(),
1125 };
1126 if let Some(cached) = self.cached_dir_scan(dir, &dir_md) {
1127 let any_expired = cached
1134 .iter()
1135 .any(|e| meta_expired(meta_activity_nanos(&e.meta), self.opts.ttl, now_nanos));
1136 if !any_expired {
1137 return cached;
1138 }
1139 let mut survivors = Vec::with_capacity(cached.len());
1140 for entry in cached {
1141 if meta_expired(meta_activity_nanos(&entry.meta), self.opts.ttl, now_nanos) {
1142 fs::remove_file(&entry.meta_path).ok();
1143 fs::remove_file(&entry.bin_path).ok();
1144 self.metadata_index_remove(&entry.meta_path);
1145 } else {
1146 survivors.push(entry);
1147 }
1148 }
1149 if let Some(mtime) = fs::metadata(dir).ok().and_then(|m| m.modified().ok()) {
1150 self.store_dir_scan(dir, mtime, &survivors);
1151 }
1152 return survivors;
1153 }
1154 let read_dir = match fs::read_dir(dir) {
1155 Ok(rd) => rd,
1156 Err(_) => return Vec::new(),
1157 };
1158 let mut entries = Vec::new();
1159 let mut mutated = false;
1160 for f in read_dir {
1161 let path = match f {
1162 Ok(e) => e.path(),
1163 Err(_) => continue,
1164 };
1165 let name = match path.file_name().and_then(|s| s.to_str()) {
1166 Some(s) => s,
1167 None => continue,
1168 };
1169 if name.contains(".tmp.") {
1170 if let Some(pid) = parse_tmp_pid(name)
1171 && pid != std::process::id()
1172 {
1173 fs::remove_file(&path).ok();
1174 mutated = true;
1175 }
1176 continue;
1177 }
1178 if path.extension().and_then(|s| s.to_str()) != Some("json") {
1179 continue;
1180 }
1181 let meta_md = match fs::metadata(&path) {
1182 Ok(m) => m,
1183 Err(_) => continue,
1184 };
1185 let bin = path.with_extension("bin");
1186 let bin_md = match fs::metadata(&bin) {
1187 Ok(m) => m,
1188 Err(_) => {
1189 fs::remove_file(&path).ok();
1190 self.metadata_index_remove(&path);
1191 mutated = true;
1192 continue;
1193 }
1194 };
1195 let meta = match self.read_meta_indexed(&path, &meta_md, &bin_md) {
1196 Ok(m) => m,
1197 Err(_) => {
1198 fs::remove_file(&path).ok();
1199 fs::remove_file(&bin).ok();
1200 self.metadata_index_remove(&path);
1201 mutated = true;
1202 continue;
1203 }
1204 };
1205 if meta.schema_version != SCHEMA_VERSION {
1206 fs::remove_file(&path).ok();
1207 fs::remove_file(&bin).ok();
1208 self.metadata_index_remove(&path);
1209 mutated = true;
1210 continue;
1211 }
1212 if meta_expired(meta_activity_nanos(&meta), self.opts.ttl, now_nanos) {
1213 fs::remove_file(&path).ok();
1214 fs::remove_file(&bin).ok();
1215 self.metadata_index_remove(&path);
1216 mutated = true;
1217 continue;
1218 }
1219 entries.push(ScannedEntry {
1220 meta_path: path,
1221 bin_path: bin,
1222 meta_len: meta_md.len(),
1223 bin_len: bin_md.len(),
1224 meta_mtime: meta_md.modified().ok(),
1225 bin_mtime: bin_md.modified().ok(),
1226 meta,
1227 });
1228 }
1229 let final_mtime = if mutated {
1233 fs::metadata(dir).ok().and_then(|m| m.modified().ok())
1234 } else {
1235 dir_md.modified().ok()
1236 };
1237 if let Some(mtime) = final_mtime {
1238 self.store_dir_scan(dir, mtime, &entries);
1239 }
1240 entries
1241 }
1242
1243 fn test_time_offset_ns(&self) -> u64 {
1244 self.test_time_offset_ns.load(Ordering::Relaxed)
1245 }
1246
1247 fn unix_now_parts(&self) -> (u64, u32) {
1248 let base = SystemTime::now()
1249 .duration_since(UNIX_EPOCH)
1250 .map(|d| d.as_nanos())
1251 .unwrap_or(0);
1252 let total = base.saturating_add(u128::from(self.test_time_offset_ns()));
1253 let secs = (total / 1_000_000_000u128) as u64;
1254 let nanos = (total % 1_000_000_000u128) as u32;
1255 (secs, nanos)
1256 }
1257
1258 fn nanos_now(&self) -> u128 {
1259 let base = SystemTime::now()
1260 .duration_since(UNIX_EPOCH)
1261 .map(|d| d.as_nanos())
1262 .unwrap_or(0);
1263 base.saturating_add(u128::from(self.test_time_offset_ns()))
1264 }
1265
1266 fn fresh_run_id(&self) -> String {
1267 let pid = std::process::id();
1268 let nanos = self.nanos_now();
1269 format!("r{pid:x}-{nanos:x}")
1270 }
1271}
1272
1273#[cfg(test)]
1274mod tests {
1275 use super::*;
1276 use crate::warm_start::key::Fingerprinter;
1277
1278 impl WarmStartStore {
1279 fn test_advance_time(&self, dur: Duration) {
1283 self.test_time_offset_ns
1284 .fetch_add(dur.as_nanos() as u64, Ordering::Relaxed);
1285 }
1286 }
1287
1288 fn temp_store() -> (tempfile::TempDir, WarmStartStore) {
1289 let dir = tempfile::tempdir().unwrap();
1290 let store = WarmStartStore::open(
1291 dir.path().to_path_buf(),
1292 StoreOptions {
1293 size_budget_bytes: 1024 * 1024,
1294 ttl: Duration::from_secs(60),
1295 },
1296 )
1297 .unwrap();
1298 (dir, store)
1299 }
1300
1301 fn key_for(s: &str) -> Fingerprint {
1302 let mut fp = Fingerprinter::new();
1303 fp.absorb_str(b"test", s);
1304 fp.finalize()
1305 }
1306
1307 #[test]
1308 fn roundtrip_save_then_lookup() {
1309 let (_d, store) = temp_store();
1310 let key = key_for("roundtrip");
1311 store
1312 .save(
1313 &key,
1314 b"hello-warm",
1315 Some(1.5),
1316 Some(7),
1317 EntryKind::Checkpoint,
1318 )
1319 .unwrap();
1320 let got = store.lookup(&key).unwrap().unwrap();
1321 assert_eq!(got.payload, b"hello-warm");
1322 assert_eq!(got.objective, Some(1.5));
1323 assert_eq!(got.iteration, Some(7));
1324 assert_eq!(got.kind, EntryKind::Checkpoint);
1325 }
1326
1327 #[test]
1328 fn lookup_picks_lowest_objective() {
1329 let (_d, store) = temp_store();
1330 let key = key_for("multi");
1331 store
1332 .save(&key, b"worse", Some(3.0), Some(1), EntryKind::Checkpoint)
1333 .unwrap();
1334 store
1335 .save(&key, b"better", Some(1.0), Some(2), EntryKind::Checkpoint)
1336 .unwrap();
1337 store
1338 .save(&key, b"mid", Some(2.0), Some(3), EntryKind::Checkpoint)
1339 .unwrap();
1340 let got = store.lookup(&key).unwrap().unwrap();
1341 assert_eq!(got.payload, b"better");
1342 assert_eq!(got.objective, Some(1.0));
1343 }
1344
1345 #[test]
1346 fn lookup_latest_ignores_objective_ordering() {
1347 let (_d, store) = temp_store();
1348 let key = key_for("latest-vs-best");
1349 store
1350 .save(&key, b"low-objective", Some(1.0), Some(1), EntryKind::Final)
1351 .unwrap();
1352 store.test_advance_time(Duration::from_millis(2));
1353 store
1354 .save(
1355 &key,
1356 b"newer-higher-objective",
1357 Some(10.0),
1358 Some(2),
1359 EntryKind::Checkpoint,
1360 )
1361 .unwrap();
1362
1363 let best = store.lookup(&key).unwrap().unwrap();
1364 assert_eq!(best.payload, b"low-objective");
1365
1366 let latest = store.lookup_latest(&key).unwrap().unwrap();
1367 assert_eq!(latest.payload, b"newer-higher-objective");
1368 assert_eq!(latest.iteration, Some(2));
1369 }
1370
1371 #[test]
1372 fn tiebreak_final_beats_checkpoint() {
1373 let (_d, store) = temp_store();
1374 let key = key_for("tie");
1375 store
1376 .save(&key, b"ckpt", Some(1.0), None, EntryKind::Checkpoint)
1377 .unwrap();
1378 store
1380 .save(&key, b"final", Some(1.0), None, EntryKind::Final)
1381 .unwrap();
1382 let got = store.lookup(&key).unwrap().unwrap();
1383 assert_eq!(got.payload, b"final");
1384 assert_eq!(got.kind, EntryKind::Final);
1385 }
1386
1387 #[test]
1388 fn tiebreak_latest_mtime_when_no_objective() {
1389 let (_d, store) = temp_store();
1390 let key = key_for("latest");
1391 store
1392 .save(&key, b"first", None, None, EntryKind::Checkpoint)
1393 .unwrap();
1394 store.test_advance_time(Duration::from_millis(1_100));
1395 store
1396 .save(&key, b"second", None, None, EntryKind::Checkpoint)
1397 .unwrap();
1398 let got = store.lookup(&key).unwrap().unwrap();
1399 assert_eq!(got.payload, b"second");
1400 }
1401
1402 #[test]
1403 fn corrupt_payload_is_cleaned_up() {
1404 let (_d, store) = temp_store();
1405 let key = key_for("corrupt");
1406 store
1407 .save(&key, b"original", Some(0.0), None, EntryKind::Checkpoint)
1408 .unwrap();
1409 let dir = store.key_dir(&key);
1411 for entry in fs::read_dir(&dir).unwrap() {
1412 let p = entry.unwrap().path();
1413 if p.extension().and_then(|s| s.to_str()) == Some("bin") {
1414 fs::write(&p, b"tampered!").unwrap();
1415 }
1416 }
1417 let got = store.lookup(&key).unwrap();
1418 assert!(got.is_none(), "tampered entry must be rejected");
1419 let remaining: Vec<_> = fs::read_dir(&dir).unwrap().collect();
1421 assert!(remaining.is_empty(), "corrupt entry should be removed");
1422 }
1423
1424 #[test]
1425 fn corrupt_meta_json_is_cleaned_up() {
1426 let (_d, store) = temp_store();
1427 let key = key_for("badjson");
1428 store
1429 .save(&key, b"x", None, None, EntryKind::Checkpoint)
1430 .unwrap();
1431 let dir = store.key_dir(&key);
1432 for entry in fs::read_dir(&dir).unwrap() {
1433 let p = entry.unwrap().path();
1434 if p.extension().and_then(|s| s.to_str()) == Some("json") {
1435 fs::write(&p, b"{not valid json").unwrap();
1436 }
1437 }
1438 let got = store.lookup(&key).unwrap();
1439 assert!(got.is_none());
1440 }
1441
1442 #[test]
1443 fn schema_mismatched_entry_is_cleaned_up() {
1444 let (_d, store) = temp_store();
1445 let key = key_for("schema");
1446 store
1447 .save(&key, b"x", None, None, EntryKind::Checkpoint)
1448 .unwrap();
1449 let dir = store.key_dir(&key);
1450 for entry in fs::read_dir(&dir).unwrap() {
1451 let p = entry.unwrap().path();
1452 if p.extension().and_then(|s| s.to_str()) == Some("json") {
1453 let raw = fs::read(&p).unwrap();
1454 let mut parsed: serde_json::Value = serde_json::from_slice(&raw).unwrap();
1455 parsed["schema_version"] = serde_json::json!(SCHEMA_VERSION + 99);
1456 fs::write(&p, serde_json::to_vec_pretty(&parsed).unwrap()).unwrap();
1457 }
1458 }
1459 assert!(store.lookup(&key).unwrap().is_none());
1460 let remaining: Vec<_> = fs::read_dir(&dir).unwrap().collect();
1461 assert!(
1462 remaining.is_empty(),
1463 "schema-mismatched entry should be removed"
1464 );
1465 }
1466
1467 #[test]
1468 fn schema_mismatched_entry_is_removed_during_save_eviction_path() {
1469 let dir = tempfile::tempdir().unwrap();
1470 let store = WarmStartStore::open(
1471 dir.path().to_path_buf(),
1472 StoreOptions {
1473 size_budget_bytes: 6 * 1024,
1474 ttl: Duration::from_secs(3600),
1475 },
1476 )
1477 .unwrap();
1478 let stale_key = key_for("schema-size-stale");
1479 store
1480 .save(
1481 &stale_key,
1482 &vec![0u8; 4 * 1024],
1483 None,
1484 None,
1485 EntryKind::Checkpoint,
1486 )
1487 .unwrap();
1488
1489 let stale_dir = store.key_dir(&stale_key);
1490 let mut stale_meta = None;
1491 let mut stale_bin = None;
1492 for entry in fs::read_dir(&stale_dir).unwrap() {
1493 let p = entry.unwrap().path();
1494 match p.extension().and_then(|s| s.to_str()) {
1495 Some("json") => {
1496 let raw = fs::read(&p).unwrap();
1497 let mut parsed: serde_json::Value = serde_json::from_slice(&raw).unwrap();
1498 parsed["schema_version"] = serde_json::json!(SCHEMA_VERSION + 99);
1499 fs::write(&p, serde_json::to_vec_pretty(&parsed).unwrap()).unwrap();
1500 stale_meta = Some(p);
1501 }
1502 Some("bin") => stale_bin = Some(p),
1503 _ => {}
1504 }
1505 }
1506 let stale_meta = stale_meta.expect("saved entry should have metadata");
1507 let stale_bin = stale_bin.expect("saved entry should have payload");
1508
1509 let fresh_key = key_for("schema-size-fresh");
1510 store
1511 .save(
1512 &fresh_key,
1513 &vec![1u8; 2 * 1024],
1514 None,
1515 None,
1516 EntryKind::Checkpoint,
1517 )
1518 .unwrap();
1519
1520 assert!(
1521 !stale_meta.exists(),
1522 "schema-mismatched metadata should be removed during eviction scan"
1523 );
1524 assert!(
1525 !stale_bin.exists(),
1526 "schema-mismatched payload should be removed during eviction scan"
1527 );
1528
1529 let mut total = 0u64;
1530 for key_dir in fs::read_dir(store.root()).unwrap() {
1531 let key_dir = key_dir.unwrap().path();
1532 if key_dir.is_dir() {
1533 for entry in fs::read_dir(key_dir).unwrap() {
1534 total += fs::metadata(entry.unwrap().path()).unwrap().len();
1535 }
1536 }
1537 }
1538 assert!(
1539 total <= store.options().size_budget_bytes,
1540 "schema-mismatched bytes must not leak past size accounting (got {total})"
1541 );
1542 assert!(store.lookup(&stale_key).unwrap().is_none());
1543 assert!(store.lookup(&fresh_key).unwrap().is_some());
1544 }
1545
1546 #[test]
1547 fn missing_bin_treated_as_missing() {
1548 let (_d, store) = temp_store();
1549 let key = key_for("nobin");
1550 store
1551 .save(&key, b"x", None, None, EntryKind::Checkpoint)
1552 .unwrap();
1553 let dir = store.key_dir(&key);
1554 for entry in fs::read_dir(&dir).unwrap() {
1555 let p = entry.unwrap().path();
1556 if p.extension().and_then(|s| s.to_str()) == Some("bin") {
1557 fs::remove_file(&p).unwrap();
1558 }
1559 }
1560 assert!(store.lookup(&key).unwrap().is_none());
1561 }
1562
1563 #[test]
1564 fn missing_key_returns_none() {
1565 let (_d, store) = temp_store();
1566 let key = key_for("absent");
1567 assert!(store.lookup(&key).unwrap().is_none());
1568 }
1569
1570 #[test]
1571 fn lru_eviction_under_size_budget() {
1572 let dir = tempfile::tempdir().unwrap();
1573 let store = WarmStartStore::open(
1575 dir.path().to_path_buf(),
1576 StoreOptions {
1577 size_budget_bytes: 4 * 1024,
1578 ttl: Duration::from_secs(3600),
1579 },
1580 )
1581 .unwrap();
1582 let mut keys = Vec::new();
1583 for i in 0..20 {
1584 let mut fp = Fingerprinter::new();
1585 fp.absorb_u64(b"i", i);
1586 let key = fp.finalize();
1587 keys.push(key);
1588 let payload = vec![0u8; 256];
1589 store
1590 .save(&key, &payload, Some(i as f64), None, EntryKind::Checkpoint)
1591 .unwrap();
1592 }
1593 let mut total = 0u64;
1595 for kd in fs::read_dir(store.root()).unwrap() {
1596 let kd = kd.unwrap().path();
1597 if kd.is_dir() {
1598 for f in fs::read_dir(&kd).unwrap() {
1599 total += fs::metadata(f.unwrap().path()).unwrap().len();
1600 }
1601 }
1602 }
1603 assert!(
1604 total <= 8 * 1024,
1605 "eviction failed to bound size (got {total})"
1606 );
1607 assert!(store.lookup(&keys[0]).unwrap().is_none());
1609 assert!(store.lookup(keys.last().unwrap()).unwrap().is_some());
1610 }
1611
1612 #[test]
1613 fn ttl_drops_old_entries() {
1614 let dir = tempfile::tempdir().unwrap();
1622 let ttl = Duration::from_secs(60);
1623 let store = WarmStartStore::open(
1624 dir.path().to_path_buf(),
1625 StoreOptions {
1626 size_budget_bytes: 1024 * 1024,
1627 ttl,
1628 },
1629 )
1630 .unwrap();
1631 let key = key_for("ttl");
1632 store
1633 .save(&key, b"x", None, None, EntryKind::Checkpoint)
1634 .unwrap();
1635 assert!(store.lookup(&key).unwrap().is_some());
1636 store.test_advance_time(ttl + Duration::from_secs(5));
1637 let other = key_for("ttl-other");
1639 store
1640 .save(&other, b"y", None, None, EntryKind::Checkpoint)
1641 .unwrap();
1642 assert!(store.lookup(&key).unwrap().is_none());
1644 assert!(store.lookup(&other).unwrap().is_some());
1645 }
1646
1647 #[test]
1648 fn orphan_temp_files_from_dead_processes_are_swept() {
1649 let (_d, store) = temp_store();
1650 let key = key_for("tmp");
1651 let dir = store.key_dir(&key);
1652 fs::create_dir_all(&dir).unwrap();
1653 let orphan_other = dir.join("r0-0.json.tmp.1.0");
1655 let mine = dir.join(format!("r0-0.bin.tmp.{}.0", std::process::id()));
1656 fs::write(&orphan_other, b"orphan").unwrap();
1657 fs::write(&mine, b"mine").unwrap();
1658 store.evict_overflow().unwrap();
1659 assert!(!orphan_other.exists(), "other-PID tmp file should be swept");
1660 assert!(mine.exists(), "same-PID tmp file must be left alone");
1661 }
1662
1663 #[test]
1664 fn tmp_filenames_without_pid_are_skipped() {
1665 let (_d, store) = temp_store();
1667 let key = key_for("malformed");
1668 let dir = store.key_dir(&key);
1669 fs::create_dir_all(&dir).unwrap();
1670 let weird = dir.join("garbage.tmp.notapid.suffix");
1671 fs::write(&weird, b"x").unwrap();
1672 store.evict_overflow().unwrap();
1674 assert!(weird.exists());
1675 }
1676
1677 #[test]
1678 fn save_overwrite_keeps_single_entry() {
1679 let (_d, store) = temp_store();
1680 let key = key_for("overwrite");
1681 let id = store
1682 .save(&key, b"v1", Some(2.0), Some(1), EntryKind::Checkpoint)
1683 .unwrap();
1684 store
1685 .save_overwrite(&key, &id, b"v2", Some(1.0), Some(2), EntryKind::Checkpoint)
1686 .unwrap();
1687 let dir = store.key_dir(&key);
1689 let files: Vec<_> = fs::read_dir(&dir).unwrap().collect();
1690 assert_eq!(files.len(), 2, "overwrite should not create a new run-id");
1691 let got = store.lookup(&key).unwrap().unwrap();
1692 assert_eq!(got.payload, b"v2");
1693 assert_eq!(got.objective, Some(1.0));
1694 }
1695
1696 #[test]
1697 fn write_and_promote_recreates_dir_removed_before_write() {
1698 let (_d, store) = temp_store();
1702 let key = key_for("race-recreate");
1703 let dir = store.key_dir(&key);
1704 assert!(!dir.exists());
1707 let bin_tmp = dir.join("r0.bin.tmp.1.0.0");
1708 let meta_tmp = dir.join("r0.json.tmp.1.0.0");
1709 let bin_final = dir.join("r0.bin");
1710 let meta_final = dir.join("r0.json");
1711 let stamp_fn = || (0u64, 0u32);
1712 let build_meta_json = |_: u64, _: u32| -> io::Result<Vec<u8>> { Ok(b"{}".to_vec()) };
1713 write_and_promote_entry(&EntryWrite {
1714 dir: &dir,
1715 bin_tmp: &bin_tmp,
1716 meta_tmp: &meta_tmp,
1717 payload: b"payload",
1718 bin_final: &bin_final,
1719 meta_final: &meta_final,
1720 stamp_fn: &stamp_fn,
1721 build_meta_json: &build_meta_json,
1722 })
1723 .expect("promote into a missing dir must recreate it and succeed");
1724 assert!(bin_final.exists() && meta_final.exists());
1725 assert_eq!(fs::read(&bin_final).unwrap(), b"payload");
1726 }
1727
1728 #[test]
1729 fn save_survives_concurrent_eviction_removing_key_dir() {
1730 use std::sync::Arc;
1737 use std::sync::atomic::AtomicBool;
1738
1739 let dir = tempfile::tempdir().unwrap();
1740 let store = Arc::new(
1744 WarmStartStore::open(
1745 dir.path().to_path_buf(),
1746 StoreOptions {
1747 size_budget_bytes: 0,
1748 ttl: Duration::from_secs(60),
1749 },
1750 )
1751 .unwrap(),
1752 );
1753 let key = key_for("concurrent-evict");
1754 let stop = Arc::new(AtomicBool::new(false));
1755
1756 let evictor = {
1757 let store = Arc::clone(&store);
1758 let stop = Arc::clone(&stop);
1759 std::thread::spawn(move || {
1760 while !stop.load(Ordering::Relaxed) {
1761 store.evict_overflow().ok();
1762 }
1763 })
1764 };
1765
1766 let writers: Vec<_> = (0..4)
1767 .map(|w| {
1768 let store = Arc::clone(&store);
1769 std::thread::spawn(move || {
1770 for i in 0..200u32 {
1771 let payload = format!("w{w}-i{i}");
1772 store
1773 .save(
1774 &key,
1775 payload.as_bytes(),
1776 Some(i as f64),
1777 Some(i as u64),
1778 EntryKind::Checkpoint,
1779 )
1780 .expect("save must not fail with ENOENT under concurrent eviction");
1781 }
1782 })
1783 })
1784 .collect();
1785
1786 for h in writers {
1787 h.join().unwrap();
1788 }
1789 stop.store(true, Ordering::Relaxed);
1790 evictor.join().unwrap();
1791 }
1792
1793 #[test]
1794 fn keys_are_isolated() {
1795 let (_d, store) = temp_store();
1796 let a = key_for("a");
1797 let b = key_for("b");
1798 store
1799 .save(&a, b"AAA", Some(1.0), None, EntryKind::Final)
1800 .unwrap();
1801 store
1802 .save(&b, b"BBB", Some(1.0), None, EntryKind::Final)
1803 .unwrap();
1804 assert_eq!(store.lookup(&a).unwrap().unwrap().payload, b"AAA");
1805 assert_eq!(store.lookup(&b).unwrap().unwrap().payload, b"BBB");
1806 }
1807}