use std::sync::Arc;
use std::sync::RwLock;
use std::sync::atomic::{AtomicU32, Ordering};
const CHUNK: usize = 1024;
const INLINE: usize = 4;
fn new_chunk<T: Default>() -> Arc<[T]> {
(0..CHUNK).map(|_| T::default()).collect()
}
pub struct ChunkedVec<T> {
backbone: RwLock<Arc<Vec<Arc<[T]>>>>,
len: AtomicU32,
}
impl<T: Default> Default for ChunkedVec<T> {
fn default() -> Self {
ChunkedVec {
backbone: RwLock::new(Arc::new(Vec::new())),
len: AtomicU32::new(0),
}
}
}
impl<T: Default> ChunkedVec<T> {
pub fn new() -> Self {
Self::default()
}
pub fn len(&self) -> u32 {
self.len.load(Ordering::Relaxed)
}
fn ensure(&self, i: usize) {
let need = i / CHUNK + 1;
if self.backbone.read().unwrap().len() >= need {
return;
}
let mut guard = self.backbone.write().unwrap();
if guard.len() >= need {
return;
}
let target = need.max(guard.len().saturating_mul(2));
let mut next: Vec<Arc<[T]>> = (**guard).clone();
while next.len() < target {
next.push(new_chunk());
}
*guard = Arc::new(next);
}
pub fn push_with(&self, f: impl FnOnce(&T)) -> u32 {
let i = self.len.load(Ordering::Relaxed) as usize;
self.ensure(i);
{
let bb = self.backbone.read().unwrap();
f(&bb[i / CHUNK][i % CHUNK]);
}
self.len.store((i + 1) as u32, Ordering::Relaxed);
i as u32
}
pub fn with(&self, i: u32, f: impl FnOnce(&T)) {
let i = i as usize;
let bb = self.backbone.read().unwrap();
f(&bb[i / CHUNK][i % CHUNK]);
}
pub fn snapshot(&self) -> Snap<T> {
Snap {
backbone: Arc::clone(&self.backbone.read().unwrap()),
}
}
}
pub struct Snap<T> {
backbone: Arc<Vec<Arc<[T]>>>,
}
impl<T> Snap<T> {
pub fn get(&self, i: u32) -> Option<&T> {
let i = i as usize;
self.backbone.get(i / CHUNK).map(|chunk| &chunk[i % CHUNK])
}
pub fn covered(&self) -> u32 {
(self.backbone.len() * CHUNK) as u32
}
}
#[derive(Default)]
pub struct PostingSlot {
len: AtomicU32,
inline: [AtomicU32; INLINE],
spill: std::sync::OnceLock<ChunkedVec<AtomicU32>>,
}
impl PostingSlot {
pub fn push_local(&self, local: u32) {
let n = self.len.load(Ordering::Relaxed) as usize;
if n < INLINE {
self.inline[n].store(local, Ordering::Relaxed);
} else {
let spill = self.spill.get_or_init(ChunkedVec::new);
spill.push_with(|slot| slot.store(local, Ordering::Relaxed));
}
self.len.store((n + 1) as u32, Ordering::Release);
}
pub fn len(&self) -> u32 {
self.len.load(Ordering::Acquire)
}
pub fn collect_below(&self, upto: u32, out: &mut Vec<u32>) {
let n = self.len.load(Ordering::Acquire) as usize;
let inline_n = n.min(INLINE);
for slot in &self.inline[..inline_n] {
let local = slot.load(Ordering::Relaxed);
if local >= upto {
return;
}
out.push(local);
}
if n > INLINE
&& let Some(spill) = self.spill.get()
{
let snap = spill.snapshot();
for j in 0..(n - INLINE) as u32 {
let Some(slot) = snap.get(j) else { break };
let local = slot.load(Ordering::Relaxed);
if local >= upto {
return;
}
out.push(local);
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicU16, AtomicU64};
use std::thread;
#[test]
fn grows_across_chunk_boundaries_and_reads_back() {
let col: ChunkedVec<AtomicU32> = ChunkedVec::new();
let n = CHUNK as u32 * 3 + 7; for v in 0..n {
col.push_with(|slot| slot.store(v * 2, Ordering::Relaxed));
}
assert_eq!(col.len(), n);
let snap = col.snapshot();
assert!(snap.covered() >= n);
for v in 0..n {
assert_eq!(snap.get(v).unwrap().load(Ordering::Relaxed), v * 2);
}
assert!(snap.get(snap.covered()).is_none());
}
#[test]
fn reader_sees_a_consistent_prefix_under_a_concurrent_writer() {
let col: Arc<ChunkedVec<AtomicU32>> = Arc::new(ChunkedVec::new());
let watermark = Arc::new(AtomicU64::new(0));
let total = 50_000u32;
let writer = {
let col = Arc::clone(&col);
let watermark = Arc::clone(&watermark);
thread::spawn(move || {
for v in 0..total {
col.push_with(|slot| slot.store(v + 1, Ordering::Relaxed));
watermark.store(v as u64 + 1, Ordering::Release);
}
})
};
let reader = {
let col = Arc::clone(&col);
let watermark = Arc::clone(&watermark);
thread::spawn(move || {
loop {
let wm = watermark.load(Ordering::Acquire);
let snap = col.snapshot();
for i in 0..wm as u32 {
assert_eq!(snap.get(i).unwrap().load(Ordering::Relaxed), i + 1);
}
if wm >= total as u64 {
return wm;
}
}
})
};
writer.join().unwrap();
assert_eq!(reader.join().unwrap(), total as u64);
}
#[test]
fn posting_slot_stays_inline_then_spills() {
let slot = PostingSlot::default();
slot.push_local(0);
slot.push_local(3);
assert!(slot.spill.get().is_none());
let mut out = Vec::new();
slot.collect_below(u32::MAX, &mut out);
assert_eq!(out, vec![0, 3]);
for local in 4..20 {
slot.push_local(local);
}
assert!(slot.spill.get().is_some());
out.clear();
slot.collect_below(u32::MAX, &mut out);
let mut expected = vec![0, 3];
expected.extend(4..20);
assert_eq!(out, expected);
}
#[test]
fn posting_slot_truncates_by_upto_across_the_spill_boundary() {
let slot = PostingSlot::default();
for local in 0..20 {
slot.push_local(local);
}
let mut out = Vec::new();
slot.collect_below(3, &mut out);
assert_eq!(out, vec![0, 1, 2]);
out.clear();
slot.collect_below(10, &mut out);
assert_eq!(out, (0..10).collect::<Vec<_>>());
}
#[test]
fn type_column_slots_are_u16() {
let col: ChunkedVec<AtomicU16> = ChunkedVec::new();
for v in 0..2000u16 {
col.push_with(|slot| slot.store(v, Ordering::Relaxed));
}
let snap = col.snapshot();
assert_eq!(snap.get(0).unwrap().load(Ordering::Relaxed), 0);
assert_eq!(snap.get(1999).unwrap().load(Ordering::Relaxed), 1999);
}
}