use crate::try_lock;
use super::TimeSeriesRecord;
use alloc::{sync::Arc, vec::Vec};
use core::time::Duration;
#[cfg(feature = "std")]
use std::sync::Mutex;
#[cfg(not(feature = "std"))]
use crate::memory::allocator::Mutex;
#[derive(Debug, Clone, Default)]
pub struct PartitionStats {
pub record_count: usize,
pub uncompressed_size: usize,
pub compressed_size: usize,
pub last_access_time: u64,
}
#[derive(Debug)]
pub struct TimeSeriesPartition {
pub start_time: u64,
pub end_time: u64,
pub records: Vec<TimeSeriesRecord>,
pub compressed: bool,
pub stats: PartitionStats,
}
impl TimeSeriesPartition {
pub fn new(start_time: u64, end_time: u64) -> Self {
Self {
start_time,
end_time,
records: Vec::new(),
compressed: false,
stats: PartitionStats::default(),
}
}
pub fn calculate_size(&self) -> usize {
self.records.len() * core::mem::size_of::<TimeSeriesRecord>()
}
pub fn clear(&mut self) {
self.records.clear();
self.stats.record_count = 0;
self.stats.uncompressed_size = 0;
self.stats.compressed_size = 0;
}
}
pub struct PartitionManager {
partitions: Vec<Arc<Mutex<TimeSeriesPartition>>>,
partition_duration: u64,
max_partitions: usize,
}
impl PartitionManager {
pub fn new(partition_duration: Duration, max_partitions: usize) -> Self {
Self {
partitions: Vec::new(),
partition_duration: partition_duration.as_secs(),
max_partitions,
}
}
pub fn get_or_create_partition(&mut self, timestamp: u64) -> Arc<Mutex<TimeSeriesPartition>> {
let partition_key = timestamp / self.partition_duration;
let start_time = partition_key * self.partition_duration;
let end_time = start_time + self.partition_duration;
for partition in &self.partitions {
let p = try_lock!(partition);
if p.start_time == start_time {
return partition.clone();
}
}
let new_partition = Arc::new(Mutex::new(TimeSeriesPartition::new(start_time, end_time)));
self.partitions.push(new_partition.clone());
if self.partitions.len() > self.max_partitions {
self.partitions.remove(0);
}
new_partition
}
pub fn get_partitions_in_range(
&self,
start_time: u64,
end_time: u64,
) -> Vec<Arc<Mutex<TimeSeriesPartition>>> {
let mut result = Vec::new();
for partition in &self.partitions {
let p = try_lock!(partition);
if p.start_time <= end_time && p.end_time >= start_time {
result.push(partition.clone());
}
}
result
}
pub fn get_oldest_partition(&self) -> Option<Arc<Mutex<TimeSeriesPartition>>> {
self.partitions.first().cloned()
}
pub fn get_newest_partition(&self) -> Option<Arc<Mutex<TimeSeriesPartition>>> {
self.partitions.last().cloned()
}
pub fn cleanup_expired_partitions(&mut self, current_time: u64, retention_period: Duration) {
let expire_time = current_time - retention_period.as_secs();
self.partitions.retain(|partition| {
let p = try_lock!(partition);
p.end_time > expire_time
});
}
pub fn get_partition_count(&self) -> usize {
self.partitions.len()
}
pub fn get_partition(&self, timestamp: u64) -> Option<Arc<Mutex<TimeSeriesPartition>>> {
let partition_key = timestamp / self.partition_duration;
let start_time = partition_key * self.partition_duration;
for partition in &self.partitions {
let p = try_lock!(partition);
if p.start_time == start_time {
return Some(partition.clone());
}
}
None
}
pub fn compress_all_partitions(&self) {
for partition in &self.partitions {
let mut p = try_lock!(partition);
if !p.compressed {
p.compressed = true;
}
}
}
}