use std::sync::Arc;
use std::time::Duration;
use crate::Position;
use crate::event::Event;
use crate::query::Query;
use super::{ReadConfig, ReadCore, ReadError, Reads, WaitOutcome};
pub const DEFAULT_MAX_BATCH_EVENTS: usize = 1024;
pub struct Subscription {
core: Arc<ReadCore>,
config: ReadConfig,
query: Query,
cursor: Position,
max_batch_events: usize,
}
impl Subscription {
pub(super) fn new(
core: Arc<ReadCore>,
config: ReadConfig,
query: Query,
after: Position,
) -> Subscription {
core.register_subscriber();
Subscription {
core,
config,
query,
cursor: after,
max_batch_events: DEFAULT_MAX_BATCH_EVENTS,
}
}
pub fn with_max_batch_events(mut self, max_batch_events: usize) -> Subscription {
self.max_batch_events = max_batch_events.max(1);
self
}
pub fn position(&self) -> Position {
self.cursor
}
pub fn poll_batch(&mut self) -> Result<Vec<(Position, Event)>, ReadError> {
let (watermark, snapshot) = self.core.load();
if self.cursor >= watermark {
return Ok(Vec::new());
}
let mut reads = Reads::plan(
snapshot,
&self.query,
self.cursor,
watermark,
&self.config,
None,
);
let mut out = Vec::new();
while let Some(item) = reads.next() {
let seq = item?;
out.push((seq.position, seq.event.to_owned()));
if out.len() >= self.max_batch_events {
self.cursor = out.last().unwrap().0;
return Ok(out);
}
}
if reads.is_exhausted() {
self.cursor = watermark;
} else if let Some((position, _)) = out.last() {
self.cursor = *position;
}
Ok(out)
}
pub fn wait(&self) -> bool {
matches!(
self.core.wait_past(self.cursor, None),
WaitOutcome::Advanced
)
}
pub fn wait_timeout(&self, timeout: Duration) -> WaitOutcome {
self.core.wait_past(self.cursor, Some(timeout))
}
pub fn next_batch(&mut self) -> Option<Result<Vec<(Position, Event)>, ReadError>> {
loop {
match self.poll_batch() {
Ok(batch) if !batch.is_empty() => return Some(Ok(batch)),
Ok(_) => {
if !self.wait() {
return match self.poll_batch() {
Ok(batch) if !batch.is_empty() => Some(Ok(batch)),
Ok(_) => None,
Err(err) => Some(Err(err)),
};
}
}
Err(err) => return Some(Err(err)),
}
}
}
#[cfg(feature = "async")]
pub async fn wait_async(&self) -> bool {
matches!(
self.core.wait_past_async(self.cursor).await,
WaitOutcome::Advanced
)
}
#[cfg(feature = "async")]
pub async fn next_batch_async(&mut self) -> Option<Result<Vec<(Position, Event)>, ReadError>> {
loop {
match self.poll_batch() {
Ok(batch) if !batch.is_empty() => return Some(Ok(batch)),
Ok(_) => {
if !self.wait_async().await {
return match self.poll_batch() {
Ok(batch) if !batch.is_empty() => Some(Ok(batch)),
Ok(_) => None,
Err(err) => Some(Err(err)),
};
}
}
Err(err) => return Some(Err(err)),
}
}
}
}
impl Drop for Subscription {
fn drop(&mut self) {
self.core.deregister_subscriber();
}
}