use core::iter;
use core::num::NonZeroU32;
use mnesis::{DomainEvent, Id, Version};
use crate::decoded::Decoded;
use crate::state::{Hydrated, PersistTrigger, SnapshotStore};
pub trait Projector: Send + Sync + 'static {
type Event: DomainEvent;
type State: Send + Sync + 'static;
type Error: core::error::Error + Send + Sync + 'static;
fn initial(&self) -> Self::State;
fn apply(&self, state: Self::State, event: &Self::Event) -> Result<Self::State, Self::Error>;
}
pub struct Projection<I, P: Projector, Trig, SS> {
id: I,
projector: P,
trigger: Trig,
snapshot_store: SS,
schema_version: NonZeroU32,
checkpoint: Option<Version>,
pending: Option<Version>,
rebuilt_from: Option<NonZeroU32>,
}
impl<I, P, Trig, SS> Projection<I, P, Trig, SS>
where
I: Id,
P: Projector,
Trig: PersistTrigger,
SS: SnapshotStore<P::State, Version>,
{
pub async fn load(
id: I,
projector: P,
trigger: Trig,
snapshot_store: SS,
schema_version: NonZeroU32,
) -> Result<(Self, P::State), SS::Error> {
let (state, checkpoint, rebuilt_from) = match snapshot_store
.hydrate(&id, schema_version)
.await?
{
Hydrated::Found { position, state } => (state, Some(position), None),
Hydrated::Absent => (projector.initial(), None, None),
Hydrated::Stale { stored_schema } => (projector.initial(), None, Some(stored_schema)),
};
Ok((
Self {
id,
projector,
trigger,
snapshot_store,
schema_version,
checkpoint,
pending: None,
rebuilt_from,
},
state,
))
}
#[must_use]
pub const fn rebuilding_from(&self) -> Option<NonZeroU32> {
self.rebuilt_from
}
pub const fn id(&self) -> &I {
&self.id
}
pub const fn checkpoint(&self) -> Option<Version> {
self.checkpoint
}
pub async fn advance(
&mut self,
state: P::State,
decoded: Decoded<P::Event>,
) -> Result<P::State, ProjectionError<P::Error, SS::Error>> {
let position = decoded.version;
let folded = self
.projector
.apply(state, &decoded.event)
.map_err(ProjectionError::Apply)?;
if self
.trigger
.should_persist(self.checkpoint, position, iter::once(decoded.event.name()))
{
self.commit(position, &folded).await?;
} else {
self.pending = Some(position);
}
Ok(folded)
}
pub async fn flush(
&mut self,
state: &P::State,
) -> Result<(), ProjectionError<P::Error, SS::Error>> {
match self.pending {
Some(position) => self.commit(position, state).await,
None => Ok(()),
}
}
async fn commit(
&mut self,
position: Version,
state: &P::State,
) -> Result<(), ProjectionError<P::Error, SS::Error>> {
self.snapshot_store
.commit(&self.id, self.schema_version, position, state)
.await
.map_err(ProjectionError::Commit)?;
self.checkpoint = Some(position);
self.pending = None;
Ok(())
}
}
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum ProjectionError<PErr, SErr> {
#[error("projector failed to apply event")]
Apply(#[source] PErr),
#[error("snapshot commit failed")]
Commit(#[source] SErr),
}