use std::any::Any;
use std::ops::Deref;
use std::sync::Arc;
pub type Pin = Arc<dyn Any + Send + Sync>;
#[derive(Debug, Clone)]
pub struct Buffer<T> {
store: Store<T>,
}
#[derive(Debug, Clone)]
enum Store<T> {
Owned(Vec<T>),
Shared(Arc<Vec<T>>),
}
impl<T> Buffer<T> {
#[must_use]
pub fn new() -> Self {
Self { store: Store::Owned(Vec::new()) }
}
#[must_use]
pub fn with_capacity(capacity: usize) -> Self {
Self { store: Store::Owned(Vec::with_capacity(capacity)) }
}
#[must_use]
pub fn from_vec(values: Vec<T>) -> Self {
Self { store: Store::Owned(values) }
}
#[must_use]
pub fn from_arc(page: Arc<Vec<T>>) -> Self {
Self { store: Store::Shared(page) }
}
#[must_use]
pub fn is_shared(&self) -> bool {
matches!(self.store, Store::Shared(_))
}
#[must_use]
pub fn footprint(&self) -> usize {
match &self.store {
Store::Owned(values) => values.capacity() * size_of::<T>(),
Store::Shared(page) => page.capacity() * size_of::<T>() / Arc::strong_count(page),
}
}
#[must_use]
#[inline]
pub fn as_slice(&self) -> &[T] {
match &self.store {
Store::Owned(values) => values,
Store::Shared(page) => page,
}
}
}
impl<T: Clone> Buffer<T> {
#[must_use]
pub fn into_vec(self) -> Vec<T> {
match self.store {
Store::Owned(values) => values,
Store::Shared(page) => Arc::try_unwrap(page).unwrap_or_else(|page| page.to_vec()),
}
}
#[inline]
pub fn to_mut(&mut self) -> &mut Vec<T> {
if let Store::Shared(page) = &self.store {
self.store = Store::Owned(page.to_vec());
}
match &mut self.store {
Store::Owned(values) => values,
Store::Shared(_) => unreachable!("a shared page was copied out one statement ago"),
}
}
#[inline]
pub fn push(&mut self, value: T) {
self.to_mut().push(value);
}
pub fn reserve(&mut self, additional: usize) {
self.to_mut().reserve(additional);
}
pub fn extend_from_slice(&mut self, values: &[T]) {
self.to_mut().extend_from_slice(values);
}
}
impl<T: PartialEq> PartialEq for Buffer<T> {
fn eq(&self, other: &Self) -> bool {
self.as_slice() == other.as_slice()
}
}
impl<T: Eq> Eq for Buffer<T> {}
impl<T> Default for Buffer<T> {
fn default() -> Self {
Self::new()
}
}
impl<T> Deref for Buffer<T> {
type Target = [T];
#[inline]
fn deref(&self) -> &[T] {
self.as_slice()
}
}
impl<T> From<Vec<T>> for Buffer<T> {
fn from(values: Vec<T>) -> Self {
Self::from_vec(values)
}
}
impl<T> FromIterator<T> for Buffer<T> {
fn from_iter<I: IntoIterator<Item = T>>(iter: I) -> Self {
Self::from_vec(iter.into_iter().collect())
}
}
impl<T: Clone> Extend<T> for Buffer<T> {
fn extend<I: IntoIterator<Item = T>>(&mut self, iter: I) {
self.to_mut().extend(iter);
}
}
impl<T: Clone> IntoIterator for Buffer<T> {
type Item = T;
type IntoIter = std::vec::IntoIter<T>;
fn into_iter(self) -> Self::IntoIter {
self.into_vec().into_iter()
}
}
#[cfg(test)]
mod tests {
use std::sync::Arc;
use super::{Buffer, Pin};
#[test]
fn a_buffer_reads_back_as_a_slice() {
let buffer: Buffer<i32> = vec![1, 2, 3].into();
assert_eq!(buffer.as_slice(), &[1, 2, 3]);
assert_eq!(buffer.len(), 3);
assert_eq!(buffer[1], 2);
assert_eq!(buffer.iter().sum::<i32>(), 6);
assert_eq!(buffer.clone().into_vec(), vec![1, 2, 3]);
}
#[test]
fn writing_goes_through_one_function() {
let mut buffer = Buffer::with_capacity(4);
buffer.push(1u8);
buffer.extend_from_slice(&[2, 3]);
buffer.extend([4u8]);
buffer.to_mut().sort_unstable_by(|a, b| b.cmp(a));
assert_eq!(buffer.as_slice(), &[4, 3, 2, 1]);
}
#[test]
fn an_empty_buffer_is_the_default_and_collects_like_a_vector() {
assert!(Buffer::<u64>::default().is_empty());
assert!(Buffer::<u64>::new().is_empty());
let collected: Buffer<u64> = (0..4).collect();
assert_eq!(collected.as_slice(), &[0, 1, 2, 3]);
assert_eq!(collected.into_iter().count(), 4);
}
#[test]
fn a_shared_page_reads_like_an_owned_run_and_compares_equal_to_one() {
let page = Arc::new(vec![1u8, 2, 3]);
let buffer = Buffer::from_arc(Arc::clone(&page));
assert!(buffer.is_shared());
assert_eq!(buffer.as_slice(), &[1, 2, 3]);
assert_eq!(buffer[2], 3);
assert_eq!(buffer, Buffer::from_vec(vec![1u8, 2, 3]));
assert_eq!(Buffer::from_vec(vec![1u8, 2, 3]), buffer);
assert_ne!(buffer, Buffer::from_vec(vec![1u8, 2]));
}
#[test]
fn the_producer_gets_the_page_back_once_the_last_buffer_over_it_is_gone() {
let mut page = Arc::new(vec![0u8; 64]);
let address = page.as_ptr();
let first = Buffer::from_arc(Arc::clone(&page));
let second = Buffer::from_arc(Arc::clone(&page));
assert!(Arc::get_mut(&mut page).is_none());
drop(first);
assert!(Arc::get_mut(&mut page).is_none());
drop(second);
let run = Arc::get_mut(&mut page).expect("the last handle");
assert_eq!(run.as_ptr(), address, "the page was reallocated rather than reused");
}
#[test]
fn writing_through_a_shared_page_copies_it_and_leaves_the_page_alone() {
let page = Arc::new(vec![1u8, 2, 3]);
let mut buffer = Buffer::from_arc(Arc::clone(&page));
buffer.push(4);
assert!(!buffer.is_shared());
assert_eq!(buffer.as_slice(), &[1, 2, 3, 4]);
assert_eq!(page.as_slice(), &[1, 2, 3]);
assert_eq!(Buffer::from_arc(Arc::clone(&page)).into_vec(), vec![1, 2, 3]);
}
#[test]
fn the_last_buffer_over_a_page_takes_the_run_without_copying_it() {
let page = Arc::new(vec![5u8; 32]);
let address = page.as_ptr();
let run = Buffer::from_arc(page).into_vec();
assert_eq!(run.as_ptr(), address);
}
#[test]
fn a_shared_page_is_charged_once_across_the_buffers_over_it() {
let page = Arc::new(vec![0u64; 100]);
let over: Vec<_> = (0..4).map(|_| Buffer::from_arc(Arc::clone(&page))).collect();
let charged: usize = over.iter().map(Buffer::footprint).sum();
assert!(
charged <= page.capacity() * 8,
"{charged} charged for a {} byte page",
page.len() * 8
);
assert!(charged > 0);
assert_eq!(Buffer::from_vec(vec![0u64; 100]).footprint(), 800);
}
#[test]
fn a_buffer_and_a_pin_both_cross_a_thread_boundary() {
const fn assert_send<T: Send>() {}
assert_send::<Buffer<i64>>();
assert_send::<Pin>();
let pin: Pin = Arc::new(vec![0u8; 8]);
let buffer: Buffer<i64> = vec![7; 2].into();
let handle = std::thread::spawn(move || (buffer.len(), Arc::strong_count(&pin)));
assert_eq!(handle.join().expect("the thread"), (2, 1));
}
}