#![no_std]
extern crate alloc;
use alloc::{sync::Arc, vec::Vec};
use core::mem::ManuallyDrop;
use core::ops::Deref;
use core::sync::atomic::{AtomicUsize, Ordering};
use core::{ptr, slice};
pub struct OnceArray<T> {
data: *mut T,
len: AtomicUsize,
cap: usize,
}
unsafe impl<T> Send for OnceArray<T> where T: Send {}
unsafe impl<T> Sync for OnceArray<T> where T: Sync {}
impl<T> Drop for OnceArray<T> {
fn drop(&mut self) {
unsafe {
drop(Vec::from_raw_parts(
self.data,
*self.len.get_mut(),
self.cap,
))
}
}
}
impl<T> OnceArray<T> {
fn from_vec(v: Vec<T>) -> Self {
let mut v = ManuallyDrop::new(v);
OnceArray {
data: v.as_mut_ptr(),
cap: v.capacity(),
len: AtomicUsize::new(v.len()),
}
}
fn into_vec(self) -> Vec<T> {
unsafe {
let mut v = ManuallyDrop::new(self);
Vec::from_raw_parts(v.data, *v.len.get_mut(), v.cap)
}
}
pub fn capacity(&self) -> usize {
self.cap
}
pub fn len(&self) -> usize {
self.len.load(Ordering::Acquire)
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
pub fn is_full(&self) -> bool {
self.len() == self.cap
}
pub fn as_slice(&self) -> &[T] {
unsafe {
slice::from_raw_parts(self.data, self.len())
}
}
}
impl<T> Deref for OnceArray<T> {
type Target = [T];
fn deref(&self) -> &Self::Target {
self.as_slice()
}
}
impl<T> AsRef<[T]> for OnceArray<T> {
fn as_ref(&self) -> &[T] {
self.as_slice()
}
}
impl<T> core::borrow::Borrow<[T]> for OnceArray<T> {
fn borrow(&self) -> &[T] {
self.as_slice()
}
}
impl<T> From<Vec<T>> for OnceArray<T> {
fn from(val: Vec<T>) -> Self {
OnceArray::from_vec(val)
}
}
impl<T> From<OnceArray<T>> for Vec<T> {
fn from(v: OnceArray<T>) -> Self {
v.into_vec()
}
}
pub struct OnceArrayWriter<T> {
inner: Arc<OnceArray<T>>,
uncommitted_len: usize,
}
impl<T> OnceArrayWriter<T> {
fn from_vec(v: Vec<T>) -> OnceArrayWriter<T> {
Self {
uncommitted_len: v.len(),
inner: Arc::new(OnceArray::from_vec(v)),
}
}
pub fn with_capacity(n: usize) -> OnceArrayWriter<T> {
Self::from_vec(Vec::with_capacity(n))
}
pub fn reader(&self) -> &Arc<OnceArray<T>> {
&self.inner
}
pub fn remaining_capacity(&self) -> usize {
self.inner.cap - self.uncommitted_len
}
pub fn as_slice(&self) -> &[T] {
unsafe {
slice::from_raw_parts(self.inner.data, self.uncommitted_len)
}
}
pub fn uncommitted_mut(&mut self) -> &mut [T] {
unsafe {
let committed_len = self.inner.len.load(Ordering::Relaxed);
slice::from_raw_parts_mut(
self.inner.data.add(committed_len),
self.uncommitted_len - committed_len,
)
}
}
unsafe fn push_unchecked(&mut self, val: T) {
unsafe {
self.inner.data.add(self.uncommitted_len).write(val);
self.uncommitted_len += 1;
}
}
pub fn try_push(&mut self, val: T) -> Result<(), T> {
if self.uncommitted_len < self.inner.cap {
unsafe {
self.push_unchecked(val);
}
Ok(())
} else {
Err(val)
}
}
pub fn extend<I: IntoIterator<Item = T>>(&mut self, iter: I) -> Result<(), I::IntoIter> {
let mut iter = iter.into_iter();
while self.uncommitted_len < self.inner.cap {
if let Some(val) = iter.next() {
unsafe {
self.push_unchecked(val);
}
} else {
return Ok(());
}
}
Err(iter)
}
pub fn extend_from_slice<'a>(&mut self, src: &'a [T]) -> &'a [T]
where
T: Copy,
{
let count = self.remaining_capacity().min(src.len());
unsafe {
self.inner
.data
.add(self.uncommitted_len)
.copy_from_nonoverlapping(src.as_ptr(), count);
}
self.uncommitted_len += count;
&src[count..]
}
pub fn commit(&mut self) {
self.inner
.len
.store(self.uncommitted_len, Ordering::Release);
}
pub fn commit_partial(&mut self, n: usize) {
let committed_len = self.inner.len.load(Ordering::Relaxed);
assert!(
n <= self.uncommitted_len - committed_len,
"Cannot commit more elements than have been initialized"
);
self.inner.len.store(committed_len + n, Ordering::Release);
}
pub fn revert(&mut self) {
let committed_len = self.inner.len.load(Ordering::Relaxed);
let uncommitted_len = self.uncommitted_len;
self.uncommitted_len = committed_len;
unsafe {
ptr::drop_in_place(ptr::slice_from_raw_parts_mut(
self.inner.data.add(committed_len),
uncommitted_len - committed_len,
));
}
}
}
impl<T> Drop for OnceArrayWriter<T> {
fn drop(&mut self) {
self.revert();
}
}
impl<T> From<Vec<T>> for OnceArrayWriter<T> {
fn from(vec: Vec<T>) -> OnceArrayWriter<T> {
OnceArrayWriter::from_vec(vec)
}
}
#[test]
fn test_to_from_vec() {
let v = OnceArray::from(alloc::vec![1, 2, 3]);
assert_eq!(v.as_slice(), &[1, 2, 3]);
let v = Vec::from(v);
assert_eq!(v.as_slice(), &[1, 2, 3]);
}
#[test]
fn test_push() {
let mut writer = OnceArrayWriter::with_capacity(4);
let reader = writer.reader().clone();
assert_eq!(reader.capacity(), 4);
assert_eq!(reader.len(), 0);
assert_eq!(writer.try_push(1), Ok(()));
assert_eq!(reader.len(), 0);
writer.commit();
assert_eq!(reader.len(), 1);
assert_eq!(reader.as_slice(), &[1]);
assert_eq!(writer.try_push(2), Ok(()));
assert_eq!(writer.try_push(3), Ok(()));
assert_eq!(writer.try_push(4), Ok(()));
assert_eq!(writer.try_push(5), Err(5));
writer.commit();
assert_eq!(reader.len(), 4);
assert_eq!(reader.as_slice(), &[1, 2, 3, 4]);
}
#[test]
fn test_extend_from_slice() {
let mut writer = OnceArrayWriter::with_capacity(4);
let reader = writer.reader().clone();
assert_eq!(reader.capacity(), 4);
assert_eq!(reader.len(), 0);
assert_eq!(writer.extend_from_slice(&[1, 2]), &[]);
assert_eq!(reader.len(), 0);
writer.commit();
assert_eq!(reader.len(), 2);
assert_eq!(reader.as_slice(), &[1, 2]);
assert_eq!(writer.extend_from_slice(&[3, 4, 5, 6]), &[5, 6]);
writer.commit();
assert_eq!(reader.len(), 4);
assert_eq!(reader.as_slice(), &[1, 2, 3, 4]);
}
#[test]
fn test_commit_revert() {
let mut writer = OnceArrayWriter::with_capacity(4);
let reader = writer.reader().clone();
assert_eq!(writer.try_push(1), Ok(()));
assert_eq!(writer.try_push(2), Ok(()));
assert_eq!(writer.as_slice(), &[1, 2]);
assert_eq!(writer.uncommitted_mut(), &mut [1, 2]);
writer.commit();
assert_eq!(reader.as_slice(), &[1, 2]);
assert_eq!(writer.uncommitted_mut(), &mut []);
assert_eq!(writer.try_push(3), Ok(()));
assert_eq!(writer.try_push(4), Ok(()));
writer.revert();
assert_eq!(reader.as_slice(), &[1, 2]);
assert_eq!(writer.uncommitted_mut(), &mut []);
assert_eq!(writer.try_push(5), Ok(()));
assert_eq!(writer.try_push(6), Ok(()));
assert_eq!(writer.as_slice(), &[1, 2, 5, 6]);
writer.commit_partial(1);
assert_eq!(reader.as_slice(), &[1, 2, 5]);
assert_eq!(writer.uncommitted_mut(), &[6]);
drop(writer);
assert_eq!(reader.as_slice(), &[1, 2, 5]);
}
#[test]
#[should_panic(expected = "Cannot commit more elements than have been initialized")]
fn test_commit_partial_panic() {
let mut writer = OnceArrayWriter::with_capacity(4);
assert_eq!(writer.try_push(1), Ok(()));
writer.commit_partial(2);
}
#[test]
fn test_extend() {
let mut writer = OnceArrayWriter::with_capacity(4);
let reader = writer.reader().clone();
assert!(writer.extend([1, 2, 3]).is_ok());
assert_eq!(writer.as_slice(), &[1, 2, 3]);
writer.commit();
assert_eq!(reader.as_slice(), &[1, 2, 3]);
let mut remainder = writer.extend([4, 5]).unwrap_err();
assert_eq!(writer.as_slice(), &[1, 2, 3, 4]);
assert_eq!(remainder.next(), Some(5));
}
#[test]
fn test_drop() {
struct DropCounter<'a>(&'a AtomicUsize);
impl<'a> Drop for DropCounter<'a> {
fn drop(&mut self) {
self.0.fetch_add(1, Ordering::Relaxed);
}
}
let drop_count = &AtomicUsize::new(0);
let mut writer = OnceArrayWriter::with_capacity(4);
let reader = writer.reader().clone();
assert!(writer.try_push(DropCounter(drop_count)).is_ok());
assert!(writer.try_push(DropCounter(drop_count)).is_ok());
writer.commit();
assert!(writer.try_push(DropCounter(drop_count)).is_ok());
writer.revert();
assert_eq!(drop_count.load(Ordering::Relaxed), 1);
assert!(writer.try_push(DropCounter(drop_count)).is_ok());
drop(writer);
assert_eq!(drop_count.load(Ordering::Relaxed), 2);
drop(reader);
assert_eq!(drop_count.load(Ordering::Relaxed), 4);
}
#[test]
fn test_concurrent_read() {
extern crate std;
use std::thread;
let mut writer = OnceArrayWriter::<usize>::with_capacity(1024);
let reader = writer.reader().clone();
let handle = thread::spawn(move || {
while reader.len() < 1024 {
let slice = reader.as_slice();
for (i, &v) in slice.iter().enumerate() {
assert_eq!(v, i);
}
}
});
for i in 0..1024 {
writer.try_push(i).unwrap();
writer.commit();
}
handle.join().unwrap();
}