1use std::fs::{File, OpenOptions};
62use std::marker::PhantomData;
63use std::mem::size_of;
64use std::path::Path;
65use std::sync::atomic::{AtomicU32, AtomicU64, Ordering};
66
67use memmap2::{MmapMut, MmapOptions};
68
69pub const REGION_MAGIC: u32 = 0x4150_5247;
70pub const NIL_INDEX: u32 = u32::MAX;
71
72#[repr(C, align(64))]
73pub struct RegionHeader {
74 pub magic: u32,
75 pub capacity: u32,
76 pub slot_size: u32,
77 _pad1: u32,
78 pub bump_next: AtomicU64,
79 pub free_head: AtomicU64,
80 _pad2: [u8; 32],
81}
82
83const _: () = {
84 assert!(size_of::<RegionHeader>() == 64);
85};
86
87pub const fn region_file_size(capacity: usize, slot_size: usize) -> usize {
88 size_of::<RegionHeader>()
89 + capacity * size_of::<AtomicU32>() + capacity * slot_size }
92
93#[derive(Debug, Clone, Copy, PartialEq, Eq)]
94pub enum RegionError {
95 Full,
96 InvalidPtr,
97 PayloadTooLarge,
98 LayoutMismatch,
99 IoError(std::io::ErrorKind),
100}
101
102impl From<std::io::Error> for RegionError {
103 fn from(e: std::io::Error) -> Self { Self::IoError(e.kind()) }
104}
105
106#[derive(Debug)]
111#[repr(C)]
112pub struct OffsetPtr<T> {
113 pub index: u32,
114 _phantom: PhantomData<T>,
115}
116
117impl<T> Clone for OffsetPtr<T> {
118 fn clone(&self) -> Self { *self }
119}
120impl<T> Copy for OffsetPtr<T> {}
121impl<T> PartialEq for OffsetPtr<T> {
122 fn eq(&self, other: &Self) -> bool { self.index == other.index }
123}
124impl<T> Eq for OffsetPtr<T> {}
125impl<T> std::hash::Hash for OffsetPtr<T> {
126 fn hash<H: std::hash::Hasher>(&self, state: &mut H) {
127 self.index.hash(state);
128 }
129}
130
131impl<T> OffsetPtr<T> {
132 pub const NIL: Self = Self { index: NIL_INDEX, _phantom: PhantomData };
133
134 #[inline]
135 pub fn new(index: u32) -> Self {
136 Self { index, _phantom: PhantomData }
137 }
138
139 #[inline]
140 pub fn is_nil(self) -> bool { self.index == NIL_INDEX }
141}
142
143#[inline]
144fn pack(counter: u32, index: u32) -> u64 {
145 ((counter as u64) << 32) | (index as u64)
146}
147#[inline]
148fn unpack(v: u64) -> (u32, u32) {
149 ((v >> 32) as u32, v as u32)
150}
151
152pub struct SharedRegion<T: Copy + 'static> {
153 _file: File,
154 mmap: MmapMut,
155 capacity: usize,
156 next_offset: usize,
157 slots_offset: usize,
158 _phantom: PhantomData<T>,
159 header_sidecar: subetha_core::HandshakeHeader,
160 ring_sidecar: Box<subetha_core::ObservationRing>,
161}
162
163unsafe impl<T: Copy + Send + 'static> Send for SharedRegion<T> {}
164unsafe impl<T: Copy + Sync + 'static> Sync for SharedRegion<T> {}
165
166impl<T: Copy + Send + Sync + 'static> subetha_sidecar::AdaptiveInstance for SharedRegion<T> {
167 fn header(&self) -> &subetha_core::HandshakeHeader { &self.header_sidecar }
168 fn ring(&self) -> &subetha_core::ObservationRing { &self.ring_sidecar }
169 fn make_policy(&self) -> Box<dyn subetha_sidecar::Policy> {
170 Box::new(subetha_sidecar::NoMigrationPolicy)
171 }
172}
173
174impl<T: Copy + 'static> SharedRegion<T> {
175 pub fn create(
176 path: impl AsRef<Path>, capacity: usize,
177 ) -> Result<Self, RegionError> {
178 assert!(capacity >= 1);
179 assert!(capacity < NIL_INDEX as usize, "capacity must be < u32::MAX");
180 let slot_size = size_of::<T>();
181 let total = region_file_size(capacity, slot_size);
182 let file = OpenOptions::new()
183 .read(true).write(true).create(true).truncate(true)
184 .open(path.as_ref())?;
185 file.set_len(total as u64)?;
186 let mut mmap = unsafe { MmapOptions::new().len(total).map_mut(&file)? };
187 crate::mmf_warm::warm_mmap(&mut mmap);
188 let hdr = mmap.as_mut_ptr() as *mut RegionHeader;
189 unsafe {
190 std::ptr::write_bytes(hdr as *mut u8, 0, size_of::<RegionHeader>());
191 (*hdr).magic = REGION_MAGIC;
192 (*hdr).capacity = capacity as u32;
193 (*hdr).slot_size = slot_size as u32;
194 (*hdr).bump_next.store(0, Ordering::Release);
195 (*hdr).free_head.store(pack(0, NIL_INDEX), Ordering::Release);
196 }
197 let next_offset = size_of::<RegionHeader>();
201 let slots_offset = next_offset + capacity * size_of::<AtomicU32>();
202 Ok(Self {
203 _file: file, mmap, capacity, next_offset, slots_offset,
204 _phantom: PhantomData,
205 header_sidecar: subetha_core::HandshakeHeader::new(),
206 ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
207 })
208 }
209
210 pub fn open(
211 path: impl AsRef<Path>, expected_capacity: usize,
212 ) -> Result<Self, RegionError> {
213 let slot_size = size_of::<T>();
214 let total = region_file_size(expected_capacity, slot_size);
215 let file = OpenOptions::new().read(true).write(true).open(path.as_ref())?;
216 if file.metadata()?.len() < total as u64 {
217 return Err(RegionError::LayoutMismatch);
218 }
219 let mut mmap = unsafe { MmapOptions::new().len(total).map_mut(&file)? };
220 crate::mmf_warm::warm_mmap(&mut mmap);
221 let hdr = unsafe { &*(mmap.as_ptr() as *const RegionHeader) };
222 if hdr.magic != REGION_MAGIC
223 || hdr.capacity != expected_capacity as u32
224 || hdr.slot_size != slot_size as u32
225 {
226 return Err(RegionError::LayoutMismatch);
227 }
228 let next_offset = size_of::<RegionHeader>();
229 let slots_offset = next_offset + expected_capacity * size_of::<AtomicU32>();
230 Ok(Self {
231 _file: file, mmap, capacity: expected_capacity,
232 next_offset, slots_offset,
233 _phantom: PhantomData,
234 header_sidecar: subetha_core::HandshakeHeader::new(),
235 ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
236 })
237 }
238
239 #[inline]
240 pub fn capacity(&self) -> usize { self.capacity }
241
242 #[inline]
247 pub fn mmap_ptr(&self) -> *const u8 { self.mmap.as_ptr() }
248
249 fn header(&self) -> &RegionHeader {
250 unsafe { &*(self.mmap.as_ptr() as *const RegionHeader) }
251 }
252
253 fn next_link(&self, idx: usize) -> &AtomicU32 {
254 let base = unsafe { self.mmap.as_ptr().add(self.next_offset) };
255 unsafe { &*(base.add(idx * size_of::<AtomicU32>()) as *const AtomicU32) }
256 }
257
258 fn slot_ptr(&self, idx: usize) -> *mut T {
259 let base = unsafe { self.mmap.as_ptr().add(self.slots_offset) };
260 unsafe { base.add(idx * size_of::<T>()) as *mut T }
261 }
262
263 pub fn len(&self) -> usize {
267 let bump = self.header().bump_next.load(Ordering::Acquire) as usize;
268 bump.saturating_sub(self.free_count())
269 }
270
271 pub fn is_empty(&self) -> bool { self.len() == 0 }
272
273 pub fn free_count(&self) -> usize {
276 let mut count = 0usize;
277 let (_, mut idx) = unpack(self.header().free_head.load(Ordering::Acquire));
278 let mut visited = 0;
279 while idx != NIL_INDEX && (visited as usize) < self.capacity {
280 count += 1;
281 visited += 1;
282 idx = self.next_link(idx as usize).load(Ordering::Acquire);
283 }
284 count
285 }
286
287 pub fn allocate(&self, value: T) -> Result<OffsetPtr<T>, RegionError> {
290 loop {
292 let head = self.header().free_head.load(Ordering::Acquire);
293 let (counter, idx) = unpack(head);
294 if idx == NIL_INDEX { break; } let next_idx = self.next_link(idx as usize).load(Ordering::Acquire);
296 let new_head = pack(counter.wrapping_add(1), next_idx);
297 if self.header().free_head.compare_exchange(
298 head, new_head, Ordering::AcqRel, Ordering::Acquire,
299 ).is_ok() {
300 unsafe { std::ptr::write(self.slot_ptr(idx as usize), value); }
301 self.ring_sidecar
302 .push_op(crate::sidecar_ops::region::OP_ALLOCATE, 0);
303 return Ok(OffsetPtr::new(idx));
304 }
305 }
307 let idx = self.header().bump_next.fetch_add(1, Ordering::AcqRel);
309 if idx >= self.capacity as u64 {
310 self.header().bump_next.fetch_sub(1, Ordering::AcqRel);
311 self.ring_sidecar
312 .push_op(crate::sidecar_ops::region::OP_ALLOCATE, 1); return Err(RegionError::Full);
314 }
315 let idx = idx as u32;
316 unsafe { std::ptr::write(self.slot_ptr(idx as usize), value); }
317 self.ring_sidecar
318 .push_op(crate::sidecar_ops::region::OP_ALLOCATE, 0);
319 Ok(OffsetPtr::new(idx))
320 }
321
322 pub fn free(&self, ptr: OffsetPtr<T>) -> Result<T, RegionError> {
325 if ptr.is_nil() || (ptr.index as usize) >= self.capacity {
326 self.ring_sidecar
327 .push_op(crate::sidecar_ops::region::OP_FREE, 1); return Err(RegionError::InvalidPtr);
329 }
330 let value = unsafe { std::ptr::read(self.slot_ptr(ptr.index as usize)) };
331 loop {
332 let head = self.header().free_head.load(Ordering::Acquire);
333 let (counter, old_top) = unpack(head);
334 self.next_link(ptr.index as usize).store(old_top, Ordering::Release);
335 let new_head = pack(counter.wrapping_add(1), ptr.index);
336 if self.header().free_head.compare_exchange(
337 head, new_head, Ordering::AcqRel, Ordering::Acquire,
338 ).is_ok() {
339 self.ring_sidecar
340 .push_op(crate::sidecar_ops::region::OP_FREE, 0);
341 return Ok(value);
342 }
343 }
345 }
346
347 pub fn get(&self, ptr: OffsetPtr<T>) -> Result<T, RegionError> {
352 if ptr.is_nil() || (ptr.index as usize) >= self.capacity {
353 self.ring_sidecar
354 .push_op(crate::sidecar_ops::region::OP_GET, 1); return Err(RegionError::InvalidPtr);
356 }
357 let v = unsafe { std::ptr::read(self.slot_ptr(ptr.index as usize)) };
358 self.ring_sidecar
359 .push_op(crate::sidecar_ops::region::OP_GET, 0);
360 Ok(v)
361 }
362
363 pub fn set(&self, ptr: OffsetPtr<T>, value: T) -> Result<(), RegionError> {
366 if ptr.is_nil() || (ptr.index as usize) >= self.capacity {
367 self.ring_sidecar
368 .push_op(crate::sidecar_ops::region::OP_SET, 1); return Err(RegionError::InvalidPtr);
370 }
371 unsafe { std::ptr::write(self.slot_ptr(ptr.index as usize), value); }
372 self.ring_sidecar
373 .push_op(crate::sidecar_ops::region::OP_SET, 0);
374 Ok(())
375 }
376
377 pub fn clear(&self) {
386 self.header().bump_next.store(0, Ordering::Release);
387 self.header().free_head.store(pack(0, NIL_INDEX), Ordering::Release);
388 }
389
390 pub fn flush(&self) -> Result<(), RegionError> {
391 self.mmap.flush()?;
392 Ok(())
393 }
394
395 pub fn flush_async(&self) -> Result<(), RegionError> {
399 self.mmap.flush_async()?;
400 Ok(())
401 }
402}
403
404#[cfg(test)]
405mod tests {
406 use super::*;
407 use std::sync::Arc;
408 use std::thread;
409
410 fn tmp(name: &str) -> std::path::PathBuf {
411 let mut p = std::env::temp_dir();
412 let pid = std::process::id();
413 p.push(format!("subetha-region-{name}-{pid}.bin"));
414 p
415 }
416
417 #[test]
418 fn create_initial_state_is_empty() {
419 let p = tmp("init");
420 let r: SharedRegion<u64> = SharedRegion::create(&p, 16).unwrap();
421 assert_eq!(r.capacity(), 16);
422 assert_eq!(r.len(), 0);
423 assert!(r.is_empty());
424 assert_eq!(r.free_count(), 0);
425 std::fs::remove_file(&p).ok();
426 }
427
428 #[test]
429 fn allocate_and_get_round_trip() {
430 let p = tmp("rt");
431 let r: SharedRegion<u64> = SharedRegion::create(&p, 8).unwrap();
432 let p1 = r.allocate(100).unwrap();
433 let p2 = r.allocate(200).unwrap();
434 let p3 = r.allocate(300).unwrap();
435 assert_eq!(p1.index, 0);
436 assert_eq!(p2.index, 1);
437 assert_eq!(p3.index, 2);
438 assert_eq!(r.get(p1).unwrap(), 100);
439 assert_eq!(r.get(p2).unwrap(), 200);
440 assert_eq!(r.get(p3).unwrap(), 300);
441 assert_eq!(r.len(), 3);
442 std::fs::remove_file(&p).ok();
443 }
444
445 #[test]
446 fn full_region_returns_error_on_bump() {
447 let p = tmp("full");
448 let r: SharedRegion<u32> = SharedRegion::create(&p, 4).unwrap();
449 for i in 0..4u32 { r.allocate(i).unwrap(); }
450 assert_eq!(r.allocate(99).err(), Some(RegionError::Full));
451 assert_eq!(r.len(), 4); std::fs::remove_file(&p).ok();
453 }
454
455 #[test]
456 fn free_returns_value_and_decrements_len() {
457 let p = tmp("free");
458 let r: SharedRegion<u64> = SharedRegion::create(&p, 8).unwrap();
459 let ptr1 = r.allocate(111).unwrap();
460 let ptr2 = r.allocate(222).unwrap();
461 assert_eq!(r.len(), 2);
462 let v = r.free(ptr1).unwrap();
463 assert_eq!(v, 111);
464 assert_eq!(r.len(), 1);
465 assert_eq!(r.free_count(), 1);
466 assert_eq!(r.get(ptr2).unwrap(), 222);
468 std::fs::remove_file(&p).ok();
469 }
470
471 #[test]
472 fn free_then_allocate_reuses_slot() {
473 let p = tmp("reuse");
474 let r: SharedRegion<u64> = SharedRegion::create(&p, 8).unwrap();
475 let p1 = r.allocate(100).unwrap();
476 let _p2 = r.allocate(200).unwrap();
477 r.free(p1).unwrap();
478 let p3 = r.allocate(999).unwrap();
480 assert_eq!(p3.index, p1.index, "free list should have returned slot 0");
481 assert_eq!(r.get(p3).unwrap(), 999);
482 std::fs::remove_file(&p).ok();
483 }
484
485 #[test]
486 fn free_invalid_ptr_returns_error() {
487 let p = tmp("invalid");
488 let r: SharedRegion<u64> = SharedRegion::create(&p, 4).unwrap();
489 assert_eq!(r.free(OffsetPtr::NIL).err(), Some(RegionError::InvalidPtr));
490 let oob = OffsetPtr::<u64>::new(100);
491 assert_eq!(r.free(oob).err(), Some(RegionError::InvalidPtr));
492 std::fs::remove_file(&p).ok();
493 }
494
495 #[test]
496 fn set_overwrites_value_in_place() {
497 let p = tmp("set");
498 let r: SharedRegion<u64> = SharedRegion::create(&p, 4).unwrap();
499 let ptr = r.allocate(42).unwrap();
500 r.set(ptr, 999).unwrap();
501 assert_eq!(r.get(ptr).unwrap(), 999);
502 std::fs::remove_file(&p).ok();
503 }
504
505 #[test]
506 fn cross_handle_visibility() {
507 let p = tmp("cross-handle");
508 let writer: SharedRegion<u64> = SharedRegion::create(&p, 16).unwrap();
509 let reader: SharedRegion<u64> = SharedRegion::open(&p, 16).unwrap();
510 let ptr = writer.allocate(7777).unwrap();
511 assert_eq!(reader.get(ptr).unwrap(), 7777);
512 writer.set(ptr, 8888).unwrap();
513 assert_eq!(reader.get(ptr).unwrap(), 8888);
514 std::fs::remove_file(&p).ok();
515 }
516
517 #[test]
518 fn concurrent_allocations_get_distinct_indices() {
519 let p = tmp("concurrent");
520 let r: Arc<SharedRegion<u64>> = Arc::new(SharedRegion::create(&p, 1024).unwrap());
521 let n_threads = 4;
522 let per_thread = 100;
523 let mut handles = vec![];
524 for t in 0..n_threads as u64 {
525 let r = r.clone();
526 handles.push(thread::spawn(move || {
527 let mut ptrs = vec![];
528 for i in 0..per_thread as u64 {
529 let v = t * 1000 + i;
530 let ptr = r.allocate(v).unwrap();
531 ptrs.push((v, ptr));
532 }
533 ptrs
534 }));
535 }
536 let all: Vec<(u64, OffsetPtr<u64>)> = handles.into_iter()
537 .flat_map(|h| h.join().unwrap()).collect();
538
539 let mut indices: Vec<u32> = all.iter().map(|(_, p)| p.index).collect();
541 indices.sort();
542 for w in indices.windows(2) {
543 assert_ne!(w[0], w[1], "two threads got the same slot index");
544 }
545 for (expected, ptr) in &all {
547 assert_eq!(r.get(*ptr).unwrap(), *expected);
548 }
549 std::fs::remove_file(&p).ok();
550 }
551
552 #[test]
553 fn concurrent_free_and_realloc_no_corruption() {
554 let p = tmp("free-realloc");
555 let r: Arc<SharedRegion<u64>> = Arc::new(SharedRegion::create(&p, 64).unwrap());
556 let initial: Vec<OffsetPtr<u64>> = (0..64u64)
558 .map(|i| r.allocate(i * 10).unwrap()).collect();
559 let r_a = r.clone();
560 let _init = initial.clone();
561 let freer = thread::spawn(move || {
563 for (i, ptr) in _init.iter().enumerate() {
564 if i % 2 == 0 { r_a.free(*ptr).ok(); }
565 }
566 });
567 let r_b = r.clone();
569 let alloc = thread::spawn(move || {
570 let mut new_ptrs = vec![];
571 for i in 100u64..132 {
572 if let Ok(p) = r_b.allocate(i) { new_ptrs.push((i, p)); }
573 }
574 new_ptrs
575 });
576 freer.join().unwrap();
577 let new_ptrs = alloc.join().unwrap();
578 for (val, ptr) in new_ptrs {
580 assert_eq!(r.get(ptr).unwrap(), val);
581 }
582 std::fs::remove_file(&p).ok();
583 }
584
585 #[test]
586 fn offset_ptr_is_position_independent() {
587 let p = tmp("position-indep");
590 let writer: SharedRegion<u64> = SharedRegion::create(&p, 4).unwrap();
591 let ptr = writer.allocate(0xCAFE_BABE).unwrap();
592 let reader: SharedRegion<u64> = SharedRegion::open(&p, 4).unwrap();
593 let same_ptr = OffsetPtr::<u64>::new(ptr.index);
595 assert_eq!(reader.get(same_ptr).unwrap(), 0xCAFE_BABE);
596 std::fs::remove_file(&p).ok();
597 }
598
599 #[test]
600 fn struct_payload_round_trip() {
601 #[derive(Clone, Copy, Debug, PartialEq)]
602 #[repr(C)]
603 struct Node { left: u32, right: u32, key: u64 }
604 let p = tmp("struct");
605 let r: SharedRegion<Node> = SharedRegion::create(&p, 16).unwrap();
606 let n = Node { left: 1, right: 2, key: 42 };
607 let ptr = r.allocate(n).unwrap();
608 assert_eq!(r.get(ptr).unwrap(), n);
609 std::fs::remove_file(&p).ok();
610 }
611
612 #[test]
613 fn nil_ptr_round_trips() {
614 let p: OffsetPtr<u64> = OffsetPtr::NIL;
615 assert!(p.is_nil());
616 assert_eq!(p.index, NIL_INDEX);
617 }
618
619 #[test]
620 fn offset_ptr_equality_and_hash() {
621 use std::collections::HashSet;
622 let a: OffsetPtr<u64> = OffsetPtr::new(5);
623 let b: OffsetPtr<u64> = OffsetPtr::new(5);
624 let c: OffsetPtr<u64> = OffsetPtr::new(6);
625 assert_eq!(a, b);
626 assert_ne!(a, c);
627 let mut s = HashSet::new();
628 s.insert(a);
629 assert!(s.contains(&b));
630 assert!(!s.contains(&c));
631 }
632
633 #[test]
634 fn disk_persistence_survives_reopen() {
635 let p = tmp("disk");
636 let saved_ptr_index;
637 {
638 let r: SharedRegion<u64> = SharedRegion::create(&p, 16).unwrap();
639 let p1 = r.allocate(1111).unwrap();
640 let _p2 = r.allocate(2222).unwrap();
641 r.flush().unwrap();
642 saved_ptr_index = p1.index;
643 }
644 let r2: SharedRegion<u64> = SharedRegion::open(&p, 16).unwrap();
645 let restored = OffsetPtr::<u64>::new(saved_ptr_index);
646 assert_eq!(r2.get(restored).unwrap(), 1111);
647 let p3 = r2.allocate(3333).unwrap();
649 assert_eq!(p3.index, 2);
650 assert_eq!(r2.get(p3).unwrap(), 3333);
651 std::fs::remove_file(&p).ok();
652 }
653}