use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex, PoisonError};
pub const SNAPSHOT_PREALLOC: usize = 256;
pub struct SnapshotSlot {
bytes: Mutex<Vec<u8>>,
swapped: Mutex<Option<Arc<Vec<u8>>>>,
supported: AtomicBool,
}
impl SnapshotSlot {
#[must_use]
pub fn new() -> Arc<Self> {
Self::with_capacity(SNAPSHOT_PREALLOC)
}
#[must_use]
pub fn with_capacity(cap: usize) -> Arc<Self> {
let slot = Self {
bytes: Mutex::new(Vec::with_capacity(cap)),
swapped: Mutex::new(None),
supported: AtomicBool::new(false),
};
drop(slot.bytes.lock().unwrap_or_else(PoisonError::into_inner));
drop(slot.swapped.lock().unwrap_or_else(PoisonError::into_inner));
Arc::new(slot)
}
pub fn publish(&self, write: impl FnOnce(&mut Vec<u8>) -> bool) -> bool {
if let Ok(mut guard) = self.bytes.try_lock() {
guard.clear();
if write(&mut guard) {
self.supported.store(true, Ordering::Release);
}
true
} else {
false
}
}
pub fn store_arc(&self, bytes: Vec<u8>) {
*self.swapped.lock().unwrap_or_else(PoisonError::into_inner) = Some(Arc::new(bytes));
self.supported.store(true, Ordering::Release);
}
#[must_use]
pub fn is_supported(&self) -> bool {
self.supported.load(Ordering::Acquire)
}
#[must_use]
pub fn read(&self) -> Option<Vec<u8>> {
if !self.supported.load(Ordering::Acquire) {
return None;
}
if let Some(arc) = self
.swapped
.lock()
.unwrap_or_else(PoisonError::into_inner)
.as_ref()
{
return Some((**arc).clone());
}
Some(
self.bytes
.lock()
.unwrap_or_else(PoisonError::into_inner)
.clone(),
)
}
#[must_use]
pub fn read_offthread(&self) -> Option<Vec<u8>> {
self.swapped
.lock()
.unwrap_or_else(PoisonError::into_inner)
.as_ref()
.map(|arc| (**arc).clone())
}
}
#[derive(Clone)]
pub struct SnapshotPublisher {
slot: Arc<SnapshotSlot>,
}
impl SnapshotPublisher {
#[must_use]
pub fn new(slot: &Arc<SnapshotSlot>) -> Self {
Self {
slot: Arc::clone(slot),
}
}
pub fn publish(&self, bytes: Vec<u8>) {
self.slot.store_arc(bytes);
}
}
#[cfg(test)]
mod tests {
use super::{SNAPSHOT_PREALLOC, SnapshotPublisher, SnapshotSlot};
#[test]
fn buffer_is_prewarmed_so_first_publish_does_not_allocate() {
let slot = SnapshotSlot::new();
let landed = slot.publish(|buf| {
buf.extend_from_slice(&[0xAB; SNAPSHOT_PREALLOC]);
true
});
assert!(landed);
let published = slot.read().expect("published");
assert_eq!(published.len(), SNAPSHOT_PREALLOC);
}
#[test]
fn read_is_none_until_first_publish() {
let slot = SnapshotSlot::new();
assert!(slot.read().is_none());
slot.publish(|buf| {
buf.extend_from_slice(&[1, 2, 3]);
true
});
assert_eq!(slot.read(), Some(vec![1, 2, 3]));
}
#[test]
fn unsupported_publish_never_marks_supported() {
let slot = SnapshotSlot::new();
assert!(slot.publish(|_| false)); assert!(slot.read().is_none());
}
#[test]
fn is_supported_latches_on_first_true_publish() {
let slot = SnapshotSlot::new();
assert!(!slot.is_supported());
slot.publish(|_| false);
assert!(!slot.is_supported());
slot.publish(|buf| {
buf.push(1);
true
});
assert!(slot.is_supported());
slot.publish(|_| true);
assert!(slot.is_supported());
assert_eq!(slot.read(), Some(vec![]));
}
#[test]
fn writer_that_only_appends_does_not_accumulate_across_publishes() {
let slot = SnapshotSlot::new();
slot.publish(|buf| {
buf.extend_from_slice(&[9; 64]);
true
});
assert_eq!(slot.read(), Some(vec![9; 64]));
slot.publish(|buf| {
buf.extend_from_slice(&[7, 7]);
true
});
assert_eq!(
slot.read(),
Some(vec![7, 7]),
"second publish must replace, not append onto, the first"
);
}
#[test]
fn publish_reports_landed_vs_skipped_on_contention() {
let slot = SnapshotSlot::new();
slot.publish(|buf| {
buf.push(1);
true
});
let held = slot.bytes.lock().unwrap();
let landed = slot.publish(|buf| {
buf.clear();
buf.push(2);
true
});
assert!(!landed, "publish must report a skipped (unlanded) write");
drop(held);
assert_eq!(slot.read(), Some(vec![1]), "the prior snapshot stands");
}
#[test]
fn off_thread_lane_takes_read_preference() {
let slot = SnapshotSlot::new();
slot.publish(|buf| {
buf.extend_from_slice(&[1, 1]);
true
});
assert_eq!(slot.read(), Some(vec![1, 1]));
SnapshotPublisher::new(&slot).publish(vec![2; 4096]);
assert_eq!(slot.read(), Some(vec![2; 4096]));
}
#[test]
fn read_offthread_returns_only_the_publisher_lane() {
let slot = SnapshotSlot::new();
slot.publish(|buf| {
buf.extend_from_slice(&[1, 1]);
true
});
assert!(slot.read_offthread().is_none());
SnapshotPublisher::new(&slot).publish(vec![9; 3]);
assert_eq!(slot.read_offthread(), Some(vec![9; 3]));
}
#[test]
fn store_arc_latches_supported_without_inline_publish() {
let slot = SnapshotSlot::new();
assert!(!slot.is_supported());
slot.store_arc(vec![5, 5, 5]);
assert!(slot.is_supported());
assert_eq!(slot.read(), Some(vec![5, 5, 5]));
}
#[test]
fn with_capacity_prewarms_the_requested_size() {
let slot = SnapshotSlot::with_capacity(8192);
let landed = slot.publish(|buf| {
buf.extend_from_slice(&[0xCD; 8192]);
true
});
assert!(landed);
assert_eq!(slot.read().expect("published").len(), 8192);
}
}