use std::error::Error as StdError;
use std::fmt;
use std::marker::PhantomData;
use std::ops::Deref;
use crate::codec::{Codec, CodecError};
use crate::error::{OpenError, PushError};
use crate::queue::{Builder, Consumer, Producer, Reserved};
use crate::store::Store;
pub type TypedEnds<S, T, C> = (TypedProducer<S, T, C>, TypedConsumer<S, T, C>);
impl<S: Store> Builder<S> {
pub fn open_typed<T, C>(self, codec: C) -> Result<TypedEnds<S, T, C>, OpenError<S::Error>>
where
C: Codec<T> + Clone,
{
let (producer, consumer) = self.open()?;
Ok((
TypedProducer {
inner: producer,
codec: codec.clone(),
_marker: PhantomData,
},
TypedConsumer {
inner: consumer,
codec,
_marker: PhantomData,
},
))
}
}
pub struct TypedProducer<S, T, C> {
inner: Producer<S>,
codec: C,
_marker: PhantomData<fn(T)>,
}
impl<S, T, C: Clone> Clone for TypedProducer<S, T, C> {
fn clone(&self) -> Self {
Self {
inner: self.inner.clone(),
codec: self.codec.clone(),
_marker: PhantomData,
}
}
}
impl<S: Store, T, C: Codec<T>> TypedProducer<S, T, C> {
pub fn push(&self, value: &T) -> Result<(), TypedPushError<S::Error>> {
let bytes = self.codec.encode(value).map_err(TypedPushError::Encode)?;
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 TypedConsumer<S, T, C> {
inner: Consumer<S>,
codec: C,
_marker: PhantomData<fn() -> T>,
}
impl<S: Store, T, C: Codec<T>> TypedConsumer<S, T, C> {
pub fn reserve(&self) -> Result<Option<TypedReserved<S, T>>, ReserveError<S::Error>> {
match self.inner.reserve().map_err(ReserveError::Store)? {
Some(reserved) => {
let value = self.codec.decode(&reserved).map_err(ReserveError::Decode)?;
Ok(Some(TypedReserved {
inner: reserved,
value,
}))
}
None => Ok(None),
}
}
}
pub struct TypedReserved<S: Store, T> {
inner: Reserved<S>,
value: T,
}
impl<S: Store, T> TypedReserved<S, T> {
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();
}
}
impl<S: Store, T> Deref for TypedReserved<S, T> {
type Target = T;
fn deref(&self) -> &T {
&self.value
}
}
#[derive(Debug)]
pub enum TypedPushError<E> {
Encode(CodecError),
Closed,
Store(E),
}
impl<E: fmt::Display> fmt::Display for TypedPushError<E> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
TypedPushError::Encode(e) => write!(f, "{e}"),
TypedPushError::Closed => write!(f, "queue is closed"),
TypedPushError::Store(e) => write!(f, "store error: {e}"),
}
}
}
impl<E: StdError + 'static> StdError for TypedPushError<E> {
fn source(&self) -> Option<&(dyn StdError + 'static)> {
match self {
TypedPushError::Encode(e) => Some(e),
TypedPushError::Store(e) => Some(e),
TypedPushError::Closed => None,
}
}
}
#[derive(Debug)]
pub enum ReserveError<E> {
Store(E),
Decode(CodecError),
}
impl<E: fmt::Display> fmt::Display for ReserveError<E> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
ReserveError::Store(e) => write!(f, "store error: {e}"),
ReserveError::Decode(e) => write!(f, "{e}"),
}
}
}
impl<E: StdError + 'static> StdError for ReserveError<E> {
fn source(&self) -> Option<&(dyn StdError + 'static)> {
match self {
ReserveError::Store(e) => Some(e),
ReserveError::Decode(e) => Some(e),
}
}
}