Skip to main content

subetha_cxc/
shared_vec.rs

1//! `SharedVec<T>` - cross-process bounded indexable sequence.
2//!
3//! Distinct from [`SharedRing`](crate::SharedRing) (FIFO; drain
4//! semantics): SharedVec is RANDOM-ACCESS, accumulates monotonically
5//! up to capacity, and supports `get(i)` for any prior index.
6//!
7//! # Layout
8//!
9//! Single MMF file:
10//!
11//! ```text
12//! +---------------------------+
13//! | VecHeader  (64B aligned)  |  magic, capacity, len, slot_size
14//! +---------------------------+
15//! | Slot[0]    (64B = cache)  |  version + payload[VEC_PAYLOAD_BYTES]
16//! | Slot[1]                   |
17//! | ...                       |
18//! | Slot[capacity - 1]        |
19//! +---------------------------+
20//! ```
21//!
22//! Each slot is its own SeqLock cell (same shape as SharedCell);
23//! per-slot writes never false-share because each is its own cache
24//! line.
25//!
26//! # Concurrency
27//!
28//! - `push_back`: atomic `len.fetch_add(1)` claims a slot index; if
29//!   it exceeds capacity, rollback with `fetch_sub(1)` and return
30//!   `Full`. On success, write the payload under the slot's SeqLock
31//!   (version bump odd → write → bump even).
32//! - `get(i)`: load `len` (Acquire). If `i >= len`, return None.
33//!   Otherwise SeqLock-read `slot[i]`: spin if version is odd
34//!   (writer in progress), reread on version change.
35//! - `pop_back`: `compare_exchange` on `len` to decrement; if
36//!   successful, read the now-popped slot's payload at the old
37//!   index. The slot bytes remain in place but are no longer
38//!   addressable via `len`-bounded access.
39//! - `set(i, v)`: bounds-check against `len`, then SeqLock-write.
40//! - `clear`: store `len = 0` (Release). Previously-pushed slot
41//!   payloads remain on disk but become unreachable through the
42//!   bounded indexing.
43//!
44//! # Capacity
45//!
46//! Fixed at create time. The MMF is pre-allocated to the full size;
47//! no resize-on-grow protocol. The unbounded variant (with
48//! coordinator-mediated MMF resize) is a separate primitive.
49
50use std::fs::{File, OpenOptions};
51use std::marker::PhantomData;
52use std::mem::size_of;
53use std::path::Path;
54use std::sync::atomic::{AtomicU32, AtomicU64, Ordering};
55
56use memmap2::{MmapMut, MmapOptions};
57
58pub const VEC_MAGIC: u32 = 0x4150_5656;
59pub const VEC_PAYLOAD_BYTES: usize = 52;
60
61#[repr(C, align(64))]
62pub struct VecHeader {
63    pub magic: u32,
64    pub slot_payload_size: u32,
65    pub capacity: u64,
66    pub len: AtomicU64,
67    _pad: [u8; 40],
68}
69
70#[repr(C, align(64))]
71pub struct VecSlot {
72    pub version: AtomicU32,
73    _pad: [u8; 4],
74    pub payload: [u8; VEC_PAYLOAD_BYTES],
75}
76
77const _: () = {
78    assert!(size_of::<VecHeader>() == 64);
79    assert!(size_of::<VecSlot>() == 64);
80};
81
82pub const fn vec_file_size(capacity: usize) -> usize {
83    size_of::<VecHeader>() + capacity * size_of::<VecSlot>()
84}
85
86#[derive(Debug, Clone, Copy, PartialEq, Eq)]
87pub enum VecError {
88    Full,
89    OutOfBounds,
90    LayoutMismatch,
91    PayloadTooLarge,
92    IoError(std::io::ErrorKind),
93}
94
95impl From<std::io::Error> for VecError {
96    fn from(e: std::io::Error) -> Self { Self::IoError(e.kind()) }
97}
98
99pub struct SharedVec<T: Copy + 'static> {
100    _file: File,
101    mmap: MmapMut,
102    capacity: usize,
103    _phantom: PhantomData<T>,
104    header_sidecar: subetha_core::HandshakeHeader,
105    ring_sidecar: Box<subetha_core::ObservationRing>,
106}
107
108unsafe impl<T: Copy + Send + 'static> Send for SharedVec<T> {}
109unsafe impl<T: Copy + Sync + 'static> Sync for SharedVec<T> {}
110
111impl<T: Copy + Send + Sync + 'static> subetha_sidecar::AdaptiveInstance for SharedVec<T> {
112    fn header(&self) -> &subetha_core::HandshakeHeader { &self.header_sidecar }
113    fn ring(&self) -> &subetha_core::ObservationRing { &self.ring_sidecar }
114    fn make_policy(&self) -> Box<dyn subetha_sidecar::Policy> {
115        Box::new(subetha_sidecar::NoMigrationPolicy)
116    }
117}
118
119impl<T: Copy + 'static> SharedVec<T> {
120    pub fn create(
121        path: impl AsRef<Path>, capacity: usize,
122    ) -> Result<Self, VecError> {
123        if size_of::<T>() > VEC_PAYLOAD_BYTES {
124            return Err(VecError::PayloadTooLarge);
125        }
126        assert!(capacity >= 1);
127        let total = vec_file_size(capacity);
128        let file = OpenOptions::new()
129            .read(true).write(true).create(true).truncate(true)
130            .open(path.as_ref())?;
131        file.set_len(total as u64)?;
132        let mut mmap = unsafe { MmapOptions::new().len(total).map_mut(&file)? };
133        let hdr = mmap.as_mut_ptr() as *mut VecHeader;
134        unsafe {
135            std::ptr::write(hdr, VecHeader {
136                magic: VEC_MAGIC,
137                slot_payload_size: VEC_PAYLOAD_BYTES as u32,
138                capacity: capacity as u64,
139                len: AtomicU64::new(0),
140                _pad: [0; 40],
141            });
142        }
143        for i in 0..capacity {
144            let slot_ptr = unsafe {
145                mmap.as_mut_ptr()
146                    .add(size_of::<VecHeader>())
147                    .add(i * size_of::<VecSlot>())
148            } as *mut VecSlot;
149            unsafe {
150                std::ptr::write(slot_ptr, VecSlot {
151                    version: AtomicU32::new(0),
152                    _pad: [0; 4],
153                    payload: [0u8; VEC_PAYLOAD_BYTES],
154                });
155            }
156        }
157        Ok(Self {
158            _file: file, mmap, capacity, _phantom: PhantomData,
159            header_sidecar: subetha_core::HandshakeHeader::new(),
160            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
161        })
162    }
163
164    pub fn open(
165        path: impl AsRef<Path>, expected_capacity: usize,
166    ) -> Result<Self, VecError> {
167        if size_of::<T>() > VEC_PAYLOAD_BYTES {
168            return Err(VecError::PayloadTooLarge);
169        }
170        let file = OpenOptions::new().read(true).write(true).open(path.as_ref())?;
171        let total = vec_file_size(expected_capacity);
172        if file.metadata()?.len() < total as u64 {
173            return Err(VecError::LayoutMismatch);
174        }
175        let mmap = unsafe { MmapOptions::new().len(total).map_mut(&file)? };
176        let hdr = unsafe { &*(mmap.as_ptr() as *const VecHeader) };
177        if hdr.magic != VEC_MAGIC || hdr.capacity != expected_capacity as u64 {
178            return Err(VecError::LayoutMismatch);
179        }
180        Ok(Self {
181            _file: file, mmap, capacity: expected_capacity, _phantom: PhantomData,
182            header_sidecar: subetha_core::HandshakeHeader::new(),
183            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
184        })
185    }
186
187    #[inline]
188    pub fn capacity(&self) -> usize { self.capacity }
189
190    #[inline]
191    pub fn len(&self) -> usize {
192        self.header().len.load(Ordering::Acquire) as usize
193    }
194
195    #[inline]
196    pub fn is_empty(&self) -> bool { self.len() == 0 }
197
198    fn header(&self) -> &VecHeader {
199        unsafe { &*(self.mmap.as_ptr() as *const VecHeader) }
200    }
201
202    fn slot(&self, i: usize) -> &VecSlot {
203        assert!(i < self.capacity, "slot index {i} out of bounds for cap {}", self.capacity);
204        let base = unsafe { self.mmap.as_ptr().add(size_of::<VecHeader>()) };
205        unsafe { &*(base.add(i * size_of::<VecSlot>()) as *const VecSlot) }
206    }
207
208    /// SeqLock write of a payload into a slot. Caller is responsible
209    /// for ensuring `i < capacity`.
210    fn write_slot(&self, i: usize, v: T) {
211        let slot = self.slot(i);
212        // Bump version odd; subsequent readers will spin.
213        slot.version.fetch_add(1, Ordering::AcqRel);
214        // Memcpy the value.
215        let dst = unsafe {
216            let base = self.mmap.as_ptr().add(size_of::<VecHeader>())
217                .add(i * size_of::<VecSlot>())
218                .add(std::mem::offset_of!(VecSlot, payload));
219            base as *mut u8
220        };
221        unsafe {
222            std::ptr::copy_nonoverlapping(
223                &v as *const T as *const u8,
224                dst,
225                size_of::<T>(),
226            );
227        }
228        // Bump version even; readers resume.
229        slot.version.fetch_add(1, Ordering::AcqRel);
230    }
231
232    /// SeqLock read of a slot. Spins if version is odd; rereads on
233    /// version change.
234    fn read_slot(&self, i: usize) -> T {
235        let slot = self.slot(i);
236        loop {
237            let v1 = slot.version.load(Ordering::Acquire);
238            if v1 & 1 != 0 {
239                std::hint::spin_loop();
240                continue;
241            }
242            let mut out = std::mem::MaybeUninit::<T>::uninit();
243            let src = unsafe {
244                self.mmap.as_ptr().add(size_of::<VecHeader>())
245                    .add(i * size_of::<VecSlot>())
246                    .add(std::mem::offset_of!(VecSlot, payload))
247            };
248            unsafe {
249                std::ptr::copy_nonoverlapping(
250                    src, out.as_mut_ptr() as *mut u8, size_of::<T>(),
251                );
252            }
253            let v2 = slot.version.load(Ordering::Acquire);
254            if v1 == v2 {
255                return unsafe { out.assume_init() };
256            }
257        }
258    }
259
260    /// Append a value. Returns the index it landed at.
261    /// Returns `Err(Full)` when the vec is at capacity.
262    pub fn push_back(&self, v: T) -> Result<usize, VecError> {
263        let idx = self.header().len.fetch_add(1, Ordering::AcqRel) as usize;
264        if idx >= self.capacity {
265            self.header().len.fetch_sub(1, Ordering::AcqRel);
266            self.ring_sidecar
267                .push_op(crate::sidecar_ops::ordered::OP_INSERT, 1); // full
268            return Err(VecError::Full);
269        }
270        self.write_slot(idx, v);
271        self.ring_sidecar
272            .push_op(crate::sidecar_ops::ordered::OP_INSERT, 0);
273        Ok(idx)
274    }
275
276    /// Remove and return the last element. Returns None when empty.
277    pub fn pop_back(&self) -> Option<T> {
278        loop {
279            let cur = self.header().len.load(Ordering::Acquire);
280            if cur == 0 {
281                self.ring_sidecar
282                    .push_op(crate::sidecar_ops::ordered::OP_POP, 2); // empty
283                return None;
284            }
285            let new = cur - 1;
286            if self.header().len.compare_exchange(
287                cur, new, Ordering::AcqRel, Ordering::Acquire,
288            ).is_ok() {
289                let v = self.read_slot(new as usize);
290                self.ring_sidecar
291                    .push_op(crate::sidecar_ops::ordered::OP_POP, 0);
292                return Some(v);
293            }
294        }
295    }
296
297    /// Read the value at index `i`. Returns None when `i >= len`.
298    pub fn get(&self, i: usize) -> Option<T> {
299        if i >= self.len() {
300            self.ring_sidecar
301                .push_op(crate::sidecar_ops::ordered::OP_GET, 2); // out of bounds / absent
302            return None;
303        }
304        let v = self.read_slot(i);
305        self.ring_sidecar
306            .push_op(crate::sidecar_ops::ordered::OP_GET, 0);
307        Some(v)
308    }
309
310    /// Overwrite the value at index `i`. Returns `Err(OutOfBounds)`
311    /// when `i >= len`.
312    pub fn set(&self, i: usize, v: T) -> Result<(), VecError> {
313        if i >= self.len() {
314            self.ring_sidecar
315                .push_op(crate::sidecar_ops::ordered::OP_INSERT, 1); // out of bounds (positional write rejected)
316            return Err(VecError::OutOfBounds);
317        }
318        self.write_slot(i, v);
319        self.ring_sidecar
320            .push_op(crate::sidecar_ops::ordered::OP_INSERT, 0);
321        Ok(())
322    }
323
324    /// Clear the vec by resetting len to 0. Slot payloads are not
325    /// zeroed; they become unreachable through bounded indexing.
326    pub fn clear(&self) {
327        self.header().len.store(0, Ordering::Release);
328        self.ring_sidecar
329            .push_op(crate::sidecar_ops::ordered::OP_REMOVE, 0);
330    }
331
332    /// Snapshot all current values into a Vec. Best-effort: under
333    /// concurrent writers, the snapshot is a consistent prefix at
334    /// the moment of the `len` load, with each slot read under its
335    /// own SeqLock.
336    pub fn snapshot(&self) -> Vec<T> {
337        let n = self.len();
338        let mut out = Vec::with_capacity(n);
339        for i in 0..n {
340            out.push(self.read_slot(i));
341        }
342        self.ring_sidecar
343            .push_op(crate::sidecar_ops::ordered::OP_ITER, 0);
344        out
345    }
346
347    pub fn flush(&self) -> Result<(), VecError> {
348        self.mmap.flush()?;
349        Ok(())
350    }
351
352    /// Non-blocking flush: schedules a writeback via the OS.
353    /// Note: Windows is only partially async (sync to page cache,
354    /// not to disk).
355    pub fn flush_async(&self) -> Result<(), VecError> {
356        self.mmap.flush_async()?;
357        Ok(())
358    }
359}
360
361#[cfg(test)]
362mod tests {
363    use super::*;
364    use std::sync::Arc;
365    use std::thread;
366
367    fn tmp(name: &str) -> std::path::PathBuf {
368        let mut p = std::env::temp_dir();
369        let pid = std::process::id();
370        p.push(format!("subetha-vec-{name}-{pid}.bin"));
371        p
372    }
373
374    #[test]
375    fn create_initial_state_is_empty() {
376        let p = tmp("init");
377        let v: SharedVec<u32> = SharedVec::create(&p, 16).unwrap();
378        assert_eq!(v.capacity(), 16);
379        assert_eq!(v.len(), 0);
380        assert!(v.is_empty());
381        assert_eq!(v.get(0), None);
382        std::fs::remove_file(&p).ok();
383    }
384
385    #[test]
386    fn push_back_advances_len_and_get_round_trip() {
387        let p = tmp("push");
388        let v: SharedVec<u32> = SharedVec::create(&p, 8).unwrap();
389        for i in 0..5u32 {
390            let idx = v.push_back(i * 10).unwrap();
391            assert_eq!(idx, i as usize);
392        }
393        assert_eq!(v.len(), 5);
394        for i in 0..5 {
395            assert_eq!(v.get(i), Some((i as u32) * 10));
396        }
397        assert_eq!(v.get(5), None);
398        std::fs::remove_file(&p).ok();
399    }
400
401    #[test]
402    fn full_capacity_returns_error() {
403        let p = tmp("full");
404        let v: SharedVec<u32> = SharedVec::create(&p, 4).unwrap();
405        for i in 0..4u32 { v.push_back(i).unwrap(); }
406        assert_eq!(v.push_back(99).err(), Some(VecError::Full));
407        assert_eq!(v.len(), 4);  // rolled back
408        std::fs::remove_file(&p).ok();
409    }
410
411    #[test]
412    fn pop_back_returns_last_then_none() {
413        let p = tmp("pop");
414        let v: SharedVec<u32> = SharedVec::create(&p, 8).unwrap();
415        v.push_back(10).unwrap();
416        v.push_back(20).unwrap();
417        v.push_back(30).unwrap();
418        assert_eq!(v.pop_back(), Some(30));
419        assert_eq!(v.pop_back(), Some(20));
420        assert_eq!(v.len(), 1);
421        assert_eq!(v.pop_back(), Some(10));
422        assert_eq!(v.pop_back(), None);
423        std::fs::remove_file(&p).ok();
424    }
425
426    #[test]
427    fn set_replaces_value_at_index() {
428        let p = tmp("set");
429        let v: SharedVec<u32> = SharedVec::create(&p, 8).unwrap();
430        v.push_back(1).unwrap();
431        v.push_back(2).unwrap();
432        v.set(0, 100).unwrap();
433        assert_eq!(v.get(0), Some(100));
434        assert_eq!(v.get(1), Some(2));
435        assert_eq!(v.set(2, 200).err(), Some(VecError::OutOfBounds));
436        std::fs::remove_file(&p).ok();
437    }
438
439    #[test]
440    fn clear_resets_len_to_zero() {
441        let p = tmp("clear");
442        let v: SharedVec<u32> = SharedVec::create(&p, 8).unwrap();
443        for i in 0..5u32 { v.push_back(i).unwrap(); }
444        assert_eq!(v.len(), 5);
445        v.clear();
446        assert_eq!(v.len(), 0);
447        assert_eq!(v.get(0), None);
448        // After clear, push works fresh.
449        v.push_back(42).unwrap();
450        assert_eq!(v.get(0), Some(42));
451        std::fs::remove_file(&p).ok();
452    }
453
454    #[test]
455    fn snapshot_returns_consistent_prefix() {
456        let p = tmp("snapshot");
457        let v: SharedVec<u32> = SharedVec::create(&p, 16).unwrap();
458        for i in 0..7u32 { v.push_back(i + 100).unwrap(); }
459        let snap = v.snapshot();
460        assert_eq!(snap, vec![100, 101, 102, 103, 104, 105, 106]);
461        std::fs::remove_file(&p).ok();
462    }
463
464    #[test]
465    fn cross_handle_visibility() {
466        let p = tmp("cross-handle");
467        let writer: SharedVec<u32> = SharedVec::create(&p, 8).unwrap();
468        let reader: SharedVec<u32> = SharedVec::open(&p, 8).unwrap();
469        writer.push_back(777).unwrap();
470        assert_eq!(reader.get(0), Some(777));
471        reader.push_back(888).unwrap();
472        assert_eq!(writer.get(1), Some(888));
473        assert_eq!(writer.len(), 2);
474        std::fs::remove_file(&p).ok();
475    }
476
477    #[test]
478    fn concurrent_pushers_all_land_at_distinct_indices() {
479        let p = tmp("concurrent");
480        let v: Arc<SharedVec<u32>> = Arc::new(SharedVec::create(&p, 1024).unwrap());
481        let n_threads = 4;
482        let per_thread = 50u32;
483        let mut handles = vec![];
484        for t in 0..n_threads {
485            let v = v.clone();
486            handles.push(thread::spawn(move || {
487                let mut indices = vec![];
488                for i in 0..per_thread {
489                    let val = (t as u32) * per_thread + i;
490                    let idx = v.push_back(val).unwrap();
491                    indices.push(idx);
492                }
493                indices
494            }));
495        }
496        let mut all_indices: Vec<usize> = handles.into_iter()
497            .flat_map(|h| h.join().unwrap())
498            .collect();
499        all_indices.sort();
500        for (expected, actual) in all_indices.iter().enumerate() {
501            assert_eq!(*actual, expected,
502                "indices must form a contiguous 0..N sequence");
503        }
504        assert_eq!(v.len(), (n_threads * per_thread as usize));
505        std::fs::remove_file(&p).ok();
506    }
507
508    #[test]
509    fn payload_too_large_at_create() {
510        #[allow(dead_code)] // size_of<Big> is the test signal, not the field
511        struct Big([u8; VEC_PAYLOAD_BYTES + 1]);
512        impl Copy for Big {}
513        impl Clone for Big { fn clone(&self) -> Self { *self } }
514        let p = tmp("too-large");
515        let r = SharedVec::<Big>::create(&p, 4);
516        assert_eq!(r.err(), Some(VecError::PayloadTooLarge));
517        std::fs::remove_file(&p).ok();
518    }
519
520    #[test]
521    fn struct_payload_round_trip() {
522        #[derive(Clone, Copy, Debug, PartialEq)]
523        #[repr(C)]
524        struct Point { x: f64, y: f64, z: f64 }
525        let p = tmp("struct");
526        let v: SharedVec<Point> = SharedVec::create(&p, 8).unwrap();
527        v.push_back(Point { x: 1.0, y: 2.0, z: 3.0 }).unwrap();
528        v.push_back(Point { x: -1.5, y: 0.0, z: 7.25 }).unwrap();
529        assert_eq!(v.get(0), Some(Point { x: 1.0, y: 2.0, z: 3.0 }));
530        assert_eq!(v.get(1), Some(Point { x: -1.5, y: 0.0, z: 7.25 }));
531        std::fs::remove_file(&p).ok();
532    }
533
534    #[test]
535    fn disk_persistence_data_survives_reopen() {
536        let p = tmp("disk");
537        {
538            let v: SharedVec<u32> = SharedVec::create(&p, 8).unwrap();
539            for i in 0..4u32 { v.push_back(i * 100).unwrap(); }
540            v.flush().unwrap();
541        }
542        let v2: SharedVec<u32> = SharedVec::open(&p, 8).unwrap();
543        assert_eq!(v2.len(), 4);
544        for i in 0..4 {
545            assert_eq!(v2.get(i), Some((i as u32) * 100));
546        }
547        std::fs::remove_file(&p).ok();
548    }
549
550    #[test]
551    fn concurrent_reader_during_writes_sees_consistent_data() {
552        let p = tmp("read-during-write");
553        let v: Arc<SharedVec<u32>> = Arc::new(SharedVec::create(&p, 256).unwrap());
554        let v_w = v.clone();
555        let writer = thread::spawn(move || {
556            for i in 0..100u32 {
557                v_w.push_back(i).unwrap();
558            }
559        });
560        let v_r = v.clone();
561        let reader = thread::spawn(move || {
562            let mut last_len = 0;
563            loop {
564                let n = v_r.len();
565                if n == 100 { break; }
566                // Read every visible slot; values must equal index.
567                for i in last_len..n {
568                    let got = v_r.get(i);
569                    assert_eq!(got, Some(i as u32),
570                        "slot {i} should hold {i}, got {got:?}");
571                }
572                last_len = n;
573                std::thread::yield_now();
574            }
575        });
576        writer.join().unwrap();
577        reader.join().unwrap();
578        std::fs::remove_file(&p).ok();
579    }
580}