use crate::reactor::reactors::Reactor;
use pravega_client_channel::{create_channel, ChannelSender};
use pravega_client_shared::*;
use tokio::sync::oneshot;
use crate::client_factory::ClientFactory;
use crate::error::*;
use crate::get_random_u128;
use crate::reactor::event::{Incoming, PendingEvent};
use tracing::info_span;
use tracing_futures::Instrument;
pub struct EventStreamWriter {
writer_id: WriterId,
sender: ChannelSender<Incoming>,
}
impl EventStreamWriter {
pub const MAX_EVENT_SIZE: usize = 8 * 1024 * 1024;
const CHANNEL_CAPACITY: usize = 16 * 1024 * 1024;
pub(crate) fn new(stream: ScopedStream, factory: ClientFactory) -> Self {
let (tx, rx) = create_channel(Self::CHANNEL_CAPACITY);
let writer_id = WriterId::from(get_random_u128());
let span = info_span!("Reactor", event_stream_writer = %writer_id);
factory
.get_runtime()
.spawn(Reactor::run(stream, tx.clone(), rx, factory.clone(), None).instrument(span));
EventStreamWriter {
writer_id,
sender: tx,
}
}
pub async fn write_event(&mut self, event: Vec<u8>) -> oneshot::Receiver<Result<(), SegmentWriterError>> {
let size = event.len();
let (tx, rx) = oneshot::channel();
if let Some(pending_event) = PendingEvent::with_header(None, event, None, tx) {
let append_event = Incoming::AppendEvent(pending_event);
self.writer_event_internal(append_event, size, rx).await
} else {
rx
}
}
pub async fn write_event_by_routing_key(
&mut self,
routing_key: String,
event: Vec<u8>,
) -> oneshot::Receiver<Result<(), SegmentWriterError>> {
let size = event.len();
let (tx, rx) = oneshot::channel();
if let Some(pending_event) = PendingEvent::with_header(Some(routing_key), event, None, tx) {
let append_event = Incoming::AppendEvent(pending_event);
self.writer_event_internal(append_event, size, rx).await
} else {
rx
}
}
async fn writer_event_internal(
&mut self,
append_event: Incoming,
size: usize,
rx: oneshot::Receiver<Result<(), SegmentWriterError>>,
) -> oneshot::Receiver<Result<(), SegmentWriterError>> {
if let Err(_e) = self.sender.send((append_event, size)).await {
let (tx_error, rx_error) = oneshot::channel();
tx_error
.send(Err(SegmentWriterError::SendToProcessor {}))
.expect("send error");
rx_error
} else {
rx
}
}
}
impl Drop for EventStreamWriter {
fn drop(&mut self) {
let _res = self.sender.send((Incoming::Close(), 0));
}
}
#[cfg(test)]
mod tests {
use tokio::runtime::Runtime;
use super::*;
use crate::reactor::event::PendingEvent;
#[test]
fn test_pending_event() {
let (tx, _rx) = oneshot::channel();
let data = vec![];
let routing_key = None;
let event = PendingEvent::without_header(routing_key, data, None, tx).expect("create pending event");
assert!(event.is_empty());
let (tx, rx) = oneshot::channel();
let data = vec![0; (PendingEvent::MAX_WRITE_SIZE + 1) as usize];
let routing_key = None;
let event = PendingEvent::without_header(routing_key, data, None, tx);
assert!(event.is_none());
let rt = Runtime::new().expect("get runtime");
let reply = rt.block_on(rx).expect("get reply");
assert!(reply.is_err());
}
}