use crate::client_factory::ClientFactory;
use crate::reader_group::reader_group_state::Offset;
use crate::segment_reader::ReaderError;
use crate::segment_slice::{SegmentDataBuffer, SegmentSlice, SliceMetadata};
use bytes::BufMut;
use im::HashMap as ImHashMap;
use pravega_client_shared::{Reader, ScopedSegment, Segment, SegmentWithRange};
use std::collections::{HashMap, HashSet};
use std::sync::Arc;
use std::time::{Duration, Instant};
use tokio::sync::mpsc::{Receiver, Sender};
use tokio::sync::oneshot;
use tokio::sync::{mpsc, Mutex};
use tokio::time::timeout;
use tracing::{debug, error, info, warn};
pub type ReaderErrorWithOffset = (ReaderError, i64);
pub type SegmentReadResult = Result<SegmentDataBuffer, ReaderErrorWithOffset>;
const REBALANCE_INTERVAL: Duration = Duration::from_secs(10);
cfg_if::cfg_if! {
if #[cfg(test)] {
use crate::reader_group::reader_group_state::MockReaderGroupState as ReaderGroupState;
} else {
use crate::reader_group::reader_group_state::ReaderGroupState;
}
}
#[derive(new)]
pub struct EventReader {
id: Reader,
factory: ClientFactory,
rx: Receiver<SegmentReadResult>,
tx: Sender<SegmentReadResult>,
meta: ReaderState,
rg_state: Arc<Mutex<ReaderGroupState>>,
}
pub struct ReaderState {
slices: HashMap<ScopedSegment, SliceMetadata>,
slices_dished_out: HashMap<ScopedSegment, i64>,
slice_release_receiver: HashMap<ScopedSegment, oneshot::Receiver<Option<SliceMetadata>>>,
slice_stop_reading: HashMap<ScopedSegment, oneshot::Sender<()>>,
last_segment_release: Instant,
last_segment_acquire: Instant,
}
impl ReaderState {
fn add_slice_release_receiver(
&mut self,
scoped_segment: ScopedSegment,
slice_return_rx: oneshot::Receiver<Option<SliceMetadata>>,
) {
self.slice_release_receiver
.insert(scoped_segment, slice_return_rx);
}
async fn wait_for_segment_slice_return(&mut self, segment: &ScopedSegment) -> Option<SliceMetadata> {
if let Some(receiver) = self.slice_release_receiver.remove(segment) {
match receiver.await {
Ok(returned_meta) => {
debug!("SegmentSlice returned {:?}", returned_meta);
returned_meta
}
Err(e) => {
error!("Error Segment slice was not returned {:?}", e);
panic!("A Segment slice was not returned to the Reader.");
}
}
} else {
warn!(
"Invalid segment {:?} provided for while waiting for segment slice return",
segment
);
None
}
}
fn close_all_slice_return_channel(&mut self) {
for (_, mut rx) in self.slice_release_receiver.drain() {
rx.close();
}
}
async fn remove_segment(&mut self, segment: ScopedSegment) -> Option<SliceMetadata> {
match self.slices.remove(&segment) {
Some(meta) => {
debug!(
"Segment slice {:?} has not been dished out for consumption",
&segment
);
Some(meta)
}
None => {
debug!(
"Segment slice for {:?} has already been dished out for consumption",
&segment
);
self.wait_for_segment_slice_return(&segment).await
}
}
}
fn add_slices(&mut self, meta: SliceMetadata) {
if self
.slices
.insert(ScopedSegment::from(meta.scoped_segment.as_str()), meta)
.is_some()
{
panic!("Pre-condition check failure. Segment slice already present");
}
}
fn add_stop_reading_tx(&mut self, segment: ScopedSegment, tx: oneshot::Sender<()>) {
if self.slice_stop_reading.insert(segment, tx).is_some() {
panic!("Pre-condition check failure. Sender used to stop fetching data is already present");
}
}
fn stop_reading(&mut self, segment: &ScopedSegment) {
if let Some(tx) = self.slice_stop_reading.remove(segment) {
if tx.send(()).is_err() {
debug!("Channel already closed, ignoring the error");
}
}
}
fn stop_reading_all(&mut self) {
for (_, tx) in self.slice_stop_reading.drain() {
if tx.send(()).is_err() {
debug!("Channel already closed, ignoring the error");
}
}
}
fn get_segment_id_with_data(&self) -> Option<ScopedSegment> {
self.slices
.iter()
.find_map(|(k, v)| if v.has_events() { Some(k.clone()) } else { None })
}
}
impl EventReader {
pub async fn init_reader(
id: String,
rg_state: Arc<Mutex<ReaderGroupState>>,
factory: ClientFactory,
) -> Self {
let reader = Reader::from(id);
let new_segments_to_acquire = rg_state
.lock()
.await
.compute_segments_to_acquire_or_release(&reader)
.await;
if new_segments_to_acquire > 0 {
for _ in 0..new_segments_to_acquire {
if let Some(seg) = rg_state
.lock()
.await
.assign_segment_to_reader(&reader)
.await
.expect("Error while waiting for segments to be assigned")
{
debug!("Acquiring segment {:?} for reader {:?}", seg, reader);
} else {
break;
}
}
}
let mut assigned_segments = rg_state
.lock()
.await
.get_segments_for_reader(&reader)
.await
.expect("Error while fetching currently assigned segments");
let mut slice_meta_map: HashMap<ScopedSegment, SliceMetadata> = HashMap::new();
slice_meta_map.extend(assigned_segments.drain().map(|(seg, offset)| {
(
seg.clone(),
SliceMetadata {
scoped_segment: seg.to_string(),
start_offset: offset.read,
read_offset: offset.read,
..Default::default()
},
)
}));
let (tx, rx) = mpsc::channel(1);
let mut stop_reading_map: HashMap<ScopedSegment, oneshot::Sender<()>> = HashMap::new();
slice_meta_map.iter().for_each(|(segment, meta)| {
let (tx_stop, rx_stop) = oneshot::channel();
stop_reading_map.insert(segment.clone(), tx_stop);
factory.get_runtime().enter();
tokio::spawn(SegmentSlice::get_segment_data(
segment.clone(),
meta.start_offset,
tx.clone(),
rx_stop,
factory.clone(),
));
});
EventReader::init_event_reader(
rg_state,
reader,
factory,
tx,
rx,
slice_meta_map,
stop_reading_map,
)
}
#[doc(hidden)]
pub fn init_event_reader(
rg_state: Arc<Mutex<ReaderGroupState>>,
id: Reader,
factory: ClientFactory,
tx: Sender<SegmentReadResult>,
rx: Receiver<SegmentReadResult>,
segment_slice_map: HashMap<ScopedSegment, SliceMetadata>,
slice_stop_reading: HashMap<ScopedSegment, oneshot::Sender<()>>,
) -> Self {
EventReader {
id,
factory,
rx,
tx,
meta: ReaderState {
slices: segment_slice_map,
slices_dished_out: Default::default(),
slice_release_receiver: HashMap::new(),
slice_stop_reading,
last_segment_release: Instant::now(),
last_segment_acquire: Instant::now(),
},
rg_state,
}
}
#[doc(hidden)]
pub fn set_last_acquire_release_time(&mut self, time: Instant) {
self.meta.last_segment_release = time;
self.meta.last_segment_acquire = time;
}
pub async fn release_segment(&mut self, mut slice: SegmentSlice) {
info!(
"releasing segment slice {} from reader {}",
slice.meta.scoped_segment, self.id
);
let scoped_segment = ScopedSegment::from(slice.meta.scoped_segment.clone().as_str());
self.meta.add_slices(slice.meta.clone());
self.meta.slices_dished_out.remove(&scoped_segment);
if self.meta.last_segment_release.elapsed() > REBALANCE_INTERVAL {
debug!("try to rebalance segments across readers");
let read_offset = slice.meta.read_offset;
self.release_segment_from_reader(slice, read_offset).await;
self.meta.last_segment_release = Instant::now();
} else {
if let Some(tx) = slice.slice_return_tx.take() {
if let Err(_e) = tx.send(Some(slice.meta.clone())) {
warn!(
"Failed to send segment slice release data for slice {:?}",
slice.meta
);
}
} else {
panic!("This is unexpected, No sender for SegmentSlice present.");
}
}
}
pub async fn release_segment_at(&mut self, slice: SegmentSlice, offset: i64) {
info!(
"releasing segment slice {} at offset {}",
slice.meta.scoped_segment, offset
);
assert!(
offset >= 0,
"the offset where the segment slice is released should be a positive number"
);
assert!(
slice.meta.start_offset <= offset,
"the offset where the segment slice is released should be greater than the start offset"
);
assert!(
slice.meta.end_offset >= offset,
"the offset where the segment slice is released should be less than the end offset"
);
let segment = ScopedSegment::from(slice.meta.scoped_segment.as_str());
if slice.meta.read_offset != offset {
self.meta.stop_reading(&segment);
let slice_meta = SliceMetadata {
start_offset: slice.meta.read_offset,
scoped_segment: slice.meta.scoped_segment.clone(),
last_event_offset: slice.meta.last_event_offset,
read_offset: offset,
end_offset: slice.meta.end_offset,
segment_data: SegmentDataBuffer::empty(),
partial_data_present: false,
};
let (tx_drop_fetch, rx_drop_fetch) = oneshot::channel();
self.factory.get_runtime().spawn(SegmentSlice::get_segment_data(
segment.clone(),
slice_meta.read_offset, self.tx.clone(),
rx_drop_fetch,
self.factory.clone(),
));
self.meta.add_stop_reading_tx(segment.clone(), tx_drop_fetch);
self.meta.add_slices(slice_meta);
self.meta.slices_dished_out.remove(&segment);
} else {
self.release_segment(slice).await;
}
}
pub async fn reader_offline(&mut self) {
info!("putting reader {} offline", self.id);
self.meta.stop_reading_all();
self.meta.close_all_slice_return_channel();
let mut offset_map: HashMap<ScopedSegment, Offset> = HashMap::new();
for (seg, off) in self.meta.slices_dished_out.drain() {
offset_map.insert(seg, Offset::new(off));
}
for (_, meta) in self.meta.slices.drain() {
offset_map.insert(
ScopedSegment::from(meta.scoped_segment.as_str()),
Offset::new(meta.read_offset),
);
}
self.rg_state
.lock()
.await
.remove_reader(&self.id, offset_map)
.await
.expect("Update ReaderGroup to ensure reader is offline");
}
async fn release_segment_from_reader(&mut self, mut slice: SegmentSlice, read_offset: i64) {
let new_segments_to_release = self
.rg_state
.lock()
.await
.compute_segments_to_acquire_or_release(&self.id)
.await;
let segment = ScopedSegment::from(slice.meta.scoped_segment.as_str());
if new_segments_to_release < 0 {
self.meta.stop_reading(&segment);
self.meta
.slices
.remove(&segment)
.expect("Segment missing in meta while releasing from reader");
if let Some(tx) = slice.slice_return_tx.take() {
if let Err(_e) = tx.send(None) {
warn!(
"Failed to send segment slice release data for slice {:?}",
slice.meta
);
}
} else {
panic!("This is unexpected, No sender for SegmentSlice present.");
}
self.rg_state
.lock()
.await
.release_segment(&self.id, &segment, &Offset::new(read_offset))
.await
.expect("Failed to release segment from RG state for reader");
}
}
pub async fn acquire_segment(&mut self) -> Option<SegmentSlice> {
info!("acquiring segment for reader {}", self.id);
if self.meta.last_segment_acquire.elapsed() > REBALANCE_INTERVAL {
info!("need to rebalance segments across readers");
if let Some(new_segments) = self.assign_segments_to_reader().await {
let current_segments = self
.rg_state
.lock()
.await
.get_segments_for_reader(&self.id)
.await
.expect("Read segments");
let new_segments: HashSet<(ScopedSegment, Offset)> = current_segments
.into_iter()
.filter(|(seg, _off)| new_segments.contains(seg))
.collect();
debug!("segments which can be read next are {:?}", new_segments);
self.initiate_segment_reads(new_segments);
self.meta.last_segment_acquire = Instant::now();
}
}
if let Some(segment_with_data) = self.meta.get_segment_id_with_data() {
info!("segment {} has data ready to read", segment_with_data);
let slice_meta = self.meta.slices.remove(&segment_with_data).unwrap();
let segment = ScopedSegment::from(slice_meta.scoped_segment.as_str());
let (slice_return_tx, slice_return_rx) = oneshot::channel();
self.meta.add_slice_release_receiver(segment, slice_return_rx);
info!(
"segment slice for {:?} is ready for consumption by reader {}",
slice_meta.scoped_segment, self.id,
);
self.meta
.slices_dished_out
.insert(segment_with_data, slice_meta.read_offset);
Some(SegmentSlice {
meta: slice_meta,
slice_return_tx: Some(slice_return_tx),
})
} else if let Ok(option) = timeout(Duration::from_millis(1000), self.rx.recv()).await {
if let Some(read_result) = option {
match read_result {
Ok(data) => {
let segment = ScopedSegment::from(data.segment.clone().as_str());
info!("new data fetched from server for segment {:?}", segment);
if let Some(mut slice_meta) = self.meta.remove_segment(segment.clone()).await {
if data.offset_in_segment
!= slice_meta.read_offset + slice_meta.segment_data.value.len() as i64
{
info!("Data from an invalid offset {:?} observed. Expected offset {:?}. Ignoring this data", data.offset_in_segment, slice_meta.read_offset);
None
} else {
EventReader::add_data_to_segment_slice(data, &mut slice_meta);
let (slice_return_tx, slice_return_rx) = oneshot::channel();
self.meta
.add_slice_release_receiver(segment.clone(), slice_return_rx);
self.meta
.slices_dished_out
.insert(segment.clone(), slice_meta.read_offset);
info!(
"segment slice for {:?} is ready for consumption by reader {}",
slice_meta.scoped_segment, self.id,
);
Some(SegmentSlice {
meta: slice_meta,
slice_return_tx: Some(slice_return_tx),
})
}
} else {
debug!("ignore the received data since None was returned");
return None;
}
}
Err((e, offset)) => {
let segment = ScopedSegment::from(e.get_segment().as_str());
debug!(
"Reader Error observed {:?} on segment {:?} at offset {:?} ",
e, segment, offset
);
if let Some(slice_meta) = self.meta.remove_segment(segment.clone()).await {
if slice_meta.read_offset != offset {
info!("Error at an invalid offset {:?} observed. Expected offset {:?}. Ignoring this data", offset, slice_meta.start_offset);
self.meta.add_slices(slice_meta);
self.meta.slices_dished_out.remove(&segment);
} else {
info!("Segment slice {:?} has received error {:?}", slice_meta, e);
self.fetch_successors(e).await;
}
}
debug!("segment Slice meta {:?}", self.meta.slices);
None
}
}
} else {
warn!("error getting updates from segment slice for reader {}", self.id);
None
}
} else {
info!(
"reader {} owns {} slices but none is ready to read",
self.id,
self.meta.slices.len()
);
None
}
}
async fn fetch_successors(&mut self, e: ReaderError) {
match e {
ReaderError::SegmentSealed {
segment,
can_retry: _,
operation: _,
error_msg: _,
}
| ReaderError::SegmentIsTruncated {
segment,
can_retry: _,
operation: _,
error_msg: _,
} => {
let completed_scoped_segment = ScopedSegment::from(segment.as_str());
self.meta.stop_reading(&completed_scoped_segment);
let successors = self
.factory
.get_controller_client()
.get_successors(&completed_scoped_segment)
.await
.expect("Failed to fetch successors of the segment")
.segment_with_predecessors;
info!("Segment Completed {:?}", segment);
self.rg_state
.lock()
.await
.segment_completed(&self.id, &completed_scoped_segment, &successors)
.await
.expect("Update segment completed");
if let Some(new_segments) = self.assign_segments_to_reader().await {
let current_segments = self
.rg_state
.lock()
.await
.get_segments_for_reader(&self.id)
.await
.expect("Read segments");
let new_segments: HashSet<(ScopedSegment, Offset)> = current_segments
.into_iter()
.filter(|(seg, _off)| new_segments.contains(seg))
.collect();
debug!("Segments which can be read next are {:?}", new_segments);
self.initiate_segment_reads(new_segments);
}
}
_ => error!("Error observed while reading from Pravega {:?}", e),
};
}
async fn assign_segments_to_reader(&self) -> Option<Vec<ScopedSegment>> {
let mut new_segments: Vec<ScopedSegment> = Vec::new();
let new_segments_to_acquire = self
.rg_state
.lock()
.await
.compute_segments_to_acquire_or_release(&self.id)
.await;
if new_segments_to_acquire <= 0 {
None
} else {
for _ in 0..new_segments_to_acquire {
if let Some(seg) = self
.rg_state
.lock()
.await
.assign_segment_to_reader(&self.id)
.await
.expect("Error while waiting for segments to be assigned")
{
debug!("Acquiring segment {:?} for reader {:?}", seg, self.id);
new_segments.push(seg);
} else {
break;
}
}
debug!("Segments acquired by reader {:?} is {:?}", self.id, new_segments);
Some(new_segments)
}
}
fn initiate_segment_reads(&mut self, new_segments: HashSet<(ScopedSegment, Offset)>) {
for (seg, offset) in new_segments {
let meta = SliceMetadata {
scoped_segment: seg.to_string(),
start_offset: offset.read,
read_offset: offset.read, ..Default::default()
};
let (tx_drop_fetch, rx_drop_fetch) = oneshot::channel();
tokio::spawn(SegmentSlice::get_segment_data(
seg.clone(),
meta.start_offset,
self.tx.clone(),
rx_drop_fetch,
self.factory.clone(),
));
self.meta.add_stop_reading_tx(seg, tx_drop_fetch);
self.meta.add_slices(meta);
}
}
fn add_data_to_segment_slice(data: SegmentDataBuffer, slice: &mut SliceMetadata) {
if slice.segment_data.value.is_empty() {
slice.segment_data = data;
} else {
slice.segment_data.value.put(data.value); slice.partial_data_present = false;
}
}
async fn get_successors(
&mut self,
completed_scoped_segment: &str,
) -> ImHashMap<SegmentWithRange, Vec<Segment>> {
let completed_scoped_segment = ScopedSegment::from(completed_scoped_segment);
self.factory
.get_controller_client()
.get_successors(&completed_scoped_segment)
.await
.expect("Failed to fetch successors of the segment")
.segment_with_predecessors
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::client_factory::ClientFactory;
use crate::error::SynchronizerError;
use crate::event_reader::{EventReader, SegmentReadResult};
use crate::reader_group::reader_group_state::ReaderGroupStateError;
use crate::segment_slice::{SegmentDataBuffer, SegmentSlice, SliceMetadata};
use bytes::{BufMut, BytesMut};
use mockall::predicate;
use mockall::predicate::*;
use pravega_client_config::{ClientConfigBuilder, MOCK_CONTROLLER_URI};
use pravega_client_shared::{Reader, Scope, ScopedSegment, ScopedStream, Stream};
use pravega_wire_protocol::commands::{Command, EventCommand};
use std::collections::HashMap;
use std::iter;
use std::sync::Arc;
use tokio::sync::mpsc::Sender;
use tokio::sync::oneshot;
use tokio::sync::oneshot::error::TryRecvError;
use tokio::sync::{mpsc, Mutex};
use tokio::time::{sleep, Duration};
use tracing::Level;
#[test]
fn test_read_events_single_segment() {
const NUM_EVENTS: usize = 100;
let (tx, rx) = mpsc::channel(1);
tracing_subscriber::fmt().with_max_level(Level::TRACE).finish();
let cf = ClientFactory::new(
ClientConfigBuilder::default()
.controller_uri(MOCK_CONTROLLER_URI)
.build()
.unwrap(),
);
let _guard = cf.get_runtime().enter();
tokio::spawn(generate_variable_size_events(
tx.clone(),
10,
NUM_EVENTS,
0,
false,
));
let init_segments = vec![create_segment_slice(0), create_segment_slice(1)];
let mut rg_mock: ReaderGroupState = ReaderGroupState::default();
rg_mock
.expect_compute_segments_to_acquire_or_release()
.return_const(0 as isize);
let mut reader = EventReader::init_event_reader(
Arc::new(Mutex::new(rg_mock)),
Reader::from("r1".to_string()),
cf.clone(),
tx.clone(),
rx,
create_slice_map(init_segments),
HashMap::new(),
);
let mut event_count = 0;
let mut event_size = 0;
while let Some(mut slice) = cf.get_runtime().block_on(reader.acquire_segment()) {
loop {
if let Some(event) = slice.next() {
println!("Read event {:?}", event);
assert_eq!(event.value.len(), event_size + 1, "Event has been missed");
assert!(is_all_same(event.value.as_slice()), "Event has been corrupted");
event_size += 1;
event_count += 1;
} else {
println!(
"Finished reading from segment {:?}, segment is auto released",
slice.meta.scoped_segment
);
break; }
}
if event_count == NUM_EVENTS {
break;
}
}
}
#[test]
fn test_acquire_segments() {
const NUM_EVENTS: usize = 10;
let (tx, rx) = mpsc::channel(1);
tracing_subscriber::fmt().with_max_level(Level::TRACE).finish();
let cf = ClientFactory::new(
ClientConfigBuilder::default()
.controller_uri(MOCK_CONTROLLER_URI)
.build()
.unwrap(),
);
let _guard = cf.get_runtime().enter();
tokio::spawn(generate_variable_size_events(
tx.clone(),
1024,
NUM_EVENTS,
0,
false,
));
let init_segments = vec![create_segment_slice(0)];
let mut rg_mock: ReaderGroupState = ReaderGroupState::default();
rg_mock
.expect_compute_segments_to_acquire_or_release()
.with(predicate::eq(Reader::from("r1".to_string())))
.return_const(1 as isize);
let res: Result<Option<ScopedSegment>, ReaderGroupStateError> =
Ok(Some(ScopedSegment::from("scope/test/1.#epoch.0")));
rg_mock
.expect_assign_segment_to_reader()
.with(predicate::eq(Reader::from("r1".to_string())))
.return_once(move |_| res);
let mut new_current_segments: HashSet<(ScopedSegment, Offset)> = HashSet::new();
new_current_segments.insert((ScopedSegment::from("scope/test/1.#epoch.0"), Offset::new(0)));
new_current_segments.insert((ScopedSegment::from("scope/test/0.#epoch.0"), Offset::new(0)));
let res: Result<HashSet<(ScopedSegment, Offset)>, SynchronizerError> = Ok(new_current_segments);
rg_mock
.expect_get_segments_for_reader()
.with(predicate::eq(Reader::from("r1".to_string())))
.return_once(move |_| res);
tokio::spawn(generate_variable_size_events(
tx.clone(),
1024,
NUM_EVENTS,
1,
false,
));
let before_time = Instant::now() - Duration::from_secs(15);
let mut reader = EventReader::init_event_reader(
Arc::new(Mutex::new(rg_mock)),
Reader::from("r1".to_string()),
cf.clone(),
tx.clone(),
rx,
create_slice_map(init_segments),
HashMap::new(),
);
reader.set_last_acquire_release_time(before_time);
let mut event_count = 0;
while let Some(mut slice) = cf.get_runtime().block_on(reader.acquire_segment()) {
loop {
if let Some(event) = slice.next() {
println!("Read event {:?}", event);
assert!(is_all_same(event.value.as_slice()), "Event has been corrupted");
event_count += 1;
} else {
println!(
"Finished reading from segment {:?}, segment is auto released",
slice.meta.scoped_segment
);
break; }
}
if event_count == NUM_EVENTS + NUM_EVENTS {
break;
}
}
assert_eq!(event_count, NUM_EVENTS + NUM_EVENTS);
}
#[test]
fn test_read_events_multiple_segments() {
const NUM_EVENTS: usize = 100;
let (tx, rx) = mpsc::channel(1);
tracing_subscriber::fmt().with_max_level(Level::TRACE).finish();
let cf = ClientFactory::new(
ClientConfigBuilder::default()
.controller_uri(MOCK_CONTROLLER_URI)
.build()
.unwrap(),
);
let _guard = cf.get_runtime().enter();
tokio::spawn(generate_variable_size_events(
tx.clone(),
100,
NUM_EVENTS,
0,
false,
));
tokio::spawn(generate_variable_size_events(
tx.clone(),
100,
NUM_EVENTS,
1,
true,
));
let init_segments = vec![create_segment_slice(0), create_segment_slice(1)];
let mut rg_mock: ReaderGroupState = ReaderGroupState::default();
rg_mock
.expect_compute_segments_to_acquire_or_release()
.return_const(0 as isize);
let mut reader = EventReader::init_event_reader(
Arc::new(Mutex::new(rg_mock)),
Reader::from("r1".to_string()),
cf.clone(),
tx.clone(),
rx,
create_slice_map(init_segments),
HashMap::new(),
);
let mut event_count_per_segment: HashMap<String, usize> = HashMap::new();
let mut total_events_read = 0;
while let Some(mut slice) = cf.get_runtime().block_on(reader.acquire_segment()) {
let segment = slice.meta.scoped_segment.clone();
println!("Received Segment Slice {:?}", segment);
let mut event_count = 0;
loop {
if let Some(event) = slice.next() {
println!("Read event {:?}", event);
assert!(is_all_same(event.value.as_slice()), "Event has been corrupted");
event_count += 1;
} else {
println!(
"Finished reading from segment {:?}, segment is auto released",
slice.meta.scoped_segment
);
break; }
}
total_events_read += event_count;
*event_count_per_segment
.entry(segment.clone())
.or_insert(event_count) += event_count;
if total_events_read == NUM_EVENTS * 2 {
break;
}
}
}
#[test]
fn test_return_slice() {
const NUM_EVENTS: usize = 2;
let (tx, rx) = mpsc::channel(1);
tracing_subscriber::fmt().with_max_level(Level::TRACE).finish();
let cf = ClientFactory::new(
ClientConfigBuilder::default()
.controller_uri(MOCK_CONTROLLER_URI)
.build()
.unwrap(),
);
let _guard = cf.get_runtime().enter();
tokio::spawn(generate_variable_size_events(
tx.clone(),
10,
NUM_EVENTS,
0,
false,
));
let init_segments = vec![create_segment_slice(0), create_segment_slice(1)];
let mut rg_mock: ReaderGroupState = ReaderGroupState::default();
rg_mock
.expect_compute_segments_to_acquire_or_release()
.return_const(0 as isize);
let mut reader = EventReader::init_event_reader(
Arc::new(Mutex::new(rg_mock)),
Reader::from("r1".to_string()),
cf.clone(),
tx.clone(),
rx,
create_slice_map(init_segments),
HashMap::new(),
);
let mut slice = cf.get_runtime().block_on(reader.acquire_segment()).unwrap();
let event = slice.next().unwrap();
assert_eq!(event.value.len(), 1);
assert!(is_all_same(event.value.as_slice()), "Event has been corrupted");
assert_eq!(event.offset_in_segment, 0);
cf.get_runtime().block_on(reader.release_segment(slice));
let slice = cf.get_runtime().block_on(reader.acquire_segment()).unwrap();
cf.get_runtime().block_on(reader.release_segment(slice));
let mut slice = cf.get_runtime().block_on(reader.acquire_segment()).unwrap();
let event = slice.next().unwrap();
assert_eq!(event.value.len(), 2);
assert!(is_all_same(event.value.as_slice()), "Event has been corrupted");
assert_eq!(event.offset_in_segment, 8 + 1); }
#[test]
fn test_return_slice_at_offset() {
const NUM_EVENTS: usize = 2;
let (tx, rx) = mpsc::channel(1);
let (stop_tx, stop_rx) = oneshot::channel();
tracing_subscriber::fmt().with_max_level(Level::TRACE).finish();
let cf = ClientFactory::new(
ClientConfigBuilder::default()
.controller_uri(MOCK_CONTROLLER_URI)
.build()
.unwrap(),
);
let _guard = cf.get_runtime().enter();
tokio::spawn(generate_constant_size_events(
tx.clone(),
20,
NUM_EVENTS,
0,
false,
stop_rx,
));
let mut stop_reading_map: HashMap<ScopedSegment, oneshot::Sender<()>> = HashMap::new();
stop_reading_map.insert(ScopedSegment::from("scope/test/0.#epoch.0"), stop_tx);
let init_segments = vec![create_segment_slice(0), create_segment_slice(1)];
let mut rg_mock: ReaderGroupState = ReaderGroupState::default();
rg_mock
.expect_compute_segments_to_acquire_or_release()
.return_const(0 as isize);
let mut reader = EventReader::init_event_reader(
Arc::new(Mutex::new(rg_mock)),
Reader::from("r1".to_string()),
cf.clone(),
tx.clone(),
rx,
create_slice_map(init_segments),
stop_reading_map,
);
let mut slice = cf.get_runtime().block_on(reader.acquire_segment()).unwrap();
let event = slice.next().unwrap();
assert_eq!(event.value.len(), 1);
assert!(is_all_same(event.value.as_slice()), "Event has been corrupted");
assert_eq!(event.offset_in_segment, 0);
let result = slice.next();
assert!(result.is_some());
let event = result.unwrap();
assert_eq!(event.value.len(), 1);
assert!(is_all_same(event.value.as_slice()), "Event has been corrupted");
assert_eq!(event.offset_in_segment, 9);
cf.get_runtime().block_on(reader.release_segment_at(slice, 0));
let (_stop_tx, stop_rx) = oneshot::channel();
tokio::spawn(generate_constant_size_events(
tx.clone(),
20,
NUM_EVENTS,
0,
false,
stop_rx,
));
let mut slice = cf.get_runtime().block_on(reader.acquire_segment()).unwrap();
let event = slice.next().unwrap();
assert_eq!(event.value.len(), 1);
assert!(is_all_same(event.value.as_slice()), "Event has been corrupted");
assert_eq!(event.offset_in_segment, 0); }
fn read_n_events(slice: &mut SegmentSlice, events_to_read: usize) {
let mut event_count = 0;
loop {
if event_count == events_to_read {
break;
}
if let Some(event) = slice.next() {
println!("Read event {:?}", event);
assert!(is_all_same(event.value.as_slice()), "Event has been corrupted");
event_count += 1;
} else {
println!(
"Finished reading from segment {:?}, segment is auto released",
slice.meta.scoped_segment
);
break;
}
}
}
fn create_slice_map(init_segments: Vec<SegmentSlice>) -> HashMap<ScopedSegment, SliceMetadata> {
let mut map = HashMap::with_capacity(init_segments.len());
for s in init_segments {
map.insert(
ScopedSegment::from(s.meta.scoped_segment.clone().as_str()),
s.meta.clone(),
);
}
map
}
fn get_scoped_stream(scope: &str, stream: &str) -> ScopedStream {
let stream: ScopedStream = ScopedStream {
scope: Scope {
name: scope.to_string(),
},
stream: Stream {
name: stream.to_string(),
},
};
stream
}
async fn generate_constant_size_events(
tx: Sender<SegmentReadResult>,
buf_size: usize,
num_events: usize,
segment_id: usize,
should_delay: bool,
mut stop_generation: oneshot::Receiver<()>,
) {
let mut segment_name = "scope/test/".to_owned();
segment_name.push_str(segment_id.to_string().as_ref());
let mut buf = BytesMut::with_capacity(buf_size);
let mut offset: i64 = 0;
for _i in 1..num_events + 1 {
if let Ok(_) | Err(TryRecvError::Closed) = stop_generation.try_recv() {
break;
}
let mut data = event_data(1); if data.len() < buf.capacity() - buf.len() {
buf.put(data);
} else {
while data.len() > 0 {
let free_space = buf.capacity() - buf.len();
if free_space == 0 {
if should_delay {
sleep(Duration::from_millis(100)).await;
}
tx.send(Ok(SegmentDataBuffer {
segment: ScopedSegment::from(segment_name.as_str()).to_string(),
offset_in_segment: offset,
value: buf,
}))
.await
.unwrap();
offset += buf_size as i64;
buf = BytesMut::with_capacity(buf_size);
} else if free_space >= data.len() {
buf.put(data.split());
} else {
buf.put(data.split_to(free_space));
}
}
}
}
tx.send(Ok(SegmentDataBuffer {
segment: ScopedSegment::from(segment_name.as_str()).to_string(),
offset_in_segment: offset,
value: buf,
}))
.await
.unwrap();
}
async fn generate_variable_size_events(
tx: Sender<SegmentReadResult>,
buf_size: usize,
num_events: usize,
segment_id: usize,
should_delay: bool,
) {
let mut segment_name = "scope/test/".to_owned();
segment_name.push_str(segment_id.to_string().as_ref());
segment_name.push_str(".#epoch.0");
let mut buf = BytesMut::with_capacity(buf_size);
let mut offset: i64 = 0;
for i in 1..num_events + 1 {
let mut data = event_data(i);
if data.len() < buf.capacity() - buf.len() {
buf.put(data);
} else {
while data.len() > 0 {
let free_space = buf.capacity() - buf.len();
if free_space == 0 {
if should_delay {
sleep(Duration::from_millis(100)).await;
}
tx.send(Ok(SegmentDataBuffer {
segment: ScopedSegment::from(segment_name.as_str()).to_string(),
offset_in_segment: offset,
value: buf,
}))
.await
.unwrap();
offset += buf_size as i64;
buf = BytesMut::with_capacity(buf_size);
} else if free_space >= data.len() {
buf.put(data.split());
} else {
buf.put(data.split_to(free_space));
}
}
}
}
tx.send(Ok(SegmentDataBuffer {
segment: ScopedSegment::from(segment_name.as_str()).to_string(),
offset_in_segment: offset,
value: buf,
}))
.await
.unwrap();
}
fn event_data(len: usize) -> BytesMut {
let mut buf = BytesMut::with_capacity(len + 8);
buf.put_i32(EventCommand::TYPE_CODE);
buf.put_i32(len as i32);
let mut data = Vec::new();
data.extend(iter::repeat(b'a').take(len));
buf.put(data.as_slice());
buf
}
fn create_segment_slice(segment_id: i64) -> SegmentSlice {
let mut segment_name = "scope/test/".to_owned();
segment_name.push_str(segment_id.to_string().as_ref());
let segment = ScopedSegment::from(segment_name.as_str());
let segment_slice = SegmentSlice {
meta: SliceMetadata {
start_offset: 0,
scoped_segment: segment.to_string(),
last_event_offset: 0,
read_offset: 0,
end_offset: i64::MAX,
segment_data: SegmentDataBuffer::empty(),
partial_data_present: false,
},
slice_return_tx: None,
};
segment_slice
}
fn is_all_same<T: Eq>(slice: &[T]) -> bool {
slice
.get(0)
.map(|first| slice.iter().all(|x| x == first))
.unwrap_or(true)
}
}