use super::journal_error::JournalError;
use super::journal_owner::JournalOwner;
use super::journal_reader::JournalReader;
use crate::event::event_envelope::EventEnvelope;
use crate::event::types::EventId;
use crate::event::vector_clock::CausalOrderingService;
use crate::event::JournalEvent;
use crate::id::JournalId;
use async_trait::async_trait;
#[async_trait]
pub trait Journal<T>: Send + Sync
where
T: JournalEvent,
{
fn id(&self) -> &JournalId;
fn owner(&self) -> Option<&JournalOwner>;
async fn append(
&self,
event: T,
parent: Option<&EventEnvelope<T>>,
) -> Result<EventEnvelope<T>, JournalError>;
async fn append_group(
&self,
group_id: &str,
events: Vec<T>,
parent: Option<&EventEnvelope<T>>,
) -> Result<Vec<EventEnvelope<T>>, JournalError> {
match events.len() {
0 => Ok(Vec::new()),
1 => Ok(vec![
self.append(events.into_iter().next().expect("one event"), parent)
.await?,
]),
count => Err(JournalError::AtomicAppendUnsupported {
group_id: group_id.to_string(),
member_count: count,
}),
}
}
async fn read_all_unordered(&self) -> Result<Vec<EventEnvelope<T>>, JournalError>;
async fn read_causally_ordered(&self) -> Result<Vec<EventEnvelope<T>>, JournalError> {
CausalOrderingService::order_envelopes_by_event_id(self.read_all_unordered().await?)
}
async fn read_causally_after(
&self,
after_event_id: &EventId,
) -> Result<Vec<EventEnvelope<T>>, JournalError> {
let all = self.read_causally_ordered().await?;
Ok(
match all.iter().position(|e| e.event.id() == after_event_id) {
Some(pos) => all.into_iter().skip(pos + 1).collect(),
None => Vec::new(),
},
)
}
async fn read_event(
&self,
event_id: &EventId,
) -> Result<Option<EventEnvelope<T>>, JournalError>;
async fn reader(&self) -> Result<Box<dyn JournalReader<T>>, JournalError> {
self.reader_from(0).await
}
async fn reader_from(&self, position: u64) -> Result<Box<dyn JournalReader<T>>, JournalError>;
async fn read_last_n(&self, count: usize) -> Result<Vec<EventEnvelope<T>>, JournalError>;
}