use std::{marker::PhantomData, sync::Arc};
use crate::{
cursor::{Args, ReadResult, Value},
Aggregate, AggregateEvent, AggregateEvents, Event, EventFilter, Executor, RoutingKey,
};
const PAGE_SIZE: usize = 100;
const DEFAULT_PAGE: u16 = 40;
pub fn read<A: Aggregate>(id: impl Into<String>) -> ReadBuilder<A> {
ReadBuilder::from_filter(EventFilter::by_id::<A>(id))
}
pub fn read_raw(aggregate_type: impl Into<String>, id: impl Into<String>) -> ReadBuilder {
ReadBuilder::from_filter(EventFilter::by_id_raw(aggregate_type, id))
}
pub struct ReadBuilder<A = ()> {
filter: EventFilter,
routing_key: Option<RoutingKey>,
limit: Option<u16>,
after: Option<Value>,
before: Option<Value>,
backward: bool,
_aggregate: PhantomData<fn() -> A>,
}
impl<A> ReadBuilder<A> {
fn from_filter(filter: EventFilter) -> Self {
Self {
filter,
routing_key: None,
limit: None,
after: None,
before: None,
backward: false,
_aggregate: PhantomData,
}
}
pub fn event<EV: AggregateEvent>(mut self) -> Self {
self.filter.name = Some(EV::event_name().to_owned());
self
}
pub fn event_raw(mut self, name: impl Into<String>) -> Self {
self.filter.name = Some(name.into());
self
}
pub fn routing_key(mut self, v: impl Into<String>) -> Self {
self.routing_key = Some(RoutingKey::Value(Some(v.into())));
self
}
pub fn no_routing_key(mut self) -> Self {
self.routing_key = Some(RoutingKey::Value(None));
self
}
pub fn limit(mut self, v: u16) -> Self {
self.limit = Some(v);
self
}
pub fn after(mut self, cursor: Value) -> Self {
self.after = Some(cursor);
self
}
pub fn before(mut self, cursor: Value) -> Self {
self.before = Some(cursor);
self.backward = true;
self
}
pub fn backward(mut self) -> Self {
self.backward = true;
self
}
pub fn args(mut self, args: Args) -> Self {
self.backward = args.is_backward();
self.limit = args.first.or(args.last);
self.after = args.after;
self.before = args.before;
self
}
fn filters(&self) -> Arc<[EventFilter]> {
Arc::from([self.filter.clone()])
}
fn args_for(&self, page_size: u16, cursor: Option<Value>) -> Args {
if self.backward {
Args::backward(page_size, cursor)
} else {
Args::forward(page_size, cursor)
}
}
pub async fn page<E: Executor>(&self, executor: &E) -> anyhow::Result<ReadResult<Event>> {
let cursor = if self.backward {
self.before.clone()
} else {
self.after.clone()
};
let args = self.args_for(self.limit.unwrap_or(DEFAULT_PAGE), cursor);
executor
.read(Some(self.filters()), self.routing_key.clone(), args, None)
.await
}
pub async fn execute<E: Executor>(&self, executor: &E) -> anyhow::Result<Vec<Event>> {
let filters = self.filters();
let mut cursor = if self.backward {
self.before.clone()
} else {
self.after.clone()
};
let mut remaining = self.limit.map(usize::from);
let mut pages: Vec<Vec<Event>> = Vec::new();
loop {
let page_size = match remaining {
Some(0) => break,
Some(r) => r.min(PAGE_SIZE),
None => PAGE_SIZE,
};
let result = executor
.read(
Some(filters.clone()),
self.routing_key.clone(),
self.args_for(page_size as u16, cursor.take()),
None,
)
.await?;
let (more, next) = if self.backward {
(
result.page_info.has_previous_page,
result.page_info.start_cursor.clone(),
)
} else {
(
result.page_info.has_next_page,
result.page_info.end_cursor.clone(),
)
};
let fetched = result.edges.len();
pages.push(result.edges.into_iter().map(|edge| edge.node).collect());
if let Some(r) = remaining.as_mut() {
*r = r.saturating_sub(fetched);
if *r == 0 {
break;
}
}
if fetched == 0 || !more {
break;
}
match next {
Some(c) => cursor = Some(c),
None => break,
}
}
if self.backward {
pages.reverse();
}
Ok(pages.into_iter().flatten().collect())
}
}
impl<A: AggregateEvents> ReadBuilder<A> {
pub async fn decode<E: Executor>(&self, executor: &E) -> anyhow::Result<Vec<A::Events>> {
let events = self.execute(executor).await?;
events
.iter()
.map(|event| A::Events::try_from(event).map_err(Into::into))
.collect()
}
}