Skip to main content

subetha_cxc/
shared_universal.rs

1//! `SharedUniversal<T>` - Layer-2 cross-process container that
2//! migrates between Shared* backings as the workload shape changes.
3//!
4//! # The architectural claim
5//!
6//! A single cross-process container that auto-swaps its backing
7//! storage when the observed operation mix favors a different shape.
8//! At creation time the container starts in `Vec` mode (cheap pushes,
9//! O(N) `contains`). When `contains` calls dominate (the common case
10//! for membership / dedup workloads), the container migrates to
11//! `HashMap` mode (O(1) `contains`, slightly more expensive insert).
12//! Subsequent peer reads observe the migration via a version bump in
13//! the shared state header and transparently re-open the new backing.
14//!
15//! # The MVP scope
16//!
17//! - **2 backings only**: `SharedVec<T>` and `SharedHashMap<T, ()>`.
18//!   The extension to 5 backings (SharedRing, SharedHandleTable,
19//!   SharedBTreeMap, SharedTreiberStack) is its own bead.
20//! - **Single-writer model**: ONE process holds the writer role and
21//!   triggers migrations. Other processes are read-only observers
22//!   that follow the strategy tag. Multi-writer voting protocol is
23//!   ap-uvj.
24//! - **Local policy**: the writer's local op histogram drives
25//!   migration decisions. Quorum / cross-process voting is ap-uvj.
26//!
27//! # File layout
28//!
29//! Three coordinated files per logical container:
30//!
31//! ```text
32//! <base>.state.bin           always; small header MMF
33//! <base>-v{N}-vec.bin        current backing if strategy == Vec
34//! <base>-v{N}-map.bin        current backing if strategy == Map
35//! ```
36//!
37//! On migration: writer creates the new `-v{N+1}-{strategy}.bin`,
38//! copies the snapshot, then bumps `state.bin`'s version+strategy
39//! with a single CAS. Readers see the bump on their next op and
40//! re-open transparently.
41//!
42//! # Concurrency model
43//!
44//! - Reader / writer ops take an INTERNAL `RwLock<Backing<T>>` on
45//!   the handle (process-local; protects against the re-open race
46//!   between two ops in the same process).
47//! - Re-open is double-checked: re-read state.version under the
48//!   write lock; if some other thread already re-opened, drop the
49//!   write lock and use the current backing.
50//! - Migration is ONLY safe from a single writer process. If two
51//!   processes both try to migrate, both will succeed locally but
52//!   race on the state CAS; the loser's new backing file is
53//!   orphaned (cleanable). The voting protocol (ap-uvj) prevents
54//!   this; the MVP documents the single-writer constraint.
55
56use std::fs::OpenOptions;
57use std::marker::PhantomData;
58use std::mem::size_of;
59use std::path::{Path, PathBuf};
60use std::sync::atomic::{AtomicU64, Ordering};
61
62use memmap2::{MmapMut, MmapOptions};
63use parking_lot::RwLock;
64
65use crate::shared_hash_map::{MapError, SharedHashMap};
66use crate::shared_vec::SharedVec;
67
68pub const UNIVERSAL_MAGIC: u32 = 0x4150_5556;
69
70/// Strategy tag: which backing is currently live.
71#[repr(u8)]
72#[derive(Debug, Clone, Copy, PartialEq, Eq)]
73pub enum Strategy {
74    Vec = 0,
75    Map = 1,
76}
77
78impl Strategy {
79    fn from_u8(b: u8) -> Option<Self> {
80        match b {
81            0 => Some(Self::Vec),
82            1 => Some(Self::Map),
83            _ => None,
84        }
85    }
86
87    fn file_suffix(self) -> &'static str {
88        match self {
89            Self::Vec => "vec",
90            Self::Map => "map",
91        }
92    }
93}
94
95#[derive(Debug, Clone, Copy, PartialEq, Eq)]
96pub enum UniversalError {
97    InvalidStrategy,
98    IoError(std::io::ErrorKind),
99    LayoutMismatch,
100    VecError,
101    MapError(MapError),
102    Full,
103    /// `current_version + 1` overflows `u32`. After ~4 billion
104    /// migrations on the same base, the container refuses further
105    /// migrations rather than wrap version back to 0 (which then
106    /// silently overwrites the v=0 backing).
107    VersionExhausted,
108}
109
110impl From<std::io::Error> for UniversalError {
111    fn from(e: std::io::Error) -> Self { Self::IoError(e.kind()) }
112}
113impl From<MapError> for UniversalError {
114    fn from(e: MapError) -> Self {
115        match e {
116            MapError::Full => Self::Full,
117            other => Self::MapError(other),
118        }
119    }
120}
121
122#[repr(C, align(64))]
123pub struct UniversalHeader {
124    pub magic: u32,
125    pub capacity: u32,
126    /// Packed state, single AtomicU64 for atomic update / atomic
127    /// reader-load:
128    /// - bits 63..32 = `version: u32` (bumps per migration within
129    ///   the current generation; wraps to 0 at u32::MAX)
130    /// - bits 31..16 = `generation: u16` (bumps when version wraps;
131    ///   ensures a reused (generation, version) pair NEVER appears
132    ///   in the same lifetime, so readers comparing the full u64
133    ///   state always observe wrap-around and re-open)
134    /// - bits 15..0  = `strategy: u16` (low byte is the Strategy
135    ///   discriminant; high byte reserved for strategy variants)
136    ///
137    /// True exhaustion: generation u16 AND version u32 both at MAX
138    /// (= 2^48 = 281 trillion migrations). Returns VersionExhausted.
139    pub state: AtomicU64,
140    /// Bumped by every `insert`; consumed by the writer's local
141    /// policy to decide when to migrate.
142    pub insert_count: AtomicU64,
143    /// Bumped by every `contains`; same role as `insert_count`.
144    pub contains_count: AtomicU64,
145    _pad: [u8; 32],
146}
147
148const _: () = {
149    assert!(size_of::<UniversalHeader>() == 64);
150};
151
152#[inline]
153fn pack(version: u32, generation: u16, strategy: u8) -> u64 {
154    ((version as u64) << 32) | ((generation as u64) << 16) | (strategy as u64)
155}
156#[inline]
157fn unpack(v: u64) -> (u32, u16, u8) {
158    let version = (v >> 32) as u32;
159    let generation = ((v >> 16) & 0xFFFF) as u16;
160    let strategy = (v & 0xFF) as u8;
161    (version, generation, strategy)
162}
163
164enum Backing<T: Copy + Eq + 'static> {
165    Vec(SharedVec<T>),
166    Map(SharedHashMap<T, ()>),
167}
168
169pub struct SharedUniversal<T: Copy + Eq + 'static> {
170    base: PathBuf,
171    capacity: usize,
172    _state_file: std::fs::File,
173    state_mmap: MmapMut,
174    /// Holds `(version, generation, Backing)`. Used to detect when
175    /// the shared state's (version, generation) pair has changed and
176    /// the local backing handle needs to be re-opened.
177    backing: RwLock<(u32, u16, Backing<T>)>,
178    _phantom: PhantomData<T>,
179    header_sidecar: subetha_core::HandshakeHeader,
180    ring_sidecar: Box<subetha_core::ObservationRing>,
181}
182
183unsafe impl<T: Copy + Eq + Send + 'static> Send for SharedUniversal<T> {}
184unsafe impl<T: Copy + Eq + Sync + 'static> Sync for SharedUniversal<T> {}
185
186impl<T: Copy + Eq + Send + Sync + 'static> subetha_sidecar::AdaptiveInstance for SharedUniversal<T> {
187    fn header(&self) -> &subetha_core::HandshakeHeader { &self.header_sidecar }
188    fn ring(&self) -> &subetha_core::ObservationRing { &self.ring_sidecar }
189    fn make_policy(&self) -> Box<dyn subetha_sidecar::Policy> {
190        Box::new(subetha_sidecar::NoMigrationPolicy)
191    }
192}
193
194impl<T: Copy + Eq + 'static> SharedUniversal<T> {
195    fn state_path(base: &Path) -> PathBuf {
196        let mut p = base.to_path_buf();
197        let stem = p.file_name().map(|s| s.to_string_lossy().to_string()).unwrap_or_default();
198        p.set_file_name(format!("{stem}.state.bin"));
199        p
200    }
201
202    fn backing_path(base: &Path, generation: u16, version: u32, strategy: Strategy) -> PathBuf {
203        let mut p = base.to_path_buf();
204        let stem = p.file_name().map(|s| s.to_string_lossy().to_string()).unwrap_or_default();
205        p.set_file_name(format!(
206            "{stem}-g{generation}-v{version}-{}.bin",
207            strategy.file_suffix(),
208        ));
209        p
210    }
211
212    /// Obtain the container at `base`: initializes a fresh one in Vec
213    /// strategy at (generation=0, version=0) if the state file does
214    /// not yet exist, and otherwise attaches to the live state and
215    /// obtains whichever backing it names. A container built with a
216    /// different capacity is a `LayoutMismatch`.
217    /// [`reset`](Self::reset) reinitializes.
218    pub fn create(base: impl AsRef<Path>, capacity: usize) -> Result<Self, UniversalError> {
219        let base = base.as_ref().to_path_buf();
220        assert!(capacity >= 1);
221        let state_p = Self::state_path(&base);
222        let (state_file, mmap) = crate::mmf_attach::create_or_attach(
223            &state_p,
224            size_of::<UniversalHeader>(),
225            |ptr| unsafe { Self::init_state(ptr, capacity) },
226            |ptr| unsafe { (*(ptr as *const UniversalHeader)).magic == UNIVERSAL_MAGIC },
227        )?;
228        let hdr = unsafe { &*(mmap.as_ptr() as *const UniversalHeader) };
229        if hdr.magic != UNIVERSAL_MAGIC || hdr.capacity != capacity as u32 {
230            return Err(UniversalError::LayoutMismatch);
231        }
232        let (version, generation, strategy_byte) = unpack(hdr.state.load(Ordering::Acquire));
233        let strategy = Strategy::from_u8(strategy_byte).ok_or(UniversalError::InvalidStrategy)?;
234        let backing = Self::obtain_backing(&base, generation, version, strategy, capacity)?;
235        Ok(Self {
236            base, capacity,
237            _state_file: state_file,
238            state_mmap: mmap,
239            backing: RwLock::new((version, generation, backing)),
240            _phantom: PhantomData,
241            header_sidecar: subetha_core::HandshakeHeader::new(),
242            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
243        })
244    }
245
246    /// Truncate the state at `base` and initialize a fresh container
247    /// in Vec strategy at (generation=0, version=0), resetting the
248    /// (0, 0) Vec backing with it. Backing files from later
249    /// (generation, version) pairs stay on disk but are unreachable
250    /// once the state names (0, 0). For a caller that knows it owns
251    /// the base path.
252    pub fn reset(base: impl AsRef<Path>, capacity: usize) -> Result<Self, UniversalError> {
253        let base = base.as_ref().to_path_buf();
254        assert!(capacity >= 1);
255        let state_p = Self::state_path(&base);
256        let (state_file, mmap) = crate::mmf_attach::reset(
257            &state_p,
258            size_of::<UniversalHeader>(),
259            |ptr| unsafe { Self::init_state(ptr, capacity) },
260        )?;
261        let backing_p = Self::backing_path(&base, 0, 0, Strategy::Vec);
262        let vec: SharedVec<T> = SharedVec::reset(&backing_p, capacity)
263            .map_err(|_| UniversalError::VecError)?;
264        Ok(Self {
265            base, capacity,
266            _state_file: state_file,
267            state_mmap: mmap,
268            backing: RwLock::new((0, 0, Backing::Vec(vec))),
269            _phantom: PhantomData,
270            header_sidecar: subetha_core::HandshakeHeader::new(),
271            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
272        })
273    }
274
275    /// Lay out a fresh state header: capacity first, magic last,
276    /// because attachers spin on it. The zeroed state word is already
277    /// (version=0, generation=0, Strategy::Vec).
278    ///
279    /// # Safety
280    /// `ptr` addresses at least `size_of::<UniversalHeader>()`
281    /// writable zeroed bytes.
282    unsafe fn init_state(ptr: *mut u8, capacity: usize) {
283        let hdr = ptr as *mut UniversalHeader;
284        unsafe {
285            (*hdr).capacity = capacity as u32;
286            std::ptr::write_volatile(&raw mut (*hdr).magic, UNIVERSAL_MAGIC);
287        }
288    }
289
290    /// Obtain the backing at `(generation, version, strategy)`:
291    /// creates an empty one if its file does not yet exist and
292    /// attaches to it if it does.
293    fn obtain_backing(
294        base: &Path, generation: u16, version: u32, strategy: Strategy, capacity: usize,
295    ) -> Result<Backing<T>, UniversalError> {
296        let p = Self::backing_path(base, generation, version, strategy);
297        match strategy {
298            Strategy::Vec => {
299                let v: SharedVec<T> = SharedVec::create(&p, capacity)
300                    .map_err(|_| UniversalError::VecError)?;
301                Ok(Backing::Vec(v))
302            }
303            Strategy::Map => {
304                let m: SharedHashMap<T, ()> = SharedHashMap::create(&p, capacity)?;
305                Ok(Backing::Map(m))
306            }
307        }
308    }
309
310    /// Open an existing container. Reads the active (generation,
311    /// version, strategy) from the state header and opens the
312    /// matching backing file.
313    pub fn open(base: impl AsRef<Path>, capacity: usize) -> Result<Self, UniversalError> {
314        let base = base.as_ref().to_path_buf();
315        let state_p = Self::state_path(&base);
316        let state_file = OpenOptions::new().read(true).write(true).open(&state_p)?;
317        if state_file.metadata()?.len() < size_of::<UniversalHeader>() as u64 {
318            return Err(UniversalError::LayoutMismatch);
319        }
320        let mmap = unsafe { MmapOptions::new().len(size_of::<UniversalHeader>()).map_mut(&state_file)? };
321        let hdr = unsafe { &*(mmap.as_ptr() as *const UniversalHeader) };
322        if hdr.magic != UNIVERSAL_MAGIC || hdr.capacity != capacity as u32 {
323            return Err(UniversalError::LayoutMismatch);
324        }
325        let (version, generation, strategy_byte) = unpack(hdr.state.load(Ordering::Acquire));
326        let strategy = Strategy::from_u8(strategy_byte).ok_or(UniversalError::InvalidStrategy)?;
327        let backing = Self::open_backing(&base, generation, version, strategy, capacity)?;
328        Ok(Self {
329            base, capacity,
330            _state_file: state_file,
331            state_mmap: mmap,
332            backing: RwLock::new((version, generation, backing)),
333            _phantom: PhantomData,
334            header_sidecar: subetha_core::HandshakeHeader::new(),
335            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
336        })
337    }
338
339    fn open_backing(
340        base: &Path, generation: u16, version: u32, strategy: Strategy, capacity: usize,
341    ) -> Result<Backing<T>, UniversalError> {
342        let p = Self::backing_path(base, generation, version, strategy);
343        match strategy {
344            Strategy::Vec => {
345                let v: SharedVec<T> = SharedVec::open(&p, capacity)
346                    .map_err(|_| UniversalError::VecError)?;
347                Ok(Backing::Vec(v))
348            }
349            Strategy::Map => {
350                let m: SharedHashMap<T, ()> = SharedHashMap::open(&p, capacity)?;
351                Ok(Backing::Map(m))
352            }
353        }
354    }
355
356    fn header(&self) -> &UniversalHeader {
357        unsafe { &*(self.state_mmap.as_ptr() as *const UniversalHeader) }
358    }
359
360    /// The strategy currently active in shared state. May differ from
361    /// the locally-held backing if another writer just migrated; the
362    /// next op call will re-open transparently.
363    pub fn strategy(&self) -> Strategy {
364        let (_, _, s) = unpack(self.header().state.load(Ordering::Acquire));
365        Strategy::from_u8(s).expect("invalid strategy byte in state header")
366    }
367
368    /// The shared strategy version. Bumps on every migration; wraps
369    /// to 0 at u32::MAX with the generation counter incrementing.
370    pub fn strategy_version(&self) -> u32 {
371        unpack(self.header().state.load(Ordering::Acquire)).0
372    }
373
374    /// The shared generation counter. Bumps each time `version`
375    /// wraps from u32::MAX back to 0. Together with `version` it
376    /// forms the true monotonic migration counter.
377    pub fn strategy_generation(&self) -> u16 {
378        unpack(self.header().state.load(Ordering::Acquire)).1
379    }
380
381    /// Re-open the local backing handle if the shared state's
382    /// (version, generation) pair differs from the locally cached
383    /// pair. Comparing both fields means a wrap-around (same version
384    /// at a new generation) ALSO triggers re-open, preventing the
385    /// stale-reader race where a reused version points at new
386    /// content. Double-checked so concurrent ops don't trample each
387    /// other.
388    fn refresh_backing_if_stale(&self) -> Result<(), UniversalError> {
389        let (shared_v, shared_g, _) = unpack(self.header().state.load(Ordering::Acquire));
390        {
391            let g = self.backing.read();
392            if g.0 == shared_v && g.1 == shared_g { return Ok(()); }
393        }
394        let mut g = self.backing.write();
395        let (shared_v2, shared_g2, shared_s_byte2) =
396            unpack(self.header().state.load(Ordering::Acquire));
397        if g.0 == shared_v2 && g.1 == shared_g2 { return Ok(()); }
398        let strategy = Strategy::from_u8(shared_s_byte2).ok_or(UniversalError::InvalidStrategy)?;
399        let new_backing = Self::open_backing(
400            &self.base, shared_g2, shared_v2, strategy, self.capacity,
401        )?;
402        *g = (shared_v2, shared_g2, new_backing);
403        Ok(())
404    }
405
406    /// Insert `value`. For Vec strategy this is push_back; for Map
407    /// strategy this is insert(value, ()).
408    pub fn insert(&self, value: T) -> Result<(), UniversalError>
409    where T: std::hash::Hash,
410    {
411        self.refresh_backing_if_stale()?;
412        let g = self.backing.read();
413        let r: Result<(), UniversalError> = match &g.2 {
414            Backing::Vec(v) => {
415                v.push_back(value).map_err(|_| UniversalError::Full).map(|_| ())
416            }
417            Backing::Map(m) => {
418                m.insert(value, ()).map(|_| ()).map_err(Into::into)
419            }
420        };
421        self.header().insert_count.fetch_add(1, Ordering::Relaxed);
422        self.ring_sidecar.push_op(
423            crate::sidecar_ops::universal::OP_INSERT,
424            if r.is_err() { 1 } else { 0 },
425        );
426        r
427    }
428
429    /// Membership check. Bumps the contains counter so the local
430    /// policy can observe contains-heavy workloads.
431    pub fn contains(&self, value: &T) -> Result<bool, UniversalError>
432    where T: std::hash::Hash,
433    {
434        self.refresh_backing_if_stale()?;
435        let g = self.backing.read();
436        let hit = match &g.2 {
437            Backing::Vec(v) => v.snapshot().iter().any(|x| x == value),
438            Backing::Map(m) => m.contains_key(value),
439        };
440        self.header().contains_count.fetch_add(1, Ordering::Relaxed);
441        self.ring_sidecar.push_op(
442            crate::sidecar_ops::universal::OP_CONTAINS,
443            if hit { 0 } else { 2 },
444        );
445        Ok(hit)
446    }
447
448    /// Number of live entries.
449    pub fn len(&self) -> Result<usize, UniversalError> {
450        self.refresh_backing_if_stale()?;
451        let g = self.backing.read();
452        Ok(match &g.2 {
453            Backing::Vec(v) => v.len(),
454            Backing::Map(m) => m.len(),
455        })
456    }
457
458    pub fn is_empty(&self) -> Result<bool, UniversalError> {
459        Ok(self.len()? == 0)
460    }
461
462    /// Reset the universal to empty: clears whichever backing is
463    /// currently live (Vec or Map). Does not change the strategy.
464    /// Useful for steady-state benches that need to reset accumulated
465    /// state between iterations. Not thread-safe with concurrent
466    /// insert/remove from other threads.
467    pub fn clear(&self) -> Result<(), UniversalError> {
468        self.refresh_backing_if_stale()?;
469        let g = self.backing.read();
470        match &g.2 {
471            Backing::Vec(v) => v.clear(),
472            Backing::Map(m) => m.clear(),
473        }
474        Ok(())
475    }
476
477    /// Snapshot all live values into a `Vec<T>`. Best-effort under
478    /// concurrent writers.
479    pub fn snapshot(&self) -> Result<Vec<T>, UniversalError> {
480        self.refresh_backing_if_stale()?;
481        let g = self.backing.read();
482        Ok(match &g.2 {
483            Backing::Vec(v) => v.snapshot(),
484            Backing::Map(m) => m.snapshot().into_iter().map(|(k, _)| k).collect(),
485        })
486    }
487
488    /// Operation counts since creation. The writer's policy code
489    /// reads these to decide when to migrate.
490    pub fn op_histogram(&self) -> (u64, u64) {
491        let hdr = self.header();
492        (
493            hdr.insert_count.load(Ordering::Acquire),
494            hdr.contains_count.load(Ordering::Acquire),
495        )
496    }
497
498    /// Force a migration to `target`. Snapshots the current backing,
499    /// creates a new backing file at version+1, restores the snapshot,
500    /// then publishes the new (version, strategy) via Release CAS.
501    ///
502    /// # Concurrency
503    ///
504    /// **Single-writer ONLY.** Two processes calling `migrate_to`
505    /// concurrently will both build new backings and race on the CAS;
506    /// the loser orphans its backing file. Use ap-uvj's voting
507    /// protocol to coordinate when multiple writers are involved.
508    pub fn migrate_to(&self, target: Strategy) -> Result<(), UniversalError>
509    where T: std::hash::Hash,
510    {
511        let r = self.migrate_to_inner(target);
512        self.ring_sidecar.push_op(
513            crate::sidecar_ops::universal::OP_MIGRATE,
514            if r.is_err() { 1 } else { 0 },
515        );
516        r
517    }
518
519    fn migrate_to_inner(&self, target: Strategy) -> Result<(), UniversalError>
520    where T: std::hash::Hash,
521    {
522        let current_state = self.header().state.load(Ordering::Acquire);
523        let (current_v, current_g, current_s) = unpack(current_state);
524        let current = Strategy::from_u8(current_s).ok_or(UniversalError::InvalidStrategy)?;
525        if current == target { return Ok(()); }
526        // Bump version; on overflow, bump generation and reset
527        // version to 0. True exhaustion (both at max) returns
528        // VersionExhausted - that ceiling is 2^48 = 281 trillion
529        // migrations on the same base.
530        let (new_v, new_g) = match current_v.checked_add(1) {
531            Some(v) => (v, current_g),
532            None => {
533                let next_g = current_g.checked_add(1)
534                    .ok_or(UniversalError::VersionExhausted)?;
535                (0, next_g)
536            }
537        };
538        let snap = self.snapshot()?;
539        let mut g = self.backing.write();
540        let new_p = Self::backing_path(&self.base, new_g, new_v, target);
541        // Build the new backing inside a closure so any error path
542        // can clean up the partially-created file before returning.
543        let build_result: Result<Backing<T>, UniversalError> = (|| {
544            match target {
545                Strategy::Vec => {
546                    let v: SharedVec<T> = SharedVec::create(&new_p, self.capacity)
547                        .map_err(|_| UniversalError::VecError)?;
548                    for x in &snap {
549                        v.push_back(*x).map_err(|_| UniversalError::Full)?;
550                    }
551                    Ok(Backing::Vec(v))
552                }
553                Strategy::Map => {
554                    let m: SharedHashMap<T, ()> = SharedHashMap::create(&new_p, self.capacity)?;
555                    for x in &snap { m.insert(*x, ())?; }
556                    Ok(Backing::Map(m))
557                }
558            }
559        })();
560        let new_backing = match build_result {
561            Ok(b) => b,
562            Err(e) => {
563                std::fs::remove_file(&new_p).ok();
564                return Err(e);
565            }
566        };
567        let new_state = pack(new_v, new_g, target as u8);
568        match self.header().state.compare_exchange(
569            current_state, new_state, Ordering::AcqRel, Ordering::Acquire,
570        ) {
571            Ok(_) => {
572                *g = (new_v, new_g, new_backing);
573                Ok(())
574            }
575            Err(_) => {
576                drop(new_backing);
577                std::fs::remove_file(&new_p).ok();
578                Err(UniversalError::VecError)
579            }
580        }
581    }
582
583    /// Local-policy migration trigger. If the observed `contains` ops
584    /// outnumber `insert` ops by at least `contains_to_insert_ratio`,
585    /// AND total ops exceed `min_total_ops`, migrate Vec → Map. If
586    /// the inverse holds, migrate Map → Vec.
587    ///
588    /// Returns `Ok(Some(new_strategy))` if a migration happened,
589    /// `Ok(None)` if no policy threshold was crossed.
590    pub fn maybe_migrate_by_policy(
591        &self,
592        contains_to_insert_ratio: f64,
593        min_total_ops: u64,
594    ) -> Result<Option<Strategy>, UniversalError>
595    where T: std::hash::Hash,
596    {
597        let (ins, cnt) = self.op_histogram();
598        if ins + cnt < min_total_ops { return Ok(None); }
599        let current = self.strategy();
600        let ratio = cnt as f64 / (ins.max(1)) as f64;
601        let want = if ratio >= contains_to_insert_ratio { Strategy::Map } else { Strategy::Vec };
602        if want == current { return Ok(None); }
603        self.migrate_to(want)?;
604        Ok(Some(want))
605    }
606}
607
608#[cfg(test)]
609mod tests {
610    use super::*;
611
612    fn tmp_base(name: &str) -> PathBuf {
613        let mut p = std::env::temp_dir();
614        let pid = std::process::id();
615        p.push(format!("subetha-universal-{name}-{pid}"));
616        p
617    }
618
619    fn cleanup(base: &Path) {
620        let stem = base.file_name().unwrap().to_string_lossy().to_string();
621        let parent = base.parent().unwrap_or_else(|| Path::new(""));
622        if let Ok(entries) = std::fs::read_dir(parent) {
623            for e in entries.flatten() {
624                let name = e.file_name().to_string_lossy().to_string();
625                if name.starts_with(&stem) {
626                    std::fs::remove_file(e.path()).ok();
627                }
628            }
629        }
630    }
631
632    #[test]
633    fn create_starts_in_vec_strategy() {
634        let base = tmp_base("starts-vec");
635        let u: SharedUniversal<u64> = SharedUniversal::create(&base, 64).unwrap();
636        assert_eq!(u.strategy(), Strategy::Vec);
637        assert_eq!(u.strategy_version(), 0);
638        assert_eq!(u.len().unwrap(), 0);
639        drop(u);
640        cleanup(&base);
641    }
642
643    #[test]
644    fn insert_and_contains_round_trip_in_vec_mode() {
645        let base = tmp_base("vec-roundtrip");
646        let u: SharedUniversal<u64> = SharedUniversal::create(&base, 64).unwrap();
647        for k in 0..10u64 { u.insert(k).unwrap(); }
648        for k in 0..10u64 { assert!(u.contains(&k).unwrap()); }
649        assert!(!u.contains(&999).unwrap());
650        assert_eq!(u.len().unwrap(), 10);
651        drop(u);
652        cleanup(&base);
653    }
654
655    #[test]
656    fn explicit_migrate_to_map_preserves_contents() {
657        let base = tmp_base("explicit-map");
658        let u: SharedUniversal<u64> = SharedUniversal::create(&base, 64).unwrap();
659        for k in 0..10u64 { u.insert(k).unwrap(); }
660        u.migrate_to(Strategy::Map).unwrap();
661        assert_eq!(u.strategy(), Strategy::Map);
662        assert_eq!(u.strategy_version(), 1);
663        for k in 0..10u64 { assert!(u.contains(&k).unwrap()); }
664        assert_eq!(u.len().unwrap(), 10);
665        drop(u);
666        cleanup(&base);
667    }
668
669    #[test]
670    fn migrate_back_and_forth_preserves_contents() {
671        let base = tmp_base("round-trip");
672        let u: SharedUniversal<u64> = SharedUniversal::create(&base, 64).unwrap();
673        for k in 0..5u64 { u.insert(k).unwrap(); }
674        u.migrate_to(Strategy::Map).unwrap();
675        u.migrate_to(Strategy::Vec).unwrap();
676        u.migrate_to(Strategy::Map).unwrap();
677        assert_eq!(u.strategy(), Strategy::Map);
678        assert_eq!(u.strategy_version(), 3);
679        let mut snap = u.snapshot().unwrap();
680        snap.sort();
681        assert_eq!(snap, vec![0, 1, 2, 3, 4]);
682        drop(u);
683        cleanup(&base);
684    }
685
686    #[test]
687    fn migrate_to_same_strategy_is_noop() {
688        let base = tmp_base("same");
689        let u: SharedUniversal<u64> = SharedUniversal::create(&base, 16).unwrap();
690        u.migrate_to(Strategy::Vec).unwrap();
691        assert_eq!(u.strategy_version(), 0);
692        u.migrate_to(Strategy::Map).unwrap();
693        u.migrate_to(Strategy::Map).unwrap();
694        assert_eq!(u.strategy_version(), 1);
695        drop(u);
696        cleanup(&base);
697    }
698
699    #[test]
700    fn local_policy_migrates_to_map_under_contains_load() {
701        let base = tmp_base("policy-map");
702        let u: SharedUniversal<u64> = SharedUniversal::create(&base, 64).unwrap();
703        for k in 0..5u64 { u.insert(k).unwrap(); }
704        // 100 contains, 5 inserts → ratio = 20, well above 0.5
705        for _ in 0..100 { u.contains(&3).unwrap(); }
706        let migrated = u.maybe_migrate_by_policy(0.5, 100).unwrap();
707        assert_eq!(migrated, Some(Strategy::Map));
708        assert_eq!(u.strategy(), Strategy::Map);
709        drop(u);
710        cleanup(&base);
711    }
712
713    /// A second create attaches to the live state - including a
714    /// post-migration Map backing - with contents in place; reset is
715    /// what strips them.
716    #[test]
717    fn second_create_attaches_and_keeps_contents() {
718        let base = tmp_base("attach");
719        cleanup(&base);
720        let u: SharedUniversal<u64> = SharedUniversal::create(&base, 64).unwrap();
721        for k in 0..10u64 { u.insert(k).unwrap(); }
722        u.migrate_to(Strategy::Map).unwrap();
723
724        let u2: SharedUniversal<u64> = SharedUniversal::create(&base, 64).unwrap();
725        assert_eq!(u2.strategy(), Strategy::Map, "attach ignored the live strategy");
726        assert_eq!(u2.len().unwrap(), 10, "attach lost contents");
727        for k in 0..10u64 { assert!(u2.contains(&k).unwrap()); }
728        assert!(matches!(
729            SharedUniversal::<u64>::create(&base, 32),
730            Err(UniversalError::LayoutMismatch),
731        ));
732
733        // Windows refuses to truncate a mapped file, so every handle
734        // goes before the reset.
735        drop(u);
736        drop(u2);
737        let fresh: SharedUniversal<u64> = SharedUniversal::reset(&base, 64).unwrap();
738        assert_eq!(fresh.strategy(), Strategy::Vec, "reset kept the migrated strategy");
739        assert_eq!(fresh.len().unwrap(), 0, "reset kept contents");
740        drop(fresh);
741        cleanup(&base);
742    }
743
744    #[test]
745    fn local_policy_keeps_vec_under_insert_load() {
746        let base = tmp_base("policy-vec");
747        let u: SharedUniversal<u64> = SharedUniversal::create(&base, 256).unwrap();
748        for k in 0..200u64 { u.insert(k).unwrap(); }
749        for _ in 0..10 { u.contains(&3).unwrap(); }
750        let migrated = u.maybe_migrate_by_policy(0.5, 100).unwrap();
751        assert_eq!(migrated, None);
752        assert_eq!(u.strategy(), Strategy::Vec);
753        drop(u);
754        cleanup(&base);
755    }
756
757    #[test]
758    fn reader_handle_observes_migration_via_version_bump() {
759        // Writer process equivalent: SharedUniversal::create.
760        // Reader process equivalent: SharedUniversal::open against
761        // the same base. After writer migrates, reader's next op
762        // must transparently re-open the new backing.
763        let base = tmp_base("cross-handle");
764        let writer: SharedUniversal<u64> = SharedUniversal::create(&base, 32).unwrap();
765        let reader: SharedUniversal<u64> = SharedUniversal::open(&base, 32).unwrap();
766        for k in 0..5u64 { writer.insert(k).unwrap(); }
767        // Reader sees the inserts in Vec mode.
768        assert_eq!(reader.strategy(), Strategy::Vec);
769        assert_eq!(reader.len().unwrap(), 5);
770        // Writer migrates.
771        writer.migrate_to(Strategy::Map).unwrap();
772        // Reader's next op transparently re-opens.
773        assert_eq!(reader.strategy(), Strategy::Map);
774        assert_eq!(reader.strategy_version(), 1);
775        assert_eq!(reader.len().unwrap(), 5);
776        for k in 0..5u64 { assert!(reader.contains(&k).unwrap()); }
777        drop(writer);
778        drop(reader);
779        cleanup(&base);
780    }
781
782    #[test]
783    fn snapshot_preserves_through_migration() {
784        let base = tmp_base("snap");
785        let u: SharedUniversal<u32> = SharedUniversal::create(&base, 32).unwrap();
786        for k in [10u32, 5, 7, 1, 99] { u.insert(k).unwrap(); }
787        let mut pre = u.snapshot().unwrap();
788        pre.sort();
789        u.migrate_to(Strategy::Map).unwrap();
790        let mut post = u.snapshot().unwrap();
791        post.sort();
792        assert_eq!(pre, post);
793        drop(u);
794        cleanup(&base);
795    }
796
797    #[test]
798    fn op_histogram_tracks_real_ops() {
799        let base = tmp_base("hist");
800        let u: SharedUniversal<u64> = SharedUniversal::create(&base, 32).unwrap();
801        u.insert(1).unwrap();
802        u.insert(2).unwrap();
803        u.insert(3).unwrap();
804        u.contains(&2).unwrap();
805        u.contains(&2).unwrap();
806        let (ins, cnt) = u.op_histogram();
807        assert_eq!(ins, 3);
808        assert_eq!(cnt, 2);
809        drop(u);
810        cleanup(&base);
811    }
812
813    #[test]
814    fn pack_unpack_round_trips_all_three_fields() {
815        // version, generation, strategy round-trip exactly through
816        // the u64 state encoding.
817        let v: u32 = 0xDEAD_BEEF;
818        let g: u16 = 0xCAFE;
819        let s: u8 = Strategy::Map as u8;
820        let packed = pack(v, g, s);
821        let (rv, rg, rs) = unpack(packed);
822        assert_eq!(rv, v);
823        assert_eq!(rg, g);
824        assert_eq!(rs, s);
825    }
826
827    #[test]
828    fn starts_at_generation_zero() {
829        let base = tmp_base("gen-zero");
830        let u: SharedUniversal<u64> = SharedUniversal::create(&base, 16).unwrap();
831        assert_eq!(u.strategy_version(), 0);
832        assert_eq!(u.strategy_generation(), 0);
833        drop(u);
834        cleanup(&base);
835    }
836
837    #[test]
838    fn migration_within_generation_keeps_generation_zero() {
839        let base = tmp_base("same-gen");
840        let u: SharedUniversal<u64> = SharedUniversal::create(&base, 16).unwrap();
841        u.migrate_to(Strategy::Map).unwrap();
842        u.migrate_to(Strategy::Vec).unwrap();
843        u.migrate_to(Strategy::Map).unwrap();
844        assert_eq!(u.strategy_version(), 3);
845        assert_eq!(u.strategy_generation(), 0);
846        drop(u);
847        cleanup(&base);
848    }
849
850    #[test]
851    fn version_wrap_bumps_generation_and_resets_version() {
852        // Synthesize a near-wrap state by creating a backing file
853        // at (g=0, v=u32::MAX, Vec) on disk, pointing the state
854        // header there, and forcing the local backing to re-open
855        // at that synthetic state. Then migrate once and verify
856        // version wraps to 0 and generation bumps to 1.
857        let base = tmp_base("wrap-gen");
858        let u: SharedUniversal<u64> = SharedUniversal::create(&base, 16).unwrap();
859        let synth_p = SharedUniversal::<u64>::backing_path(
860            u.base.as_path(), 0, u32::MAX, Strategy::Vec,
861        );
862        // Pre-create the synthetic backing file so refresh_backing
863        // can open it.
864        let synth: SharedVec<u64> = SharedVec::create(&synth_p, 16).unwrap();
865        drop(synth);
866        u.header().state.store(
867            pack(u32::MAX, 0, Strategy::Vec as u8),
868            Ordering::Release,
869        );
870        u.refresh_backing_if_stale().unwrap();
871        // Now migrate: version wraps to 0; generation bumps to 1.
872        u.migrate_to(Strategy::Map).unwrap();
873        assert_eq!(u.strategy_version(), 0);
874        assert_eq!(u.strategy_generation(), 1);
875        assert_eq!(u.strategy(), Strategy::Map);
876        drop(u);
877        cleanup(&base);
878    }
879
880    #[test]
881    fn true_exhaustion_returns_version_exhausted() {
882        // generation = u16::MAX, version = u32::MAX → next migrate
883        // can't bump either; returns VersionExhausted.
884        let base = tmp_base("exhausted");
885        let u: SharedUniversal<u64> = SharedUniversal::create(&base, 16).unwrap();
886        u.header().state.store(
887            pack(u32::MAX, u16::MAX, Strategy::Vec as u8),
888            Ordering::Release,
889        );
890        u.refresh_backing_if_stale().unwrap_or(());
891        let r = u.migrate_to(Strategy::Map);
892        assert_eq!(r.err(), Some(UniversalError::VersionExhausted));
893        drop(u);
894        cleanup(&base);
895    }
896
897    #[test]
898    fn reader_re_opens_on_generation_change_even_at_same_version() {
899        // This is the load-bearing safety property: if a writer
900        // wraps version back to 0 (bumping generation), an old
901        // reader whose cached (v, g) is (0, 0) MUST re-open when
902        // the shared state changes to (0, 1, new_strategy).
903        let base = tmp_base("reader-gen");
904        let writer: SharedUniversal<u64> = SharedUniversal::create(&base, 16).unwrap();
905        let reader: SharedUniversal<u64> = SharedUniversal::open(&base, 16).unwrap();
906        // Both at (v=0, g=0). Synthesize a wrap by writing the
907        // post-wrap state directly + creating a matching backing.
908        writer.insert(11).unwrap();
909        writer.insert(22).unwrap();
910        // Now wrap: pretend writer just completed a migration that
911        // wrapped version to 0 and bumped generation to 1, with
912        // strategy Map. Create the new-gen backing file the same
913        // way migrate_to does, then publish state.
914        let new_p = SharedUniversal::<u64>::backing_path(
915            writer.base.as_path(), 1, 0, Strategy::Map,
916        );
917        let m: SharedHashMap<u64, ()> = SharedHashMap::create(&new_p, 16).unwrap();
918        m.insert(99u64, ()).unwrap();
919        drop(m);
920        writer.header().state.store(
921            pack(0, 1, Strategy::Map as u8),
922            Ordering::Release,
923        );
924        // Reader's local backing is at (v=0, g=0). Without the
925        // generation check, it sees "v=0 == 0, no re-open
926        // needed" and return stale results. With the generation
927        // check, refresh_backing_if_stale re-opens at (0, 1, Map).
928        assert!(reader.contains(&99u64).unwrap());
929        // Old keys are NOT in the new Map backing (it was created
930        // fresh with only 99).
931        assert!(!reader.contains(&11u64).unwrap());
932        assert_eq!(reader.strategy(), Strategy::Map);
933        assert_eq!(reader.strategy_version(), 0);
934        assert_eq!(reader.strategy_generation(), 1);
935        drop(writer);
936        drop(reader);
937        cleanup(&base);
938    }
939}