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 let meta_md = fs::metadata(meta_path)?;
1034 let bin_md = fs::metadata(bin_path)?;
1035 self.read_meta_indexed(meta_path, &meta_md, &bin_md)?;
1036 if let Some(parent) = meta_path.parent()
1039 && let Ok(mut index) = self.index.lock()
1040 {
1041 index.by_key_dir.remove(parent);
1042 }
1043 Ok(())
1044 }
1045
1046 fn metadata_index_remove(&self, meta_path: &Path) {
1047 if let Ok(mut index) = self.index.lock() {
1048 index.by_meta_path.remove(meta_path);
1049 if let Some(parent) = meta_path.parent() {
1050 index.by_key_dir.remove(parent);
1051 }
1052 }
1053 }
1054
1055 fn metadata_index_remove_key(&self, key: &Fingerprint) {
1056 let dir = self.key_dir(key);
1057 if let Ok(mut index) = self.index.lock() {
1058 index.by_meta_path.retain(|path, _| !path.starts_with(&dir));
1059 index.by_key_dir.remove(&dir);
1060 }
1061 }
1062
1063 fn cached_dir_scan(&self, dir: &Path, dir_md: &fs::Metadata) -> Option<Vec<ScannedEntry>> {
1074 let dir_mtime = dir_md.modified().ok()?;
1075 let index = self.index.lock().ok()?;
1076 let cached = index.by_key_dir.get(dir)?;
1077 if cached.dir_mtime != dir_mtime {
1078 return None;
1079 }
1080 for entry in &cached.entries {
1081 let meta_md = fs::metadata(&entry.meta_path).ok()?;
1082 let bin_md = fs::metadata(&entry.bin_path).ok()?;
1083 if !entry.matches_files(&meta_md, &bin_md) {
1084 return None;
1085 }
1086 }
1087 Some(cached.entries.clone())
1088 }
1089
1090 fn store_dir_scan(&self, dir: &Path, dir_mtime: SystemTime, entries: &[ScannedEntry]) {
1091 if let Ok(mut index) = self.index.lock() {
1092 index.by_key_dir.insert(
1093 dir.to_path_buf(),
1094 ScannedDir {
1095 dir_mtime,
1096 entries: entries.to_vec(),
1097 },
1098 );
1099 }
1100 }
1101
1102 fn scan_key_dir(&self, dir: &Path, now_nanos: u128) -> Vec<ScannedEntry> {
1116 let dir_md = match fs::metadata(dir) {
1117 Ok(m) => m,
1118 Err(_) => return Vec::new(),
1119 };
1120 if let Some(cached) = self.cached_dir_scan(dir, &dir_md) {
1121 let any_expired = cached
1128 .iter()
1129 .any(|e| meta_expired(meta_activity_nanos(&e.meta), self.opts.ttl, now_nanos));
1130 if !any_expired {
1131 return cached;
1132 }
1133 let mut survivors = Vec::with_capacity(cached.len());
1134 for entry in cached {
1135 if meta_expired(meta_activity_nanos(&entry.meta), self.opts.ttl, now_nanos) {
1136 fs::remove_file(&entry.meta_path).ok();
1137 fs::remove_file(&entry.bin_path).ok();
1138 self.metadata_index_remove(&entry.meta_path);
1139 } else {
1140 survivors.push(entry);
1141 }
1142 }
1143 if let Some(mtime) = fs::metadata(dir).ok().and_then(|m| m.modified().ok()) {
1144 self.store_dir_scan(dir, mtime, &survivors);
1145 }
1146 return survivors;
1147 }
1148 let read_dir = match fs::read_dir(dir) {
1149 Ok(rd) => rd,
1150 Err(_) => return Vec::new(),
1151 };
1152 let mut entries = Vec::new();
1153 let mut mutated = false;
1154 for f in read_dir {
1155 let path = match f {
1156 Ok(e) => e.path(),
1157 Err(_) => continue,
1158 };
1159 let name = match path.file_name().and_then(|s| s.to_str()) {
1160 Some(s) => s,
1161 None => continue,
1162 };
1163 if name.contains(".tmp.") {
1164 if let Some(pid) = parse_tmp_pid(name)
1165 && pid != std::process::id()
1166 {
1167 fs::remove_file(&path).ok();
1168 mutated = true;
1169 }
1170 continue;
1171 }
1172 if path.extension().and_then(|s| s.to_str()) != Some("json") {
1173 continue;
1174 }
1175 let meta_md = match fs::metadata(&path) {
1176 Ok(m) => m,
1177 Err(_) => continue,
1178 };
1179 let bin = path.with_extension("bin");
1180 let bin_md = match fs::metadata(&bin) {
1181 Ok(m) => m,
1182 Err(_) => {
1183 fs::remove_file(&path).ok();
1184 self.metadata_index_remove(&path);
1185 mutated = true;
1186 continue;
1187 }
1188 };
1189 let meta = match self.read_meta_indexed(&path, &meta_md, &bin_md) {
1190 Ok(m) => m,
1191 Err(_) => {
1192 fs::remove_file(&path).ok();
1193 fs::remove_file(&bin).ok();
1194 self.metadata_index_remove(&path);
1195 mutated = true;
1196 continue;
1197 }
1198 };
1199 if meta.schema_version != SCHEMA_VERSION {
1200 fs::remove_file(&path).ok();
1201 fs::remove_file(&bin).ok();
1202 self.metadata_index_remove(&path);
1203 mutated = true;
1204 continue;
1205 }
1206 if meta_expired(meta_activity_nanos(&meta), self.opts.ttl, now_nanos) {
1207 fs::remove_file(&path).ok();
1208 fs::remove_file(&bin).ok();
1209 self.metadata_index_remove(&path);
1210 mutated = true;
1211 continue;
1212 }
1213 entries.push(ScannedEntry {
1214 meta_path: path,
1215 bin_path: bin,
1216 meta_len: meta_md.len(),
1217 bin_len: bin_md.len(),
1218 meta_mtime: meta_md.modified().ok(),
1219 bin_mtime: bin_md.modified().ok(),
1220 meta,
1221 });
1222 }
1223 let final_mtime = if mutated {
1227 fs::metadata(dir).ok().and_then(|m| m.modified().ok())
1228 } else {
1229 dir_md.modified().ok()
1230 };
1231 if let Some(mtime) = final_mtime {
1232 self.store_dir_scan(dir, mtime, &entries);
1233 }
1234 entries
1235 }
1236
1237 fn test_time_offset_ns(&self) -> u64 {
1238 self.test_time_offset_ns.load(Ordering::Relaxed)
1239 }
1240
1241 fn unix_now_parts(&self) -> (u64, u32) {
1242 let base = SystemTime::now()
1243 .duration_since(UNIX_EPOCH)
1244 .map(|d| d.as_nanos())
1245 .unwrap_or(0);
1246 let total = base.saturating_add(u128::from(self.test_time_offset_ns()));
1247 let secs = (total / 1_000_000_000u128) as u64;
1248 let nanos = (total % 1_000_000_000u128) as u32;
1249 (secs, nanos)
1250 }
1251
1252 fn nanos_now(&self) -> u128 {
1253 let base = SystemTime::now()
1254 .duration_since(UNIX_EPOCH)
1255 .map(|d| d.as_nanos())
1256 .unwrap_or(0);
1257 base.saturating_add(u128::from(self.test_time_offset_ns()))
1258 }
1259
1260 fn fresh_run_id(&self) -> String {
1261 let pid = std::process::id();
1262 let nanos = self.nanos_now();
1263 format!("r{pid:x}-{nanos:x}")
1264 }
1265}
1266
1267#[cfg(test)]
1268mod tests {
1269 use super::*;
1270 use crate::warm_start::key::Fingerprinter;
1271
1272 impl WarmStartStore {
1273 fn test_advance_time(&self, dur: Duration) {
1277 self.test_time_offset_ns
1278 .fetch_add(dur.as_nanos() as u64, Ordering::Relaxed);
1279 }
1280 }
1281
1282 fn temp_store() -> (tempfile::TempDir, WarmStartStore) {
1283 let dir = tempfile::tempdir().unwrap();
1284 let store = WarmStartStore::open(
1285 dir.path().to_path_buf(),
1286 StoreOptions {
1287 size_budget_bytes: 1024 * 1024,
1288 ttl: Duration::from_secs(60),
1289 },
1290 )
1291 .unwrap();
1292 (dir, store)
1293 }
1294
1295 fn key_for(s: &str) -> Fingerprint {
1296 let mut fp = Fingerprinter::new();
1297 fp.absorb_str(b"test", s);
1298 fp.finalize()
1299 }
1300
1301 #[test]
1302 fn roundtrip_save_then_lookup() {
1303 let (_d, store) = temp_store();
1304 let key = key_for("roundtrip");
1305 store
1306 .save(
1307 &key,
1308 b"hello-warm",
1309 Some(1.5),
1310 Some(7),
1311 EntryKind::Checkpoint,
1312 )
1313 .unwrap();
1314 let got = store.lookup(&key).unwrap().unwrap();
1315 assert_eq!(got.payload, b"hello-warm");
1316 assert_eq!(got.objective, Some(1.5));
1317 assert_eq!(got.iteration, Some(7));
1318 assert_eq!(got.kind, EntryKind::Checkpoint);
1319 }
1320
1321 #[test]
1322 fn lookup_picks_lowest_objective() {
1323 let (_d, store) = temp_store();
1324 let key = key_for("multi");
1325 store
1326 .save(&key, b"worse", Some(3.0), Some(1), EntryKind::Checkpoint)
1327 .unwrap();
1328 store
1329 .save(&key, b"better", Some(1.0), Some(2), EntryKind::Checkpoint)
1330 .unwrap();
1331 store
1332 .save(&key, b"mid", Some(2.0), Some(3), EntryKind::Checkpoint)
1333 .unwrap();
1334 let got = store.lookup(&key).unwrap().unwrap();
1335 assert_eq!(got.payload, b"better");
1336 assert_eq!(got.objective, Some(1.0));
1337 }
1338
1339 #[test]
1340 fn lookup_latest_ignores_objective_ordering() {
1341 let (_d, store) = temp_store();
1342 let key = key_for("latest-vs-best");
1343 store
1344 .save(&key, b"low-objective", Some(1.0), Some(1), EntryKind::Final)
1345 .unwrap();
1346 store.test_advance_time(Duration::from_millis(2));
1347 store
1348 .save(
1349 &key,
1350 b"newer-higher-objective",
1351 Some(10.0),
1352 Some(2),
1353 EntryKind::Checkpoint,
1354 )
1355 .unwrap();
1356
1357 let best = store.lookup(&key).unwrap().unwrap();
1358 assert_eq!(best.payload, b"low-objective");
1359
1360 let latest = store.lookup_latest(&key).unwrap().unwrap();
1361 assert_eq!(latest.payload, b"newer-higher-objective");
1362 assert_eq!(latest.iteration, Some(2));
1363 }
1364
1365 #[test]
1366 fn tiebreak_final_beats_checkpoint() {
1367 let (_d, store) = temp_store();
1368 let key = key_for("tie");
1369 store
1370 .save(&key, b"ckpt", Some(1.0), None, EntryKind::Checkpoint)
1371 .unwrap();
1372 store
1374 .save(&key, b"final", Some(1.0), None, EntryKind::Final)
1375 .unwrap();
1376 let got = store.lookup(&key).unwrap().unwrap();
1377 assert_eq!(got.payload, b"final");
1378 assert_eq!(got.kind, EntryKind::Final);
1379 }
1380
1381 #[test]
1382 fn tiebreak_latest_mtime_when_no_objective() {
1383 let (_d, store) = temp_store();
1384 let key = key_for("latest");
1385 store
1386 .save(&key, b"first", None, None, EntryKind::Checkpoint)
1387 .unwrap();
1388 store.test_advance_time(Duration::from_millis(1_100));
1389 store
1390 .save(&key, b"second", None, None, EntryKind::Checkpoint)
1391 .unwrap();
1392 let got = store.lookup(&key).unwrap().unwrap();
1393 assert_eq!(got.payload, b"second");
1394 }
1395
1396 #[test]
1397 fn corrupt_payload_is_cleaned_up() {
1398 let (_d, store) = temp_store();
1399 let key = key_for("corrupt");
1400 store
1401 .save(&key, b"original", Some(0.0), None, EntryKind::Checkpoint)
1402 .unwrap();
1403 let dir = store.key_dir(&key);
1405 for entry in fs::read_dir(&dir).unwrap() {
1406 let p = entry.unwrap().path();
1407 if p.extension().and_then(|s| s.to_str()) == Some("bin") {
1408 fs::write(&p, b"tampered!").unwrap();
1409 }
1410 }
1411 let got = store.lookup(&key).unwrap();
1412 assert!(got.is_none(), "tampered entry must be rejected");
1413 let remaining: Vec<_> = fs::read_dir(&dir).unwrap().collect();
1415 assert!(remaining.is_empty(), "corrupt entry should be removed");
1416 }
1417
1418 #[test]
1419 fn corrupt_meta_json_is_cleaned_up() {
1420 let (_d, store) = temp_store();
1421 let key = key_for("badjson");
1422 store
1423 .save(&key, b"x", None, None, EntryKind::Checkpoint)
1424 .unwrap();
1425 let dir = store.key_dir(&key);
1426 for entry in fs::read_dir(&dir).unwrap() {
1427 let p = entry.unwrap().path();
1428 if p.extension().and_then(|s| s.to_str()) == Some("json") {
1429 fs::write(&p, b"{not valid json").unwrap();
1430 }
1431 }
1432 let got = store.lookup(&key).unwrap();
1433 assert!(got.is_none());
1434 }
1435
1436 #[test]
1437 fn schema_mismatched_entry_is_cleaned_up() {
1438 let (_d, store) = temp_store();
1439 let key = key_for("schema");
1440 store
1441 .save(&key, b"x", None, None, EntryKind::Checkpoint)
1442 .unwrap();
1443 let dir = store.key_dir(&key);
1444 for entry in fs::read_dir(&dir).unwrap() {
1445 let p = entry.unwrap().path();
1446 if p.extension().and_then(|s| s.to_str()) == Some("json") {
1447 let raw = fs::read(&p).unwrap();
1448 let mut parsed: serde_json::Value = serde_json::from_slice(&raw).unwrap();
1449 parsed["schema_version"] = serde_json::json!(SCHEMA_VERSION + 99);
1450 fs::write(&p, serde_json::to_vec_pretty(&parsed).unwrap()).unwrap();
1451 }
1452 }
1453 assert!(store.lookup(&key).unwrap().is_none());
1454 let remaining: Vec<_> = fs::read_dir(&dir).unwrap().collect();
1455 assert!(
1456 remaining.is_empty(),
1457 "schema-mismatched entry should be removed"
1458 );
1459 }
1460
1461 #[test]
1462 fn schema_mismatched_entry_is_removed_during_save_eviction_path() {
1463 let dir = tempfile::tempdir().unwrap();
1464 let store = WarmStartStore::open(
1465 dir.path().to_path_buf(),
1466 StoreOptions {
1467 size_budget_bytes: 6 * 1024,
1468 ttl: Duration::from_secs(3600),
1469 },
1470 )
1471 .unwrap();
1472 let stale_key = key_for("schema-size-stale");
1473 store
1474 .save(
1475 &stale_key,
1476 &vec![0u8; 4 * 1024],
1477 None,
1478 None,
1479 EntryKind::Checkpoint,
1480 )
1481 .unwrap();
1482
1483 let stale_dir = store.key_dir(&stale_key);
1484 let mut stale_meta = None;
1485 let mut stale_bin = None;
1486 for entry in fs::read_dir(&stale_dir).unwrap() {
1487 let p = entry.unwrap().path();
1488 match p.extension().and_then(|s| s.to_str()) {
1489 Some("json") => {
1490 let raw = fs::read(&p).unwrap();
1491 let mut parsed: serde_json::Value = serde_json::from_slice(&raw).unwrap();
1492 parsed["schema_version"] = serde_json::json!(SCHEMA_VERSION + 99);
1493 fs::write(&p, serde_json::to_vec_pretty(&parsed).unwrap()).unwrap();
1494 stale_meta = Some(p);
1495 }
1496 Some("bin") => stale_bin = Some(p),
1497 _ => {}
1498 }
1499 }
1500 let stale_meta = stale_meta.expect("saved entry should have metadata");
1501 let stale_bin = stale_bin.expect("saved entry should have payload");
1502
1503 let fresh_key = key_for("schema-size-fresh");
1504 store
1505 .save(
1506 &fresh_key,
1507 &vec![1u8; 2 * 1024],
1508 None,
1509 None,
1510 EntryKind::Checkpoint,
1511 )
1512 .unwrap();
1513
1514 assert!(
1515 !stale_meta.exists(),
1516 "schema-mismatched metadata should be removed during eviction scan"
1517 );
1518 assert!(
1519 !stale_bin.exists(),
1520 "schema-mismatched payload should be removed during eviction scan"
1521 );
1522
1523 let mut total = 0u64;
1524 for key_dir in fs::read_dir(store.root()).unwrap() {
1525 let key_dir = key_dir.unwrap().path();
1526 if key_dir.is_dir() {
1527 for entry in fs::read_dir(key_dir).unwrap() {
1528 total += fs::metadata(entry.unwrap().path()).unwrap().len();
1529 }
1530 }
1531 }
1532 assert!(
1533 total <= store.options().size_budget_bytes,
1534 "schema-mismatched bytes must not leak past size accounting (got {total})"
1535 );
1536 assert!(store.lookup(&stale_key).unwrap().is_none());
1537 assert!(store.lookup(&fresh_key).unwrap().is_some());
1538 }
1539
1540 #[test]
1541 fn missing_bin_treated_as_missing() {
1542 let (_d, store) = temp_store();
1543 let key = key_for("nobin");
1544 store
1545 .save(&key, b"x", None, None, EntryKind::Checkpoint)
1546 .unwrap();
1547 let dir = store.key_dir(&key);
1548 for entry in fs::read_dir(&dir).unwrap() {
1549 let p = entry.unwrap().path();
1550 if p.extension().and_then(|s| s.to_str()) == Some("bin") {
1551 fs::remove_file(&p).unwrap();
1552 }
1553 }
1554 assert!(store.lookup(&key).unwrap().is_none());
1555 }
1556
1557 #[test]
1558 fn missing_key_returns_none() {
1559 let (_d, store) = temp_store();
1560 let key = key_for("absent");
1561 assert!(store.lookup(&key).unwrap().is_none());
1562 }
1563
1564 #[test]
1565 fn lru_eviction_under_size_budget() {
1566 let dir = tempfile::tempdir().unwrap();
1567 let store = WarmStartStore::open(
1569 dir.path().to_path_buf(),
1570 StoreOptions {
1571 size_budget_bytes: 4 * 1024,
1572 ttl: Duration::from_secs(3600),
1573 },
1574 )
1575 .unwrap();
1576 let mut keys = Vec::new();
1577 for i in 0..20 {
1578 let mut fp = Fingerprinter::new();
1579 fp.absorb_u64(b"i", i);
1580 let key = fp.finalize();
1581 keys.push(key);
1582 let payload = vec![0u8; 256];
1583 store
1584 .save(&key, &payload, Some(i as f64), None, EntryKind::Checkpoint)
1585 .unwrap();
1586 }
1587 let mut total = 0u64;
1589 for kd in fs::read_dir(store.root()).unwrap() {
1590 let kd = kd.unwrap().path();
1591 if kd.is_dir() {
1592 for f in fs::read_dir(&kd).unwrap() {
1593 total += fs::metadata(f.unwrap().path()).unwrap().len();
1594 }
1595 }
1596 }
1597 assert!(
1598 total <= 8 * 1024,
1599 "eviction failed to bound size (got {total})"
1600 );
1601 assert!(store.lookup(&keys[0]).unwrap().is_none());
1603 assert!(store.lookup(keys.last().unwrap()).unwrap().is_some());
1604 }
1605
1606 #[test]
1607 fn ttl_drops_old_entries() {
1608 let dir = tempfile::tempdir().unwrap();
1616 let ttl = Duration::from_secs(60);
1617 let store = WarmStartStore::open(
1618 dir.path().to_path_buf(),
1619 StoreOptions {
1620 size_budget_bytes: 1024 * 1024,
1621 ttl,
1622 },
1623 )
1624 .unwrap();
1625 let key = key_for("ttl");
1626 store
1627 .save(&key, b"x", None, None, EntryKind::Checkpoint)
1628 .unwrap();
1629 assert!(store.lookup(&key).unwrap().is_some());
1630 store.test_advance_time(ttl + Duration::from_secs(5));
1631 let other = key_for("ttl-other");
1633 store
1634 .save(&other, b"y", None, None, EntryKind::Checkpoint)
1635 .unwrap();
1636 assert!(store.lookup(&key).unwrap().is_none());
1638 assert!(store.lookup(&other).unwrap().is_some());
1639 }
1640
1641 #[test]
1642 fn orphan_temp_files_from_dead_processes_are_swept() {
1643 let (_d, store) = temp_store();
1644 let key = key_for("tmp");
1645 let dir = store.key_dir(&key);
1646 fs::create_dir_all(&dir).unwrap();
1647 let orphan_other = dir.join("r0-0.json.tmp.1.0");
1649 let mine = dir.join(format!("r0-0.bin.tmp.{}.0", std::process::id()));
1650 fs::write(&orphan_other, b"orphan").unwrap();
1651 fs::write(&mine, b"mine").unwrap();
1652 store.evict_overflow().unwrap();
1653 assert!(!orphan_other.exists(), "other-PID tmp file should be swept");
1654 assert!(mine.exists(), "same-PID tmp file must be left alone");
1655 }
1656
1657 #[test]
1658 fn tmp_filenames_without_pid_are_skipped() {
1659 let (_d, store) = temp_store();
1661 let key = key_for("malformed");
1662 let dir = store.key_dir(&key);
1663 fs::create_dir_all(&dir).unwrap();
1664 let weird = dir.join("garbage.tmp.notapid.suffix");
1665 fs::write(&weird, b"x").unwrap();
1666 store.evict_overflow().unwrap();
1668 assert!(weird.exists());
1669 }
1670
1671 #[test]
1672 fn save_overwrite_keeps_single_entry() {
1673 let (_d, store) = temp_store();
1674 let key = key_for("overwrite");
1675 let id = store
1676 .save(&key, b"v1", Some(2.0), Some(1), EntryKind::Checkpoint)
1677 .unwrap();
1678 store
1679 .save_overwrite(&key, &id, b"v2", Some(1.0), Some(2), EntryKind::Checkpoint)
1680 .unwrap();
1681 let dir = store.key_dir(&key);
1683 let files: Vec<_> = fs::read_dir(&dir).unwrap().collect();
1684 assert_eq!(files.len(), 2, "overwrite should not create a new run-id");
1685 let got = store.lookup(&key).unwrap().unwrap();
1686 assert_eq!(got.payload, b"v2");
1687 assert_eq!(got.objective, Some(1.0));
1688 }
1689
1690 #[test]
1691 fn write_and_promote_recreates_dir_removed_before_write() {
1692 let (_d, store) = temp_store();
1696 let key = key_for("race-recreate");
1697 let dir = store.key_dir(&key);
1698 assert!(!dir.exists());
1701 let bin_tmp = dir.join("r0.bin.tmp.1.0.0");
1702 let meta_tmp = dir.join("r0.json.tmp.1.0.0");
1703 let bin_final = dir.join("r0.bin");
1704 let meta_final = dir.join("r0.json");
1705 let stamp_fn = || (0u64, 0u32);
1706 let build_meta_json = |_: u64, _: u32| -> io::Result<Vec<u8>> { Ok(b"{}".to_vec()) };
1707 write_and_promote_entry(&EntryWrite {
1708 dir: &dir,
1709 bin_tmp: &bin_tmp,
1710 meta_tmp: &meta_tmp,
1711 payload: b"payload",
1712 bin_final: &bin_final,
1713 meta_final: &meta_final,
1714 stamp_fn: &stamp_fn,
1715 build_meta_json: &build_meta_json,
1716 })
1717 .expect("promote into a missing dir must recreate it and succeed");
1718 assert!(bin_final.exists() && meta_final.exists());
1719 assert_eq!(fs::read(&bin_final).unwrap(), b"payload");
1720 }
1721
1722 #[test]
1723 fn save_survives_concurrent_eviction_removing_key_dir() {
1724 use std::sync::Arc;
1731 use std::sync::atomic::AtomicBool;
1732
1733 let dir = tempfile::tempdir().unwrap();
1734 let store = Arc::new(
1738 WarmStartStore::open(
1739 dir.path().to_path_buf(),
1740 StoreOptions {
1741 size_budget_bytes: 0,
1742 ttl: Duration::from_secs(60),
1743 },
1744 )
1745 .unwrap(),
1746 );
1747 let key = key_for("concurrent-evict");
1748 let stop = Arc::new(AtomicBool::new(false));
1749
1750 let evictor = {
1751 let store = Arc::clone(&store);
1752 let stop = Arc::clone(&stop);
1753 std::thread::spawn(move || {
1754 while !stop.load(Ordering::Relaxed) {
1755 store.evict_overflow().ok();
1756 }
1757 })
1758 };
1759
1760 let writers: Vec<_> = (0..4)
1761 .map(|w| {
1762 let store = Arc::clone(&store);
1763 std::thread::spawn(move || {
1764 for i in 0..200u32 {
1765 let payload = format!("w{w}-i{i}");
1766 store
1767 .save(
1768 &key,
1769 payload.as_bytes(),
1770 Some(i as f64),
1771 Some(i as u64),
1772 EntryKind::Checkpoint,
1773 )
1774 .expect("save must not fail with ENOENT under concurrent eviction");
1775 }
1776 })
1777 })
1778 .collect();
1779
1780 for h in writers {
1781 h.join().unwrap();
1782 }
1783 stop.store(true, Ordering::Relaxed);
1784 evictor.join().unwrap();
1785 }
1786
1787 #[test]
1788 fn keys_are_isolated() {
1789 let (_d, store) = temp_store();
1790 let a = key_for("a");
1791 let b = key_for("b");
1792 store
1793 .save(&a, b"AAA", Some(1.0), None, EntryKind::Final)
1794 .unwrap();
1795 store
1796 .save(&b, b"BBB", Some(1.0), None, EntryKind::Final)
1797 .unwrap();
1798 assert_eq!(store.lookup(&a).unwrap().unwrap().payload, b"AAA");
1799 assert_eq!(store.lookup(&b).unwrap().unwrap().payload, b"BBB");
1800 }
1801}