use std::sync::{Arc, OnceLock};
use bytes::Bytes;
use ruststream::testing::{Coordinator, TestableBroker};
use ruststream::{
Broker, ConnectedBroker, DefaultPublish, OutgoingMessage, PairError, PublishPolicy, Publisher,
RawMessage, Subscribe,
};
use crate::error::PulsarError;
use crate::testing::router::AddressRouter;
use crate::testing::subscriber::PulsarTestSubscriber;
#[derive(Debug, Default)]
pub(crate) struct TestState {
pub(crate) router: AddressRouter,
coordinator: OnceLock<Coordinator>,
}
impl TestState {
fn coordinator(&self) -> Option<&Coordinator> {
self.coordinator.get()
}
pub(crate) fn publish(&self, name: &str, payload: Bytes, headers: ruststream::Headers) {
self.router
.publish(name, payload, headers, self.coordinator());
}
}
#[derive(Debug, Clone, Default)]
#[must_use]
pub struct PulsarTestBroker {
state: Arc<TestState>,
}
impl PulsarTestBroker {
pub fn new() -> Self {
Self::default()
}
#[must_use]
pub fn publisher(&self) -> PulsarTestPublisher {
PulsarTestPublisher {
state: Arc::clone(&self.state),
}
}
}
impl Broker for PulsarTestBroker {
type Error = PulsarError;
type Connected = ConnectedPulsarTestBroker;
async fn connect(self) -> Result<Self::Connected, Self::Error> {
Ok(ConnectedPulsarTestBroker { state: self.state })
}
}
#[derive(Debug, Clone)]
pub struct ConnectedPulsarTestBroker {
state: Arc<TestState>,
}
impl ConnectedPulsarTestBroker {
#[must_use]
pub fn publisher(&self) -> PulsarTestPublisher {
PulsarTestPublisher {
state: Arc::clone(&self.state),
}
}
}
impl ConnectedBroker for ConnectedPulsarTestBroker {
type Error = PulsarError;
type Closed = ();
async fn shutdown(self) -> Result<(), Self::Error> {
self.state.router.clear();
Ok(())
}
}
impl Subscribe for ConnectedPulsarTestBroker {
type Subscriber = PulsarTestSubscriber;
async fn subscribe(&self, name: &str) -> Result<Self::Subscriber, Self::Error> {
let (id, requeue, rx) = self.state.router.subscribe(name.to_owned());
Ok(PulsarTestSubscriber::new(
Arc::clone(&self.state),
id,
rx,
requeue,
self.state.coordinator().cloned(),
))
}
}
impl TestableBroker for ConnectedPulsarTestBroker {
fn install_coordinator(&self, coordinator: Coordinator) {
let _ = self.state.coordinator.set(coordinator);
}
fn inject(&self, message: OutgoingMessage<'_>) {
self.state.publish(
message.name(),
Bytes::copy_from_slice(message.payload()),
message.headers().clone(),
);
}
fn published(&self, name: &str) -> Vec<RawMessage> {
self.state.router.published(name)
}
}
ruststream::register_testable_broker!(ConnectedPulsarTestBroker);
#[derive(Debug, Clone)]
pub struct PulsarTestPublisher {
state: Arc<TestState>,
}
impl Publisher for PulsarTestPublisher {
type Error = PulsarError;
async fn publish(&self, msg: OutgoingMessage<'_>) -> Result<(), Self::Error> {
self.state.publish(
msg.name(),
Bytes::copy_from_slice(msg.payload()),
msg.headers().clone(),
);
Ok(())
}
}
#[derive(Debug, Clone, Copy, Default)]
#[must_use]
pub struct PulsarTestPublish;
impl PublishPolicy<ConnectedPulsarTestBroker> for PulsarTestPublish {
type Live = PulsarTestPublisher;
async fn pair(self, connected: &ConnectedPulsarTestBroker) -> Result<Self::Live, PairError> {
Ok(connected.publisher())
}
}
impl DefaultPublish for ConnectedPulsarTestBroker {
type Policy = PulsarTestPublish;
}