use std::fs::{File, OpenOptions};
use std::mem::size_of;
use std::path::Path;
use std::sync::atomic::{AtomicU64, Ordering};
use memmap2::{Mmap, MmapMut, MmapOptions};
pub const ARENA_MAGIC: u64 = 0x4150_5341_524E_4131;
enum Mapping {
Writable(MmapMut),
ReadOnly(Mmap),
}
impl Mapping {
#[inline]
fn as_ptr(&self) -> *const u8 {
match self {
Mapping::Writable(m) => m.as_ptr(),
Mapping::ReadOnly(m) => m.as_ptr(),
}
}
#[inline]
fn is_writable(&self) -> bool {
matches!(self, Mapping::Writable(_))
}
fn flush(&self) -> Result<(), std::io::Error> {
match self {
Mapping::Writable(m) => m.flush(),
Mapping::ReadOnly(_) => Ok(()),
}
}
fn flush_async(&self) -> Result<(), std::io::Error> {
match self {
Mapping::Writable(m) => m.flush_async(),
Mapping::ReadOnly(_) => Ok(()),
}
}
}
#[repr(C, align(64))]
pub struct ArenaHeader {
pub magic: u64,
pub capacity_bytes: u64,
pub used_bytes: AtomicU64,
_pad: [u8; 40],
}
const _: () = {
assert!(size_of::<ArenaHeader>() == 64);
};
pub const fn arena_file_size(capacity_bytes: usize) -> usize {
size_of::<ArenaHeader>() + capacity_bytes
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ArenaError {
Full,
InvalidRef,
InvalidUtf8,
LayoutMismatch,
ReadOnly,
IoError(std::io::ErrorKind),
}
impl From<std::io::Error> for ArenaError {
fn from(e: std::io::Error) -> Self { Self::IoError(e.kind()) }
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub struct StringRef {
pub offset: u32,
pub len: u32,
}
impl StringRef {
#[inline]
pub fn to_u64(self) -> u64 {
((self.offset as u64) << 32) | (self.len as u64)
}
#[inline]
pub fn from_u64(v: u64) -> Self {
Self {
offset: (v >> 32) as u32,
len: v as u32,
}
}
}
pub struct SharedStringArena {
_file: File,
mmap: Mapping,
capacity_bytes: usize,
header_sidecar: subetha_core::HandshakeHeader,
ring_sidecar: Box<subetha_core::ObservationRing>,
}
unsafe impl Send for SharedStringArena {}
unsafe impl Sync for SharedStringArena {}
impl subetha_sidecar::AdaptiveInstance for SharedStringArena {
fn header(&self) -> &subetha_core::HandshakeHeader { &self.header_sidecar }
fn ring(&self) -> &subetha_core::ObservationRing { &self.ring_sidecar }
fn make_policy(&self) -> Box<dyn subetha_sidecar::Policy> {
Box::new(subetha_sidecar::NoMigrationPolicy)
}
}
impl SharedStringArena {
pub fn create(
path: impl AsRef<Path>, capacity_bytes: usize,
) -> Result<Self, ArenaError> {
Self::check_capacity(capacity_bytes)?;
let (file, mmap) = crate::mmf_attach::create_or_attach(
path.as_ref(),
arena_file_size(capacity_bytes),
|ptr| unsafe { Self::init_region(ptr, capacity_bytes) },
|ptr| unsafe { (*(ptr as *const ArenaHeader)).magic == ARENA_MAGIC },
)?;
let this = Self {
_file: file, mmap: Mapping::Writable(mmap), capacity_bytes,
header_sidecar: subetha_core::HandshakeHeader::new(),
ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
};
this.validate(capacity_bytes)?;
Ok(this)
}
pub fn reset(
path: impl AsRef<Path>, capacity_bytes: usize,
) -> Result<Self, ArenaError> {
Self::check_capacity(capacity_bytes)?;
let (file, mmap) = crate::mmf_attach::reset(
path.as_ref(),
arena_file_size(capacity_bytes),
|ptr| unsafe { Self::init_region(ptr, capacity_bytes) },
)?;
Ok(Self {
_file: file, mmap: Mapping::Writable(mmap), capacity_bytes,
header_sidecar: subetha_core::HandshakeHeader::new(),
ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
})
}
fn check_capacity(capacity_bytes: usize) -> Result<(), ArenaError> {
if capacity_bytes < 1 || capacity_bytes > u32::MAX as usize {
return Err(ArenaError::LayoutMismatch);
}
Ok(())
}
unsafe fn init_region(ptr: *mut u8, capacity_bytes: usize) {
let hdr = ptr as *mut ArenaHeader;
unsafe {
(*hdr).capacity_bytes = capacity_bytes as u64;
std::ptr::write_volatile(&raw mut (*hdr).magic, ARENA_MAGIC);
}
}
pub fn open(
path: impl AsRef<Path>, expected_capacity_bytes: usize,
) -> Result<Self, ArenaError> {
if expected_capacity_bytes > u32::MAX as usize {
return Err(ArenaError::LayoutMismatch);
}
let total = arena_file_size(expected_capacity_bytes);
let file = OpenOptions::new().read(true).write(true).open(path.as_ref())?;
if file.metadata()?.len() < total as u64 {
return Err(ArenaError::LayoutMismatch);
}
let mmap = unsafe { MmapOptions::new().len(total).map_mut(&file)? };
let this = Self {
_file: file, mmap: Mapping::Writable(mmap),
capacity_bytes: expected_capacity_bytes,
header_sidecar: subetha_core::HandshakeHeader::new(),
ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
};
this.validate(expected_capacity_bytes)?;
Ok(this)
}
pub fn open_read_only(
path: impl AsRef<Path>, expected_capacity_bytes: usize,
) -> Result<Self, ArenaError> {
let total = arena_file_size(expected_capacity_bytes);
let file = OpenOptions::new().read(true).open(path.as_ref())?;
if file.metadata()?.len() < total as u64 {
return Err(ArenaError::LayoutMismatch);
}
let mmap = unsafe { MmapOptions::new().len(total).map(&file)? };
let this = Self {
_file: file, mmap: Mapping::ReadOnly(mmap),
capacity_bytes: expected_capacity_bytes,
header_sidecar: subetha_core::HandshakeHeader::new(),
ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
};
this.validate(expected_capacity_bytes)?;
Ok(this)
}
fn validate(&self, expected_capacity_bytes: usize) -> Result<(), ArenaError> {
let hdr = self.header();
if hdr.magic != ARENA_MAGIC || hdr.capacity_bytes != expected_capacity_bytes as u64 {
return Err(ArenaError::LayoutMismatch);
}
Ok(())
}
#[inline]
pub fn is_writable(&self) -> bool {
self.mmap.is_writable()
}
#[inline]
pub fn capacity_bytes(&self) -> usize { self.capacity_bytes }
#[inline]
pub fn used_bytes(&self) -> usize {
self.header().used_bytes.load(Ordering::Acquire) as usize
}
#[inline]
pub fn remaining_bytes(&self) -> usize {
self.capacity_bytes.saturating_sub(self.used_bytes())
}
fn header(&self) -> &ArenaHeader {
unsafe { &*(self.mmap.as_ptr() as *const ArenaHeader) }
}
pub fn intern(&self, s: &str) -> Result<StringRef, ArenaError> {
self.intern_bytes(s.as_bytes())
}
pub fn intern_bytes(&self, bytes: &[u8]) -> Result<StringRef, ArenaError> {
if !self.mmap.is_writable() {
return Err(ArenaError::ReadOnly);
}
let len = bytes.len() as u64;
if len > self.capacity_bytes as u64 {
self.ring_sidecar
.push_op(crate::sidecar_ops::string_arena::OP_INTERN, 1);
return Err(ArenaError::Full);
}
let offset = self.header().used_bytes.fetch_add(len, Ordering::AcqRel);
if offset.saturating_add(len) > self.capacity_bytes as u64 {
self.header().used_bytes.fetch_sub(len, Ordering::AcqRel);
self.ring_sidecar
.push_op(crate::sidecar_ops::string_arena::OP_INTERN, 1);
return Err(ArenaError::Full);
}
if offset.saturating_add(len) > u32::MAX as u64 {
self.header().used_bytes.fetch_sub(len, Ordering::AcqRel);
self.ring_sidecar
.push_op(crate::sidecar_ops::string_arena::OP_INTERN, 1);
return Err(ArenaError::Full);
}
let dst = unsafe {
self.mmap.as_ptr()
.add(size_of::<ArenaHeader>())
.add(offset as usize)
as *mut u8
};
unsafe {
std::ptr::copy_nonoverlapping(bytes.as_ptr(), dst, bytes.len());
}
self.ring_sidecar
.push_op(crate::sidecar_ops::string_arena::OP_INTERN, 0);
Ok(StringRef { offset: offset as u32, len: len as u32 })
}
pub fn get_bytes(&self, r: StringRef) -> Result<&[u8], ArenaError> {
let end = (r.offset as u64).saturating_add(r.len as u64);
if end > self.header().used_bytes.load(Ordering::Acquire) {
self.ring_sidecar
.push_op(crate::sidecar_ops::string_arena::OP_GET_BYTES, 1);
return Err(ArenaError::InvalidRef);
}
if end > self.capacity_bytes as u64 {
self.ring_sidecar
.push_op(crate::sidecar_ops::string_arena::OP_GET_BYTES, 1);
return Err(ArenaError::InvalidRef);
}
self.ring_sidecar
.push_op(crate::sidecar_ops::string_arena::OP_GET_BYTES, 0);
let base = unsafe {
self.mmap.as_ptr()
.add(size_of::<ArenaHeader>())
.add(r.offset as usize)
};
Ok(unsafe { std::slice::from_raw_parts(base, r.len as usize) })
}
pub fn get(&self, r: StringRef) -> Result<&str, ArenaError> {
let bytes = self.get_bytes(r)?;
std::str::from_utf8(bytes).map_err(|_| ArenaError::InvalidUtf8)
}
pub fn intern_and_get(&self, s: &str) -> Result<(StringRef, &str), ArenaError> {
let r = self.intern(s)?;
let got = self.get(r)?;
Ok((r, got))
}
pub fn clear(&self) {
if !self.mmap.is_writable() {
return;
}
self.header().used_bytes.store(0, Ordering::Release);
self.ring_sidecar
.push_op(crate::sidecar_ops::string_arena::OP_CLEAR, 0);
}
pub fn flush(&self) -> Result<(), ArenaError> {
self.mmap.flush()?;
Ok(())
}
pub fn flush_async(&self) -> Result<(), ArenaError> {
self.mmap.flush_async()?;
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
use std::thread;
fn tmp(name: &str) -> std::path::PathBuf {
let mut p = std::env::temp_dir();
let pid = std::process::id();
p.push(format!("subetha-arena-{name}-{pid}.bin"));
p
}
#[test]
fn second_create_attaches_and_keeps_strings() {
let p = tmp("attach");
std::fs::remove_file(&p).ok();
let a = SharedStringArena::create(&p, 4096).unwrap();
let r = a.intern("held").unwrap();
let a2 = SharedStringArena::create(&p, 4096).unwrap();
assert_eq!(a2.get(r).unwrap(), "held", "attach lost an interned string");
assert!(matches!(
SharedStringArena::create(&p, 2048),
Err(ArenaError::LayoutMismatch),
));
drop(a);
drop(a2);
let fresh = SharedStringArena::reset(&p, 4096).unwrap();
assert_eq!(fresh.used_bytes(), 0, "reset kept interned bytes");
drop(fresh);
std::fs::remove_file(&p).ok();
}
#[test]
fn a_read_only_arena_resolves_refs_and_refuses_interning() {
let p = tmp("readonly");
let r = {
let w = SharedStringArena::create(&p, 4096).unwrap();
let r = w.intern("notepad.exe").unwrap();
w.flush().unwrap();
r
};
let ro = SharedStringArena::open_read_only(&p, 4096).unwrap();
assert!(!ro.is_writable());
assert_eq!(ro.get(r), Ok("notepad.exe"));
assert_eq!(ro.get_bytes(r), Ok(&b"notepad.exe"[..]));
assert_eq!(ro.used_bytes(), 11);
assert_eq!(ro.intern("more"), Err(ArenaError::ReadOnly));
assert_eq!(ro.intern_bytes(b"more"), Err(ArenaError::ReadOnly));
ro.clear();
assert_eq!(ro.used_bytes(), 11, "clear on a read-only arena is inert");
ro.flush().unwrap();
std::fs::remove_file(&p).ok();
}
#[test]
fn a_read_only_open_still_validates_the_header() {
let p = tmp("readonly-mismatch");
{
let w = SharedStringArena::create(&p, 4096).unwrap();
w.intern("x").unwrap();
w.flush().unwrap();
}
assert_eq!(
SharedStringArena::open_read_only(&p, 2048).err(),
Some(ArenaError::LayoutMismatch)
);
std::fs::remove_file(&p).ok();
}
#[test]
fn create_initial_state_is_empty() {
let p = tmp("init");
let a = SharedStringArena::create(&p, 1024).unwrap();
assert_eq!(a.capacity_bytes(), 1024);
assert_eq!(a.used_bytes(), 0);
assert_eq!(a.remaining_bytes(), 1024);
std::fs::remove_file(&p).ok();
}
#[test]
fn intern_and_get_round_trip() {
let p = tmp("rt");
let a = SharedStringArena::create(&p, 1024).unwrap();
let r1 = a.intern("hello").unwrap();
let r2 = a.intern("world").unwrap();
assert_eq!(a.get(r1).unwrap(), "hello");
assert_eq!(a.get(r2).unwrap(), "world");
assert_eq!(a.used_bytes(), 10);
std::fs::remove_file(&p).ok();
}
#[test]
fn empty_string_interns_with_zero_len() {
let p = tmp("empty");
let a = SharedStringArena::create(&p, 16).unwrap();
let r = a.intern("").unwrap();
assert_eq!(r.len, 0);
assert_eq!(a.get(r).unwrap(), "");
assert_eq!(a.used_bytes(), 0);
std::fs::remove_file(&p).ok();
}
#[test]
fn full_arena_returns_error() {
let p = tmp("full");
let a = SharedStringArena::create(&p, 10).unwrap();
a.intern("hello").unwrap();
a.intern("world").unwrap();
assert_eq!(a.intern("more").err(), Some(ArenaError::Full));
assert_eq!(a.used_bytes(), 10);
std::fs::remove_file(&p).ok();
}
#[test]
fn string_too_large_returns_full() {
let p = tmp("too-large");
let a = SharedStringArena::create(&p, 8).unwrap();
let big = "x".repeat(100);
assert_eq!(a.intern(&big).err(), Some(ArenaError::Full));
assert_eq!(a.used_bytes(), 0);
std::fs::remove_file(&p).ok();
}
#[test]
fn string_ref_packs_and_unpacks() {
let r = StringRef { offset: 0x1234_5678, len: 42 };
let packed = r.to_u64();
let unpacked = StringRef::from_u64(packed);
assert_eq!(unpacked, r);
}
#[test]
fn cross_handle_visibility() {
let p = tmp("cross-handle");
let writer = SharedStringArena::create(&p, 1024).unwrap();
let reader = SharedStringArena::open(&p, 1024).unwrap();
let r = writer.intern("cross-process").unwrap();
assert_eq!(reader.get(r).unwrap(), "cross-process");
std::fs::remove_file(&p).ok();
}
#[test]
fn invalid_ref_beyond_used_rejected() {
let p = tmp("invalid");
let a = SharedStringArena::create(&p, 1024).unwrap();
a.intern("hi").unwrap(); let bad = StringRef { offset: 100, len: 5 };
assert_eq!(a.get(bad).err(), Some(ArenaError::InvalidRef));
std::fs::remove_file(&p).ok();
}
#[test]
fn concurrent_interners_get_distinct_refs() {
let p = tmp("concurrent");
let a: Arc<SharedStringArena> = Arc::new(SharedStringArena::create(&p, 4096).unwrap());
let n_threads = 4;
let per_thread = 20;
let mut handles = vec![];
for t in 0..n_threads {
let a = a.clone();
handles.push(thread::spawn(move || {
let mut refs = vec![];
for i in 0..per_thread {
let s = format!("thread-{t}-msg-{i:03}");
let r = a.intern(&s).unwrap();
refs.push((s, r));
}
refs
}));
}
let all: Vec<(String, StringRef)> = handles.into_iter()
.flat_map(|h| h.join().unwrap())
.collect();
for (expected, r) in &all {
let got = a.get(*r).unwrap();
assert_eq!(got, expected,
"ref offset={} len={} should resolve to {expected}",
r.offset, r.len);
}
let mut refs: Vec<StringRef> = all.iter().map(|(_, r)| *r).collect();
refs.sort_by_key(|r| r.offset);
for w in refs.windows(2) {
let r1_end = w[0].offset + w[0].len;
assert!(r1_end <= w[1].offset,
"ref {:?} overlaps with ref {:?}", w[0], w[1]);
}
std::fs::remove_file(&p).ok();
}
#[test]
fn intern_and_get_helper_returns_both() {
let p = tmp("intern-and-get");
let a = SharedStringArena::create(&p, 1024).unwrap();
let (r, s) = a.intern_and_get("composite").unwrap();
assert_eq!(s, "composite");
assert_eq!(a.get(r).unwrap(), "composite");
std::fs::remove_file(&p).ok();
}
#[test]
fn utf8_validation_on_get() {
let p = tmp("utf8");
let a = SharedStringArena::create(&p, 128).unwrap();
let r = a.intern("hello").unwrap();
assert!(a.get(r).is_ok());
let r2 = a.intern_bytes(&[0xFF, 0xFE, 0xFD]).unwrap();
assert_eq!(a.get(r2).err(), Some(ArenaError::InvalidUtf8));
assert_eq!(a.get_bytes(r2).unwrap(), &[0xFF, 0xFE, 0xFD]);
std::fs::remove_file(&p).ok();
}
#[test]
fn clear_resets_used_bytes() {
let p = tmp("clear");
let a = SharedStringArena::create(&p, 128).unwrap();
a.intern("first").unwrap();
a.intern("second").unwrap();
assert!(a.used_bytes() > 0);
a.clear();
assert_eq!(a.used_bytes(), 0);
let r = a.intern("after-clear").unwrap();
assert_eq!(a.get(r).unwrap(), "after-clear");
assert_eq!(r.offset, 0);
std::fs::remove_file(&p).ok();
}
#[test]
fn disk_persistence_survives_reopen() {
let p = tmp("disk");
let r_persist;
{
let a = SharedStringArena::create(&p, 1024).unwrap();
r_persist = a.intern("persisted-string").unwrap();
a.flush().unwrap();
}
let a2 = SharedStringArena::open(&p, 1024).unwrap();
assert_eq!(a2.get(r_persist).unwrap(), "persisted-string");
let r2 = a2.intern("more-after-reopen").unwrap();
assert_eq!(a2.get(r2).unwrap(), "more-after-reopen");
std::fs::remove_file(&p).ok();
}
#[test]
fn create_refuses_a_capacity_a_ref_cannot_address() {
assert!(matches!(
SharedStringArena::create(tmp("too-big"), u32::MAX as usize + 1),
Err(ArenaError::LayoutMismatch),
));
assert!(matches!(
SharedStringArena::reset(tmp("too-big-reset"), u32::MAX as usize + 1),
Err(ArenaError::LayoutMismatch),
));
assert!(matches!(
SharedStringArena::create(tmp("zero-cap"), 0),
Err(ArenaError::LayoutMismatch),
));
}
#[test]
fn open_refuses_a_capacity_a_ref_cannot_address() {
assert!(matches!(
SharedStringArena::open(tmp("open-too-big"), u32::MAX as usize + 1),
Err(ArenaError::LayoutMismatch),
));
}
#[test]
fn deduplication_via_hashmap_composition() {
use crate::SharedHashMap;
use crate::shared_hash_map::fnv1a_64;
let p_arena = tmp("dedup-arena");
let p_index = tmp("dedup-index");
let arena = SharedStringArena::create(&p_arena, 256).unwrap();
let index: SharedHashMap<u64, u64> = SharedHashMap::create(&p_index, 32).unwrap();
let s = "deduplicate-me";
let h = fnv1a_64(s.as_bytes());
let r = if let Some(packed) = index.get(&h) {
StringRef::from_u64(packed)
} else {
let r = arena.intern(s).unwrap();
index.insert(h, r.to_u64()).unwrap();
r
};
let used_after_first = arena.used_bytes();
let r2 = if let Some(packed) = index.get(&h) {
StringRef::from_u64(packed)
} else {
let r = arena.intern(s).unwrap();
index.insert(h, r.to_u64()).unwrap();
r
};
assert_eq!(r, r2, "dedup should return the same ref");
assert_eq!(arena.used_bytes(), used_after_first,
"second intern should not consume more bytes");
std::fs::remove_file(&p_arena).ok();
std::fs::remove_file(&p_index).ok();
}
}