use flume::{self as channel, Sender};
use crate::Position;
use crate::event::Event;
use crate::log::set::PositionRange;
use crate::query::{AppendCondition, Query};
use crate::read::{ReadHandle, Reads, Subscription};
use super::{AppendError, AppendReply, Message, Request};
#[derive(Clone)]
pub struct WriteHandle {
pub(super) tx: Sender<Message>,
pub(super) reader: ReadHandle,
}
impl WriteHandle {
pub fn append(
&self,
events: Vec<Event>,
condition: Option<AppendCondition>,
) -> Result<PositionRange, AppendError> {
if events.is_empty() {
return Err(AppendError::Empty);
}
let (reply, response) = channel::unbounded();
let request = Request {
events,
condition,
reply,
token: 0,
};
self.tx
.send(Message::Append(request))
.map_err(|_| AppendError::Shutdown)?;
let (_token, result) = response.recv().map_err(|_| AppendError::Shutdown)?;
result
}
pub fn append_submit(
&self,
events: Vec<Event>,
condition: Option<AppendCondition>,
token: u64,
reply: AppendReply,
) -> Result<(), AppendError> {
if events.is_empty() {
return Err(AppendError::Empty);
}
let request = Request {
events,
condition,
reply,
token,
};
self.tx
.send(Message::Append(request))
.map_err(|_| AppendError::Shutdown)?;
Ok(())
}
pub fn read(&self, query: &Query, after: Position, limit: Option<u64>) -> Reads {
self.reader.read(query, after, limit)
}
pub fn subscribe(&self, query: Query, after: Position) -> Subscription {
self.reader.subscribe(query, after)
}
pub fn reader(&self) -> ReadHandle {
self.reader.clone()
}
#[cfg(feature = "async")]
pub async fn append_async(
&self,
events: Vec<Event>,
condition: Option<AppendCondition>,
) -> Result<PositionRange, AppendError> {
if events.is_empty() {
return Err(AppendError::Empty);
}
let (reply, response) = channel::unbounded();
let request = Request {
events,
condition,
reply,
token: 0,
};
self.tx
.send_async(Message::Append(request))
.await
.map_err(|_| AppendError::Shutdown)?;
let (_token, result) = response
.recv_async()
.await
.map_err(|_| AppendError::Shutdown)?;
result
}
#[cfg(feature = "async")]
pub async fn append_submit_async(
&self,
events: Vec<Event>,
condition: Option<AppendCondition>,
token: u64,
reply: AppendReply,
) -> Result<(), AppendError> {
if events.is_empty() {
return Err(AppendError::Empty);
}
let request = Request {
events,
condition,
reply,
token,
};
self.tx
.send_async(Message::Append(request))
.await
.map_err(|_| AppendError::Shutdown)?;
Ok(())
}
}