use ruststream::{PairError, PublishPolicy};
use crate::broker::ConnectedLapinBroker;
use crate::publisher::{ConfirmsPublisher, LapinPublisher, ServerTxPublisher};
use self::sealed::Sealed;
mod sealed {
pub trait Sealed {}
impl Sealed for super::LapinPublish {}
impl Sealed for super::ConfirmsPublish {}
impl Sealed for super::ServerTxPublish {}
impl Sealed for crate::requester::LapinRequest {}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct PublishOptions {
pub(crate) exchange: String,
pub(crate) persistent: bool,
}
impl Default for PublishOptions {
fn default() -> Self {
Self {
exchange: String::new(),
persistent: true,
}
}
}
pub trait LapinPublishPolicy: PublishPolicy<ConnectedLapinBroker> + Sealed {
#[must_use]
fn bind(self, connected: &ConnectedLapinBroker) -> Self::Live;
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
#[must_use]
pub struct LapinPublish(PublishOptions);
impl LapinPublish {
pub fn exchange(mut self, exchange: impl Into<String>) -> Self {
self.0.exchange = exchange.into();
self
}
pub fn persistent(mut self, persistent: bool) -> Self {
self.0.persistent = persistent;
self
}
pub fn confirms(self) -> ConfirmsPublish {
ConfirmsPublish(self.0)
}
pub fn server_tx(self) -> ServerTxPublish {
ServerTxPublish(self.0)
}
}
impl PublishPolicy<ConnectedLapinBroker> for LapinPublish {
type Live = LapinPublisher;
async fn pair(self, connected: &ConnectedLapinBroker) -> Result<Self::Live, PairError> {
Ok(self.bind(connected))
}
}
impl LapinPublishPolicy for LapinPublish {
fn bind(self, connected: &ConnectedLapinBroker) -> Self::Live {
LapinPublisher::new(connected, self.0)
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
#[must_use]
pub struct ConfirmsPublish(PublishOptions);
impl ConfirmsPublish {
pub fn exchange(mut self, exchange: impl Into<String>) -> Self {
self.0.exchange = exchange.into();
self
}
pub fn persistent(mut self, persistent: bool) -> Self {
self.0.persistent = persistent;
self
}
}
impl PublishPolicy<ConnectedLapinBroker> for ConfirmsPublish {
type Live = ConfirmsPublisher;
async fn pair(self, connected: &ConnectedLapinBroker) -> Result<Self::Live, PairError> {
Ok(self.bind(connected))
}
}
impl LapinPublishPolicy for ConfirmsPublish {
fn bind(self, connected: &ConnectedLapinBroker) -> Self::Live {
ConfirmsPublisher::new(connected, self.0)
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
#[must_use]
pub struct ServerTxPublish(PublishOptions);
impl ServerTxPublish {
pub fn exchange(mut self, exchange: impl Into<String>) -> Self {
self.0.exchange = exchange.into();
self
}
pub fn persistent(mut self, persistent: bool) -> Self {
self.0.persistent = persistent;
self
}
}
impl PublishPolicy<ConnectedLapinBroker> for ServerTxPublish {
type Live = ServerTxPublisher;
async fn pair(self, connected: &ConnectedLapinBroker) -> Result<Self::Live, PairError> {
Ok(self.bind(connected))
}
}
impl LapinPublishPolicy for ServerTxPublish {
fn bind(self, connected: &ConnectedLapinBroker) -> Self::Live {
ServerTxPublisher::new(connected, self.0)
}
}