use std::future::Future;
use crate::entity::{Entity, EventRecord};
use crate::outbox::OutboxMessage;
use crate::read_model::{ReadModelLoadGraph, ReadModelLoadRequest, ReadModelQueryCapabilities};
use crate::snapshot::SnapshotRecord;
use crate::table::{TableAdapterCapabilities, TableCommitOutcome, TableStoreError, TableWritePlan};
use super::inbox::InboxReceipt;
use super::{RepositoryError, StreamIdentity};
pub struct StreamWrite<'a> {
pub identity: StreamIdentity,
pub entity: &'a mut Entity,
}
impl<'a> StreamWrite<'a> {
pub fn new(identity: StreamIdentity, entity: &'a mut Entity) -> Self {
Self { identity, entity }
}
}
#[derive(Clone, Debug)]
pub enum SnapshotWrite {
Save {
identity: StreamIdentity,
record: SnapshotRecord,
},
}
pub struct CommitBatch<'a> {
pub streams: Vec<StreamWrite<'a>>,
pub outbox_messages: Vec<OutboxMessage>,
pub read_model_plans: Vec<TableWritePlan>,
pub snapshots: Vec<SnapshotWrite>,
pub inbox_receipts: Vec<InboxReceipt>,
}
impl<'a> CommitBatch<'a> {
pub fn new(streams: Vec<StreamWrite<'a>>) -> Self {
Self {
streams,
outbox_messages: Vec::new(),
read_model_plans: Vec::new(),
snapshots: Vec::new(),
inbox_receipts: Vec::new(),
}
}
pub fn empty() -> Self {
Self::new(Vec::new())
}
}
#[derive(Clone, Debug)]
pub struct PreparedEventAppend<'a> {
pub identity: StreamIdentity,
pub expected_version: u64,
pub events: &'a [EventRecord],
}
impl<'a> PreparedEventAppend<'a> {
pub fn from_stream_write(write: &'a StreamWrite<'_>) -> Self {
Self {
identity: write.identity.clone(),
expected_version: write.entity.committed_version(),
events: write.entity.new_events(),
}
}
}
pub trait GetStream: Send + Sync {
fn get_stream<'a>(
&'a self,
identity: &'a StreamIdentity,
) -> impl Future<Output = Result<Option<Entity>, RepositoryError>> + Send + 'a;
fn get_streams<'a>(
&'a self,
identities: &'a [StreamIdentity],
) -> impl Future<Output = Result<Vec<Entity>, RepositoryError>> + Send + 'a {
async move {
let mut entities = Vec::with_capacity(identities.len());
for identity in identities {
if let Some(entity) = self.get_stream(identity).await? {
entities.push(entity);
}
}
Ok(entities)
}
}
fn get_stream_tail<'a>(
&'a self,
identity: &'a StreamIdentity,
after_version: u64,
) -> impl Future<Output = Result<Option<Entity>, RepositoryError>> + Send + 'a {
let _ = after_version;
self.get_stream(identity)
}
}
pub trait TransactionalCommit: Send + Sync {
fn commit_batch<'a>(
&'a self,
batch: CommitBatch<'a>,
) -> impl Future<Output = Result<(), RepositoryError>> + Send + 'a;
}
pub trait InboxStore: Send + Sync {
fn inbox_contains<'a>(
&'a self,
consumer: &'a str,
message_id: &'a str,
) -> impl Future<Output = Result<bool, RepositoryError>> + Send + 'a;
fn purge_inbox_older_than(
&self,
age: std::time::Duration,
) -> impl Future<Output = Result<u64, RepositoryError>> + Send;
}
pub trait Repository: GetStream + TransactionalCommit {}
impl<T> Repository for T where T: GetStream + TransactionalCommit {}
pub trait ReadModelWritePlanStore: Send + Sync {
fn read_model_capabilities(&self) -> TableAdapterCapabilities;
fn commit_write_plan(
&self,
plan: TableWritePlan,
) -> impl Future<Output = Result<TableCommitOutcome, TableStoreError>> + Send + '_;
}
pub trait RelationalReadModelQueryStore: Send + Sync {
fn read_model_query_capabilities(&self) -> ReadModelQueryCapabilities;
fn load_graph(
&self,
request: ReadModelLoadRequest,
) -> impl Future<Output = Result<ReadModelLoadGraph, TableStoreError>> + Send + '_;
}
pub trait SnapshotStore: Send + Sync {
fn get_snapshot<'a>(
&'a self,
identity: &'a StreamIdentity,
) -> impl Future<Output = Result<Option<SnapshotRecord>, RepositoryError>> + Send + 'a;
fn get_snapshots<'a>(
&'a self,
identities: &'a [StreamIdentity],
) -> impl Future<Output = Result<Vec<SnapshotRecord>, RepositoryError>> + Send + 'a {
async move {
let mut records = Vec::with_capacity(identities.len());
for identity in identities {
if let Some(record) = self.get_snapshot(identity).await? {
records.push(record);
}
}
Ok(records)
}
}
fn save_snapshot<'a>(
&'a self,
identity: &'a StreamIdentity,
record: SnapshotRecord,
) -> impl Future<Output = Result<(), RepositoryError>> + Send + 'a;
fn delete_snapshot<'a>(
&'a self,
identity: &'a StreamIdentity,
) -> impl Future<Output = Result<bool, RepositoryError>> + Send + 'a;
}