use circular_buffer::{CircularBuffer, Iter};
use lib::{
search::{Order, OrderBy},
KafkaRecord,
};
use rayon::prelude::*;
use tokio::sync::{
mpsc::UnboundedSender,
watch::{self, Receiver, Sender},
};
use crate::action::Action;
const BUFFER_SIZE: usize = 500;
pub struct RecordsBuffer {
buffer: CircularBuffer<BUFFER_SIZE, KafkaRecord>,
tx_action: Option<UnboundedSender<Action>>,
read: usize,
pub channels: (Sender<BufferAction>, Receiver<BufferAction>),
last_time_sorted: usize,
matched: usize,
}
macro_rules! sort_records {
($array:ident, $field: ident, $reverse: expr) => {
$array.par_sort_by(|a, b| {
let mut ordering = a.$field.cmp(&b.$field);
if $reverse {
ordering = ordering.reverse();
}
ordering
})
};
}
impl Default for RecordsBuffer {
fn default() -> Self {
Self::new()
}
}
impl RecordsBuffer {
pub fn new() -> Self {
Self {
buffer: CircularBuffer::<BUFFER_SIZE, KafkaRecord>::new(),
read: 0,
channels: watch::channel(BufferAction::Count((0, 0, 0))),
matched: 0,
last_time_sorted: 0,
tx_action: None,
}
}
pub fn register_action_handler(&mut self, tx: UnboundedSender<Action>) {
self.tx_action = Some(tx);
}
pub fn is_empty(&self) -> bool {
self.buffer.is_empty()
}
pub fn reset(&mut self) {
self.buffer.clear();
self.read = 0;
self.matched = 0;
self.dispatch_metrics();
}
pub fn matched_and_read(&self) -> (usize, usize, usize) {
(self.matched, self.read, self.buffer.len())
}
pub fn new_record_read(&mut self) {
self.read += 1;
}
pub fn len(&self) -> usize {
self.buffer.len()
}
pub fn get(&self, index: usize) -> Option<&KafkaRecord> {
self.buffer.get(index)
}
pub fn iter(&self) -> Iter<KafkaRecord> {
self.buffer.iter()
}
pub fn push(&mut self, kafka_record: KafkaRecord) -> usize {
self.buffer.push_back(kafka_record);
self.matched += 1;
self.matched
}
pub fn dispatch_metrics(&mut self) {
self.channels
.0
.send(BufferAction::Count(self.matched_and_read()))
.unwrap();
}
pub fn sort(&mut self, order_by: &OrderBy) {
let mut unsorted = self.buffer.to_vec();
if self.read == self.last_time_sorted {
return;
}
let reverse = order_by.is_descending();
match order_by.order {
Order::Timestamp => {
sort_records!(unsorted, timestamp, reverse)
}
Order::Key => {
sort_records!(unsorted, key_as_string, reverse)
}
Order::Value => sort_records!(unsorted, value_as_string, reverse),
Order::Partition => {
sort_records!(unsorted, partition, reverse)
}
Order::Offset => {
sort_records!(unsorted, offset, reverse)
}
Order::Size => unsorted.sort_by(|a, b| {
let mut ordering = a.size.cmp(&b.size);
if order_by.is_descending() {
ordering = ordering.reverse();
}
ordering
}),
Order::Topic => {
sort_records!(unsorted, topic, reverse)
}
}
self.buffer.clear();
self.buffer.extend(unsorted)
}
}
#[allow(clippy::large_enum_variant)]
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum BufferAction {
Count((usize, usize, usize)),
}