use std::sync::Arc;
use std::time::Duration;
use crate::client::GnitzClient;
use crate::connection::{DeltaCursor, ScanReply, Sent, Target};
use crate::error::ClientError;
use crate::mirror::PollOutcome;
use crate::protocol::wal_block::decode_wal_block_into;
use crate::{Schema, ZSetBatch};
pub struct Pushed {
pub sub: u64,
pub result: Result<(ScanReply, DeltaCursor), ClientError>,
}
#[must_use]
pub struct Synced {
pub pushed: Vec<Pushed>,
pub mirrored: Vec<PollOutcome>,
}
impl GnitzClient {
pub async fn subscribe(
&mut self,
view: impl Into<Target>,
from: Option<DeltaCursor>,
reply_schema: &Arc<Schema>,
spec: &[u8],
) -> Result<(u64, ScanReply, DeltaCursor), ClientError> {
let id = self.next_sub;
self.next_sub += 1;
let sent = self
.session
.submit_delta_read(view.into(), from, reply_schema, spec, Some(id));
let (rows, cursor) = self.wait(sent).await?;
self.readers.insert(id, Arc::clone(reply_schema));
Ok((id, rows, cursor))
}
pub fn unsubscribe(&mut self, id: u64) {
self.readers.remove(&id);
}
pub async fn sync(&mut self, wait: Duration) -> Result<Synced, ClientError> {
let synced = self.begin_sync(wait)?;
let synced = self.wait(synced).await;
if let Err(e @ ClientError::Interrupted(_)) = synced {
return Err(e);
}
self.finish_sync(synced).await
}
pub fn begin_sync(&mut self, wait: Duration) -> Result<Sent<()>, ClientError> {
let wait = self.mirror_hold(wait)?;
let mirrored = self.mirror.iter().flat_map(|m| m.views.values().filter_map(|v| v.sub));
let held: Vec<u64> = self.readers.keys().copied().chain(mirrored).collect();
Ok(self.session.submit_sync(&held, wait))
}
pub async fn finish_sync(&mut self, synced: Result<(), ClientError>) -> Result<Synced, ClientError> {
let mirrored = self.advance_mirror(&synced).await?;
let pushed = self.drain_readers(&synced);
Ok(Synced { pushed, mirrored })
}
fn drain_readers(&mut self, synced: &Result<(), ClientError>) -> Vec<Pushed> {
let GnitzClient { session, readers, .. } = self;
let mut pushed = Vec::with_capacity(readers.len());
readers.retain(|&sub, schema| {
let taken = synced.clone().and_then(|()| session.take_pushed(sub));
let result = taken.and_then(|(blocks, cursor)| {
let mut batch = ZSetBatch::new(schema);
for block in blocks {
decode_wal_block_into(&mut batch, block.block(), schema)?;
}
let schema = Arc::clone(schema);
Ok((ScanReply { schema, batch, lsn: None }, cursor))
});
let held = result.is_ok();
pushed.push(Pushed { sub, result });
held
});
pushed
}
}