1use 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 fn write_slot(&self, i: usize, v: T) {
211 let slot = self.slot(i);
212 slot.version.fetch_add(1, Ordering::AcqRel);
214 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 slot.version.fetch_add(1, Ordering::AcqRel);
230 }
231
232 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 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); 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 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); 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 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); 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 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); 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 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 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 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); 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 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)] 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 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}