use std::collections::VecDeque;
use std::sync::Arc;
use chrono::{DateTime, Utc};
use crate::common::progress_tracker::{ProgressTracker, ProgressView, new_progress_tracker};
use parking_lot::Mutex;
use schemars::JsonSchema;
use crate::segment::common::anonymize::Anonymize;
use serde::{Deserialize, Serialize};
use uuid::Uuid;
use crate::shard::operations::optimization::{Optimization, OptimizationSegmentInfo};
use crate::shard::segment_holder::SegmentId;
const KEEP_LAST_TRACKERS: usize = 16;
#[derive(Default, Clone, Debug)]
pub struct TrackerLog {
descriptions: VecDeque<Tracker>,
}
impl TrackerLog {
pub fn register(&mut self, description: Tracker) {
self.descriptions.push_back(description);
self.truncate();
}
fn truncate(&mut self) {
let truncate_range = self.descriptions.len().saturating_sub(KEEP_LAST_TRACKERS);
let truncate = self
.descriptions
.iter()
.enumerate()
.take(truncate_range)
.filter(|(_, tracker)| match tracker.state.lock().status {
TrackerStatus::Optimizing | TrackerStatus::Error(_) => false,
TrackerStatus::Done | TrackerStatus::Cancelled(_) => true,
})
.map(|(index, _)| index)
.collect::<Vec<_>>();
truncate.into_iter().rev().for_each(|index| {
self.descriptions.remove(index);
});
}
pub fn to_telemetry(&self) -> Vec<TrackerTelemetry> {
self.descriptions
.iter()
.rev()
.map(Tracker::to_telemetry)
.collect()
}
pub fn iter(&self) -> impl Iterator<Item = &Tracker> {
self.descriptions.iter()
}
}
#[derive(Clone, Debug)]
pub struct Tracker {
pub name: &'static str,
pub uuid: Uuid,
pub segments: Vec<TrackerSegmentInfo>,
pub state: Arc<Mutex<TrackerState>>,
pub progress_view: ProgressView,
}
#[derive(Copy, Clone, Debug)]
pub struct TrackerSegmentInfo {
pub id: SegmentId,
pub uuid: Uuid,
pub points_count: usize,
}
impl From<&TrackerSegmentInfo> for OptimizationSegmentInfo {
fn from(value: &TrackerSegmentInfo) -> Self {
let TrackerSegmentInfo {
id: _,
uuid,
points_count,
} = *value;
OptimizationSegmentInfo { uuid, points_count }
}
}
impl Tracker {
pub fn start(
name: &'static str,
uuid: Uuid,
segments: Vec<TrackerSegmentInfo>,
) -> (Tracker, ProgressTracker) {
let (progress_view, progress_tracker) = new_progress_tracker();
let tracker = Self {
name,
uuid,
segments,
state: Default::default(),
progress_view,
};
(tracker, progress_tracker)
}
pub fn handle(&self) -> TrackerHandle {
self.state.clone().into()
}
pub fn to_telemetry(&self) -> TrackerTelemetry {
let state = self.state.lock();
TrackerTelemetry {
name: self.name,
uuid: self.uuid,
segment_ids: self.segments.iter().map(|s| s.id).collect(),
segment_uuids: self.segments.iter().map(|s| s.uuid).collect(),
status: state.status.clone(),
start_at: self.progress_view.started_at(),
end_at: state.end_at,
}
}
pub fn to_optimization(&self) -> Optimization {
Optimization {
optimizer: self.name.to_string(),
uuid: self.uuid,
segments: self
.segments
.iter()
.map(OptimizationSegmentInfo::from)
.collect(),
status: self.state.lock().status.clone(),
progress: self.progress_view.snapshot("Segment Optimizing"),
}
}
}
#[derive(Serialize, Deserialize, Clone, Debug, JsonSchema, )]
pub struct TrackerTelemetry {
pub name: &'static str,
pub uuid: Uuid,
pub segment_ids: Vec<SegmentId>,
pub segment_uuids: Vec<Uuid>,
pub status: TrackerStatus,
pub start_at: DateTime<Utc>,
pub end_at: Option<DateTime<Utc>>,
}
#[derive(Clone)]
pub struct TrackerHandle {
handle: Arc<Mutex<TrackerState>>,
}
impl TrackerHandle {
pub fn update(&self, status: TrackerStatus) {
self.handle.lock().update(status);
}
}
impl From<Arc<Mutex<TrackerState>>> for TrackerHandle {
fn from(state: Arc<Mutex<TrackerState>>) -> Self {
Self { handle: state }
}
}
#[derive(Debug, Default, Clone, PartialEq, Eq)]
pub struct TrackerState {
pub status: TrackerStatus,
pub end_at: Option<DateTime<Utc>>,
}
impl TrackerState {
pub fn update(&mut self, status: TrackerStatus) {
self.end_at = if status.is_running() {
None
} else {
Some(Utc::now())
};
self.status = status;
}
}
#[derive(
Serialize, Deserialize, Clone, Debug, JsonSchema, Default, Eq, PartialEq, Hash,
)]
#[serde(rename_all = "lowercase")]
pub enum TrackerStatus {
#[default]
Optimizing,
Done,
Cancelled(String),
Error(String),
}
impl TrackerStatus {
pub fn is_running(&self) -> bool {
match self {
TrackerStatus::Optimizing => true,
TrackerStatus::Done | TrackerStatus::Cancelled(_) | TrackerStatus::Error(_) => false,
}
}
}