use crate::error::Error;
use crate::topic::backend::SubscriptionBackend;
use crate::topic::topic::InMemoryDeliveryFinalizer;
#[cfg(test)]
use crate::topic::topic::InMemoryDeliveryKind;
use crate::topic::types::{Envelope, RecvItem};
use std::sync::Arc;
pub struct Subscription<T: Send + Sync + 'static> {
inner: Box<dyn SubscriptionBackend<T>>,
}
pub enum RecvDelivery<T: Send + Sync + 'static> {
Message(Delivery<T>),
Lagged {
missed: u64,
},
}
#[allow(dead_code)]
pub(crate) trait DeliveryBackend<T: Send + Sync + 'static> {
fn envelope(&self) -> &Envelope<T>;
fn commit(&mut self);
fn abort(&mut self, reason: Arc<str>) -> Result<(), Error>;
fn abandon(&mut self);
}
pub struct Delivery<T: Send + Sync + 'static> {
envelope: Envelope<T>,
finalizer: DeliveryFinalizer<T>,
}
enum DeliveryFinalizer<T: Send + Sync + 'static> {
InMemory(InMemoryDeliveryFinalizer),
#[allow(dead_code)]
Opaque(Box<dyn DeliveryBackend<T>>),
Finished,
}
impl<T: Send + Sync + 'static> Delivery<T> {
pub(crate) fn new_in_memory(
envelope: Envelope<T>,
finalizer: InMemoryDeliveryFinalizer,
) -> Self {
Self {
envelope,
finalizer: DeliveryFinalizer::InMemory(finalizer),
}
}
#[allow(dead_code)]
pub(crate) fn new_opaque(inner: Box<dyn DeliveryBackend<T>>) -> Self {
let envelope = inner.envelope().clone();
Self {
envelope,
finalizer: DeliveryFinalizer::Opaque(inner),
}
}
#[must_use]
pub fn envelope(&self) -> &Envelope<T> {
&self.envelope
}
#[must_use]
pub fn message_id(&self) -> u64 {
self.envelope().id
}
#[must_use]
pub fn tracked(&self) -> bool {
self.envelope().tracked
}
pub fn commit(mut self) {
match std::mem::replace(&mut self.finalizer, DeliveryFinalizer::Finished) {
DeliveryFinalizer::InMemory(mut finalizer) => finalizer.commit(),
DeliveryFinalizer::Opaque(mut inner) => inner.commit(),
DeliveryFinalizer::Finished => {}
}
}
pub fn abort(mut self, reason: impl Into<Arc<str>>) -> Result<(), Error> {
let reason = reason.into();
match std::mem::replace(&mut self.finalizer, DeliveryFinalizer::Finished) {
DeliveryFinalizer::InMemory(mut finalizer) => finalizer.abort(&self.envelope, reason),
DeliveryFinalizer::Opaque(mut inner) => inner.abort(reason),
DeliveryFinalizer::Finished => Ok(()),
}
}
#[cfg(test)]
pub(crate) fn storage_kind(&self) -> DeliveryStorageKind {
match &self.finalizer {
DeliveryFinalizer::InMemory(finalizer) => match finalizer.kind() {
InMemoryDeliveryKind::Balanced => DeliveryStorageKind::Balanced,
InMemoryDeliveryKind::Broadcast => DeliveryStorageKind::Broadcast,
},
DeliveryFinalizer::Opaque(_) => DeliveryStorageKind::Opaque,
DeliveryFinalizer::Finished => panic!("finished deliveries should not be inspected"),
}
}
}
impl<T: Send + Sync + 'static> Drop for Delivery<T> {
fn drop(&mut self) {
match std::mem::replace(&mut self.finalizer, DeliveryFinalizer::Finished) {
DeliveryFinalizer::InMemory(mut finalizer) => finalizer.abandon(&self.envelope),
DeliveryFinalizer::Opaque(mut inner) => inner.abandon(),
DeliveryFinalizer::Finished => {}
}
}
}
#[cfg(test)]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum DeliveryStorageKind {
Balanced,
Broadcast,
Opaque,
}
impl<T: Send + Sync + 'static> Subscription<T> {
pub(crate) fn new(inner: Box<dyn SubscriptionBackend<T>>) -> Self {
Self { inner }
}
pub async fn recv(&mut self) -> Result<RecvItem<T>, Error> {
match self.recv_delivery().await? {
RecvDelivery::Message(delivery) => {
let envelope = delivery.envelope().clone();
delivery.commit();
Ok(RecvItem::Message(envelope))
}
RecvDelivery::Lagged { missed } => Ok(RecvItem::Lagged { missed }),
}
}
pub async fn recv_delivery(&mut self) -> Result<RecvDelivery<T>, Error> {
std::future::poll_fn(|cx| self.inner.poll_recv_delivery(cx)).await
}
pub fn ack(&self, id: u64) -> Result<(), Error> {
self.inner.ack(id)
}
pub fn nack(&self, id: u64, reason: impl Into<Arc<str>>) -> Result<(), Error> {
self.inner.nack(id, reason.into())
}
}