use std::sync::Arc;
use crate::bus::{Bus, Message, MessageKind, MessagePublisher, TransportError};
pub struct BusPublisher<B> {
bus: Arc<B>,
}
impl<B> BusPublisher<B> {
pub fn new(bus: Arc<B>) -> Self {
Self { bus }
}
pub fn bus(&self) -> &Arc<B> {
&self.bus
}
}
impl<B> Clone for BusPublisher<B> {
fn clone(&self) -> Self {
Self {
bus: Arc::clone(&self.bus),
}
}
}
impl<B: Bus> MessagePublisher for BusPublisher<B> {
async fn publish(&self, message: Message) -> Result<(), TransportError> {
match message.kind {
MessageKind::Command => self.bus.send_message(message).await,
MessageKind::Event => self.bus.publish_message(message).await,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::future::Future;
use std::sync::Mutex;
fn block_on<F: Future>(future: F) -> F::Output {
use std::ptr;
use std::task::{Context, Poll, RawWaker, RawWakerVTable, Waker};
const VTABLE: RawWakerVTable = RawWakerVTable::new(
|_| RawWaker::new(ptr::null(), &VTABLE),
|_| {},
|_| {},
|_| {},
);
let waker = unsafe { Waker::from_raw(RawWaker::new(ptr::null(), &VTABLE)) };
let mut cx = Context::from_waker(&waker);
let mut future = std::pin::pin!(future);
loop {
if let Poll::Ready(output) = future.as_mut().poll(&mut cx) {
return output;
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum Call {
Send(String),
Publish(String),
}
#[derive(Default)]
struct RecordingBus {
calls: Mutex<Vec<Call>>,
}
impl RecordingBus {
fn calls(&self) -> Vec<Call> {
self.calls.lock().unwrap().clone()
}
}
impl Bus for RecordingBus {
async fn send(&self, name: &str, payload: Vec<u8>) -> Result<(), TransportError> {
self.send_message(Message::new(name, MessageKind::Command, payload))
.await
}
async fn publish(&self, name: &str, payload: Vec<u8>) -> Result<(), TransportError> {
self.publish_message(Message::new(name, MessageKind::Event, payload))
.await
}
async fn send_message(&self, message: Message) -> Result<(), TransportError> {
self.calls.lock().unwrap().push(Call::Send(message.name));
Ok(())
}
async fn publish_message(&self, message: Message) -> Result<(), TransportError> {
self.calls.lock().unwrap().push(Call::Publish(message.name));
Ok(())
}
}
#[test]
fn routes_command_to_send_and_event_to_publish() {
let bus = Arc::new(RecordingBus::default());
let publisher = BusPublisher::new(bus.clone());
block_on(publisher.publish(Message::new(
"ship.order",
MessageKind::Command,
b"{}".to_vec(),
)))
.unwrap();
block_on(publisher.publish(Message::new(
"order.shipped",
MessageKind::Event,
b"{}".to_vec(),
)))
.unwrap();
assert_eq!(
bus.calls(),
vec![
Call::Send("ship.order".to_string()),
Call::Publish("order.shipped".to_string()),
]
);
}
}