use std::any::Any;
use crate::{BusDeviceMetrics, ChunkId, InstanceType, SegmentId};
#[derive(Debug, Copy, Clone, PartialEq)]
pub struct CollectSkipper {
pub skip: u64,
pub skipped: u64,
pub skipping: bool,
}
impl CollectSkipper {
pub fn new(skip: u64) -> Self {
CollectSkipper { skip, skipped: 0, skipping: skip > 0 }
}
#[inline(always)]
pub fn should_skip(&mut self) -> bool {
if !self.skipping {
return false;
}
if self.skip == 0 || self.skipped >= self.skip {
self.skipping = false;
return false;
}
self.skipped += 1;
true
}
#[inline(always)]
pub fn rows_to_skip(&mut self, rows: u64) -> u64 {
if !self.skipping {
return 0;
}
if self.skip == 0 || self.skipped >= self.skip {
self.skipping = false;
return 0;
}
if (self.skipped + rows) >= self.skip {
let result = self.skip - self.skipped;
self.skipped = self.skip;
self.skipping = false;
return result;
}
self.skipped += rows;
rows
}
#[inline(always)]
pub fn should_skip_query(&mut self, apply: bool) -> bool {
if !self.skipping {
return false;
}
if self.skip == 0 || self.skipped >= self.skip {
self.skipping = false;
return false;
}
if apply {
self.skipped += 1;
}
true
}
}
#[derive(Debug, Copy, Clone, PartialEq)]
pub struct CollectCounter {
pub initial_skip: u32,
pub initial_skipped: u32,
pub collect_count: u32,
pub collected: u32,
pub initial_skipping: bool,
pub final_skip_phase: bool,
}
impl CollectCounter {
pub fn new(initial_skip: u32, collect_count: u32) -> Self {
CollectCounter {
initial_skip,
initial_skipped: 0,
collect_count,
collected: 0,
initial_skipping: collect_count > 0 && initial_skip > 0,
final_skip_phase: collect_count == 0,
}
}
#[inline(always)]
pub fn should_skip(&mut self) -> bool {
if self.initial_skipping {
if self.initial_skip == 0 || self.initial_skipped >= self.initial_skip {
self.initial_skipping = false;
} else {
self.initial_skipped += 1;
return true;
}
}
if self.collected < self.collect_count {
self.collected += 1;
return false;
}
self.final_skip_phase = true;
true
}
#[inline(always)]
pub fn should_process(&mut self, rows: u32) -> Option<(u32, u32)> {
let mut skip = 0;
let mut rows = rows;
if self.initial_skipping {
if self.initial_skip == 0 {
self.initial_skipping = false;
} else if (self.initial_skipped + rows) >= self.initial_skip {
skip = self.initial_skip - self.initial_skipped;
rows -= skip;
self.initial_skipped = self.initial_skip;
self.initial_skipping = false;
if rows == 0 {
return None;
}
} else {
self.initial_skipped += rows;
return None;
}
}
if self.final_skip_phase {
None
} else if (self.collected + rows) >= self.collect_count {
let rows_to_collect = self.collect_count - self.collected;
self.final_skip_phase = true;
self.collected = self.collect_count;
if rows_to_collect == 0 {
None
} else {
Some((skip, rows_to_collect))
}
} else {
self.collected += rows;
Some((skip, rows))
}
}
pub fn reset(&mut self, initial_skip: u32, collect_count: u32) {
self.initial_skip = initial_skip;
self.initial_skipped = 0;
self.collect_count = collect_count;
self.collected = 0;
self.initial_skipping = initial_skip > 0;
self.final_skip_phase = false;
}
pub fn get_phase(&self) -> &str {
if self.initial_skipping {
"initial_skip"
} else if self.collected < self.collect_count {
"collecting"
} else {
"final_skip"
}
}
pub fn is_collecting(&self) -> bool {
!self.initial_skipping && self.collected < self.collect_count
}
pub fn is_final_skip(&self) -> bool {
self.final_skip_phase
}
pub fn remaining_to_collect(&self) -> u32 {
self.collect_count.saturating_sub(self.collected)
}
pub fn count(&self) -> u32 {
self.collect_count
}
pub fn skip(&self) -> u32 {
self.initial_skip
}
}
#[derive(Debug, Clone, PartialEq)]
pub enum CheckPoint {
None,
Single(ChunkId),
Multiple(Vec<ChunkId>),
}
#[derive(Debug)]
pub struct Plan {
pub airgroup_id: usize,
pub air_id: usize,
pub segment_id: Option<SegmentId>,
pub instance_type: InstanceType,
pub check_point: CheckPoint,
pub meta: Option<Box<dyn Any + Send + Sync>>,
pub global_id: Option<usize>,
}
impl Plan {
pub fn new(
airgroup_id: usize,
air_id: usize,
segment_id: Option<SegmentId>,
instance_type: InstanceType,
check_point: CheckPoint,
meta: Option<Box<dyn Any + Send + Sync>>,
) -> Self {
Plan { airgroup_id, air_id, segment_id, instance_type, check_point, meta, global_id: None }
}
pub fn set_global_id(&mut self, global_id: usize) {
self.global_id = Some(global_id);
}
}
pub trait Planner {
fn plan(&self, counter: Vec<(ChunkId, Box<dyn BusDeviceMetrics>)>) -> Vec<Plan>;
}