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 { page: Arc<Vec<T>>, from: usize, len: usize },
}
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 {
let len = page.len();
Self { store: Store::Shared { page, from: 0, len } }
}
#[must_use]
pub fn window(page: Arc<Vec<T>>, from: usize, len: usize) -> Self {
let from = from.min(page.len());
let len = len.min(page.len() - from);
Self { store: Store::Shared { page, from, len } }
}
#[must_use]
pub fn is_shared(&self) -> bool {
matches!(self.store, Store::Shared { .. })
}
#[must_use]
pub fn into_page(self) -> Self {
match self.store {
Store::Owned(values) => Self::from_arc(Arc::new(values)),
shared @ Store::Shared { .. } => Self { 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, from, len } => &page[*from..*from + *len],
}
}
}
impl<T: Clone> Buffer<T> {
#[must_use]
pub fn into_vec(self) -> Vec<T> {
match self.store {
Store::Owned(values) => values,
Store::Shared { page, from, len } if from == 0 && len == page.len() => {
Arc::try_unwrap(page).unwrap_or_else(|page| page.to_vec())
}
Store::Shared { page, from, len } => page[from..from + len].to_vec(),
}
}
#[must_use]
pub fn slice(&self, from: usize, len: usize) -> Self {
match &self.store {
Store::Owned(values) => {
let from = from.min(values.len());
let len = len.min(values.len() - from);
Self::from_vec(values[from..from + len].to_vec())
}
Store::Shared { page, from: start, len: held } => {
let from = from.min(*held);
let len = len.min(*held - from);
Self { store: Store::Shared { page: Arc::clone(page), from: start + from, len } }
}
}
}
#[inline]
pub fn to_mut(&mut self) -> &mut Vec<T> {
if let Store::Shared { page, from, len } = &self.store {
self.store = Store::Owned(page[*from..*from + *len].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_window_reads_as_its_own_run() {
let page = Arc::new((0u8..10).collect::<Vec<_>>());
let window = Buffer::window(Arc::clone(&page), 3, 4);
assert!(window.is_shared());
assert_eq!(window.as_slice(), &[3, 4, 5, 6]);
assert_eq!(window.len(), 4);
assert_eq!(window[0], 3);
assert_eq!(window, Buffer::from_vec(vec![3u8, 4, 5, 6]));
assert_eq!(window.iter().sum::<u8>(), 18);
assert_eq!(window.clone().into_vec(), vec![3, 4, 5, 6]);
}
#[test]
fn cutting_a_shared_buffer_points_into_the_same_page() {
let page = Arc::new((0u64..100).collect::<Vec<_>>());
let address = page.as_ptr() as usize;
let whole = Buffer::from_arc(Arc::clone(&page));
let run = whole.slice(64, 16);
assert!(run.is_shared());
assert_eq!(run.as_slice().as_ptr() as usize, address + 64 * 8);
assert_eq!(run.as_slice(), &(64u64..80).collect::<Vec<_>>()[..]);
let inner = run.slice(4, 2);
assert_eq!(inner.as_slice().as_ptr() as usize, address + 68 * 8);
assert_eq!(inner.as_slice(), &[68, 69]);
}
#[test]
fn cutting_an_owned_run_copies_and_a_page_is_how_a_producer_avoids_that() {
let owned = Buffer::from_vec((0u32..16).collect());
let run = owned.slice(4, 4);
assert!(!run.is_shared());
assert_eq!(run.as_slice(), &[4, 5, 6, 7]);
let page = Buffer::from_vec((0u32..16).collect()).into_page();
let address = page.as_slice().as_ptr() as usize;
assert!(page.is_shared(), "into_page did not share the run");
let run = page.slice(4, 4);
assert!(run.is_shared());
assert_eq!(run.as_slice().as_ptr() as usize, address + 4 * 4);
assert_eq!(run.as_slice(), &[4, 5, 6, 7]);
let again = run.into_page();
assert_eq!(again.as_slice().as_ptr() as usize, address + 4 * 4);
}
#[test]
fn a_cut_past_the_end_comes_back_short_rather_than_panicking() {
let page = Arc::new(vec![1u16, 2, 3, 4]);
let whole = Buffer::from_arc(Arc::clone(&page));
assert_eq!(whole.slice(2, 10).as_slice(), &[3, 4]);
assert!(whole.slice(9, 1).is_empty());
assert_eq!(Buffer::window(Arc::clone(&page), 3, 9).as_slice(), &[4]);
assert!(Buffer::window(Arc::clone(&page), 7, 2).is_empty());
let owned = Buffer::from_vec(vec![1u16, 2, 3, 4]);
assert_eq!(owned.slice(2, 10).as_slice(), &[3, 4]);
assert!(owned.slice(9, 1).is_empty());
let middle = Buffer::window(Arc::clone(&page), 1, 2);
assert_eq!(middle.slice(0, 10).as_slice(), &[2, 3]);
}
#[test]
fn writing_through_a_window_copies_the_window_and_leaves_the_page_alone() {
let page = Arc::new(vec![1u8, 2, 3, 4, 5]);
let mut window = Buffer::window(Arc::clone(&page), 1, 3);
window.push(9);
assert!(!window.is_shared());
assert_eq!(window.as_slice(), &[2, 3, 4, 9]);
assert_eq!(page.as_slice(), &[1, 2, 3, 4, 5]);
}
#[test]
fn taking_the_run_out_of_a_window_copies_and_out_of_a_whole_page_does_not() {
let page = Arc::new(vec![7u8; 32]);
let address = page.as_ptr();
assert_eq!(Buffer::window(Arc::clone(&page), 8, 4).into_vec(), vec![7u8; 4]);
let whole = Buffer::from_arc(page).into_vec();
assert_eq!(whole.as_ptr(), address);
}
#[test]
fn a_window_is_charged_for_the_page_it_holds_down() {
let page = Arc::new(vec![0u64; 100]);
let windows: Vec<_> =
(0..4).map(|n| Buffer::window(Arc::clone(&page), n * 25, 25)).collect();
let charged: usize = windows.iter().map(Buffer::footprint).sum();
assert!(charged <= page.capacity() * 8, "{charged} charged for a {} byte page", 800);
assert!(charged > 100 * 8 / 4, "a window was charged less than its own share of the page");
}
#[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));
}
}