use std::any::type_name;
use anyhow::Result;
use futures::{StreamExt, pin_mut};
use crate::{
engine::MessageBusEngine,
prelude::*,
view::{Query, View, Viewer},
};
pub struct MessageBus<D: MessageBusDriver> {
engine: MessageBusEngine<D>,
}
impl<D: MessageBusDriver> Clone for MessageBus<D> {
fn clone(&self) -> Self {
Self {
engine: self.engine.clone(),
}
}
}
impl<D> From<&D> for MessageBus<D>
where
D: MessageBusDriver,
D::Broker: for<'a> From<&'a D>,
D::Projector: for<'a> From<&'a D>,
D::Handler: for<'a> From<&'a D>,
D::Policy: for<'a> From<&'a D>,
D::Viewer: for<'a> From<&'a D>,
<D::PolicyContext as PolicyContext>::Factory: for<'a> From<&'a D>,
<D::UnitOfWork as UnitOfWork>::Factory: for<'a> From<&'a D>,
{
fn from(driver: &D) -> Self {
let engine = MessageBusEngine::from(driver);
Self { engine }
}
}
impl<D: MessageBusDriver> MessageBus<D> {
pub async fn dispatch<C: Command>(&self, cmd: C) -> Result<Option<D::Identifier>>
where
D::Handler: CommandHandler<C, D>,
{
println!("User provided command: {}", type_name::<C>());
let mut uow = self.engine.uow_factory.create().await?;
match self.engine.handler.handle(&mut uow, cmd).await {
Ok(res) => {
let events = uow
.commit()
.await?
.into_iter()
.map(DriverMessage::<D>::Event)
.collect();
self.engine.broker.publish_batch(events).await?;
Ok(res)
}
Err(e) => {
uow.rollback().await?;
Err(e)
}
}
}
pub async fn view<Q: Query>(&self, query: Q) -> Result<impl View>
where
D::Viewer: Viewer<Q>,
{
self.engine.viewer.view(query).await
}
pub async fn start(self) -> Result<()>
where
D::Handler: CommandHandler<D::Command, D>,
D::Policy: Policy<D::Event, D, Output = SideEffect<D::Command, D::Projection>>,
{
let stream = self.engine.broker.receiver();
pin_mut!(stream);
while let Some((id, msg)) = stream.next().await {
match self.handle_message(msg).await {
Ok(_) => {
self.engine.broker.ack(id).await?;
println!("Handled message successfully.");
}
Err(e) => {
self.engine.broker.nack(id).await?;
println!("Handled message unsuccessfully: {e:#?}");
}
};
}
Ok(())
}
async fn handle_message(&self, msg: DriverMessage<D>) -> Result<()>
where
D::Handler: CommandHandler<D::Command, D>,
D::Policy: Policy<D::Event, D, Output = SideEffect<D::Command, D::Projection>>,
{
match msg {
Message::Command(cmd) => {
println!("Executing command");
self.dispatch(cmd).await?;
}
Message::Event(event) => {
println!("Executing event");
self.handle_event(event).await?;
}
Message::Projection(projection) => {
println!("Executing projection");
self.engine.projector.project(projection).await?;
}
};
Ok(())
}
async fn handle_event(&self, event: D::Event) -> Result<()>
where
D::Policy: Policy<D::Event, D, Output = SideEffect<D::Command, D::Projection>>,
{
let mut ctx = self.engine.policy_context_factory.create().await?;
let res = match self.engine.policy.apply(&mut ctx, event).await {
Ok(events) => {
let messages = events
.into_iter()
.map(|side_effect| match side_effect {
SideEffect::Command(cmd) => Message::Command(cmd),
SideEffect::Projection(proj) => Message::Projection(proj),
})
.collect::<Vec<_>>();
let num_events = messages.len();
self.engine.broker.publish_batch(messages).await?;
println!("Published {num_events} events.");
Ok(())
}
Err(e) => Err(e),
};
ctx.close().await?;
res
}
}