use std::marker::PhantomData;
use rkyv::api::high::HighValidator;
use rkyv::bytecheck::CheckBytes;
use rkyv::rancor::{Error, Strategy};
use rkyv::ser::Serializer;
use rkyv::ser::allocator::ArenaHandle;
use rkyv::ser::sharing::Share;
use rkyv::util::AlignedVec;
use rkyv::{Archive, Portable, Serialize};
use crate::codec::CodecError;
use crate::error::{OpenError, PushError};
use crate::queue::{Builder, Consumer, Producer, Reserved};
use crate::store::Store;
use crate::typed::{ReserveError, TypedPushError};
pub trait Archivable:
Archive<Archived: Portable + for<'a> CheckBytes<HighValidator<'a, Error>>>
+ for<'a> Serialize<Strategy<Serializer<AlignedVec, ArenaHandle<'a>, Share>, Error>>
{
}
impl<T> Archivable for T where
T: Archive<Archived: Portable + for<'a> CheckBytes<HighValidator<'a, Error>>>
+ for<'a> Serialize<Strategy<Serializer<AlignedVec, ArenaHandle<'a>, Share>, Error>>
{
}
pub type ArchivedEnds<S, T> = (ArchivedProducer<S, T>, ArchivedConsumer<S, T>);
impl<S: Store> Builder<S> {
pub fn open_archived<T: Archivable>(self) -> Result<ArchivedEnds<S, T>, OpenError<S::Error>> {
let (producer, consumer) = self.open()?;
Ok((
ArchivedProducer {
inner: producer,
_marker: PhantomData,
},
ArchivedConsumer {
inner: consumer,
_marker: PhantomData,
},
))
}
}
pub struct ArchivedProducer<S, T> {
inner: Producer<S>,
_marker: PhantomData<fn(T)>,
}
impl<S, T> Clone for ArchivedProducer<S, T> {
fn clone(&self) -> Self {
Self {
inner: self.inner.clone(),
_marker: PhantomData,
}
}
}
impl<S: Store, T: Archivable> ArchivedProducer<S, T> {
pub fn push(&self, value: &T) -> Result<(), TypedPushError<S::Error>> {
let bytes = rkyv::to_bytes::<Error>(value)
.map_err(|e| TypedPushError::Encode(CodecError::new(e)))?;
self.inner.push(&bytes).map_err(|e| match e {
PushError::Closed => TypedPushError::Closed,
PushError::Store(e) => TypedPushError::Store(e),
})
}
pub fn close(&self) {
self.inner.close();
}
pub fn len(&self) -> usize {
self.inner.len()
}
pub fn is_empty(&self) -> bool {
self.inner.is_empty()
}
}
pub struct ArchivedConsumer<S, T> {
inner: Consumer<S>,
_marker: PhantomData<fn() -> T>,
}
impl<S: Store, T: Archivable> ArchivedConsumer<S, T> {
pub fn reserve(&self) -> Result<Option<ArchivedReserved<S, T>>, ReserveError<S::Error>> {
match self.inner.reserve().map_err(ReserveError::Store)? {
Some(reserved) => {
let mut aligned = AlignedVec::<16>::new();
aligned.extend_from_slice(&reserved);
rkyv::access::<T::Archived, Error>(&aligned)
.map_err(|e| ReserveError::Decode(CodecError::new(e)))?;
Ok(Some(ArchivedReserved {
inner: reserved,
aligned,
_marker: PhantomData,
}))
}
None => Ok(None),
}
}
}
pub struct ArchivedReserved<S: Store, T> {
inner: Reserved<S>,
aligned: AlignedVec,
_marker: PhantomData<fn() -> T>,
}
impl<S: Store, T: Archivable> ArchivedReserved<S, T> {
pub fn get(&self) -> &T::Archived {
unsafe { rkyv::access_unchecked::<T::Archived>(&self.aligned) }
}
pub fn seq(&self) -> u64 {
self.inner.seq()
}
pub fn ack(self) -> Result<(), S::Error> {
self.inner.ack()
}
pub fn nack(self) {
self.inner.nack();
}
}