pub mod chunk;
pub mod conf;
pub mod gorilla;
pub mod meta;
pub use chunk::{ChunkHeader, MergeStats, TSChunk};
pub use conf::{
AggregationType, Aggregator, BucketTimestampType, GroupReducerType, TSCreateOption,
TSDownStreamMeta, TSInfoResult, TSMGetOption, TSMGetResult, TSMRangeOption, TSMRangeResult,
TSRangeOption,
};
pub use gorilla::TSSample;
pub use meta::{ChunkType, DuplicatePolicy, TimeSeriesMeta, TimeSeriesMetaOptions};
use std::cmp::Ordering;
use std::collections::BinaryHeap;
use std::str;
use rapidhash::{RapidHashMap, RapidHashSet};
use crate::db::WeDb;
use crate::error::{Error, Result};
use crate::key_composer::KeyComposer;
#[derive(Debug, Clone, Default)]
pub struct TimeSeriesLabelFilter {
pub equals: RapidHashMap<String, RapidHashSet<String>>,
pub not_equals: RapidHashMap<String, RapidHashSet<String>>,
pub has_matchers: bool,
}
impl TimeSeriesLabelFilter {
pub fn new() -> Self {
Self::default()
}
pub fn parse(filters: &[String]) -> Self {
let mut filter = Self::new();
for f in filters {
filter.add_filter(f);
}
filter
}
pub fn add_filter(&mut self, expr: &str) -> bool {
let trimmed = expr.trim();
if trimmed.is_empty() {
return false;
}
let (op_pos, is_not_equal) = Self::find_operator(trimmed);
if op_pos == usize::MAX {
return false;
}
let label = trimmed[..op_pos].trim().to_string();
let value_str = if is_not_equal {
trimmed[op_pos + 2..].trim()
} else {
trimmed[op_pos + 1..].trim()
};
if is_not_equal {
let mut vals = RapidHashSet::default();
if value_str.is_empty() {
vals.insert(String::new());
} else if value_str.starts_with('(') && value_str.ends_with(')') {
for item in Self::split_value_list(&value_str[1..value_str.len() - 1]) {
let unquoted = Self::unquote(&item);
if !unquoted.is_empty() {
vals.insert(unquoted.to_string());
}
}
} else {
vals.insert(Self::unquote(value_str).to_string());
}
self.not_equals.entry(label).or_default().extend(vals);
self.has_matchers = true;
true
} else {
let mut vals = RapidHashSet::default();
if value_str.is_empty() {
self.equals.entry(label).or_default();
} else if value_str.starts_with('(') && value_str.ends_with(')') {
for item in Self::split_value_list(&value_str[1..value_str.len() - 1]) {
let unquoted = Self::unquote(&item);
if !unquoted.is_empty() {
vals.insert(unquoted.to_string());
}
}
self.equals.entry(label).or_default().extend(vals);
} else {
vals.insert(Self::unquote(value_str).to_string());
self.equals.entry(label).or_default().extend(vals);
}
self.has_matchers = true;
true
}
}
fn find_operator(expr: &str) -> (usize, bool) {
let mut quote = None;
let bytes = expr.as_bytes();
let len = bytes.len();
let mut i = 0;
while i < len {
let b = bytes[i];
if b == b'\'' || b == b'"' {
if quote == Some(b) {
quote = None;
} else if quote.is_none() {
quote = Some(b);
}
} else if quote.is_none() {
if b == b'!' && i + 1 < len && bytes[i + 1] == b'=' {
return (i, true);
} else if b == b'=' {
return (i, false);
}
}
i += 1;
}
(usize::MAX, false)
}
fn split_value_list(list: &str) -> Vec<String> {
let mut values = Vec::new();
let mut quote = None;
let mut depth = 0;
let mut start = 0;
let bytes = list.as_bytes();
let len = bytes.len();
for i in 0..=len {
if i == len {
if start < i {
let val = list[start..i].trim();
if !val.is_empty() {
values.push(val.to_string());
}
}
break;
}
let b = bytes[i];
if b == b'\'' || b == b'"' {
if quote == Some(b) {
quote = None;
} else if quote.is_none() {
quote = Some(b);
}
} else if quote.is_none() {
if b == b'(' {
depth += 1;
} else if b == b')' && depth > 0 {
depth -= 1;
} else if b == b',' && depth == 0 {
let val = list[start..i].trim();
if !val.is_empty() {
values.push(val.to_string());
}
start = i + 1;
}
}
}
values
}
#[inline]
fn unquote(s: &str) -> &str {
let s = s.trim();
if s.len() >= 2
&& ((s.starts_with('"') && s.ends_with('"'))
|| (s.starts_with('\'') && s.ends_with('\'')))
{
&s[1..s.len() - 1]
} else {
s
}
}
pub fn matches(&self, meta_labels: &[(String, String)]) -> bool {
if !self.has_matchers {
return true;
}
for (k, allowed_vals) in &self.equals {
let actual = meta_labels
.iter()
.find(|(lk, _)| lk == k)
.map(|(_, lv)| lv.as_str());
if allowed_vals.is_empty() {
if actual.is_some() {
return false;
}
} else {
match actual {
Some(actual_v) => {
if !allowed_vals.contains(actual_v) {
return false;
}
}
None => return false,
}
}
}
for (k, forbidden_vals) in &self.not_equals {
let actual = meta_labels
.iter()
.find(|(lk, _)| lk == k)
.map(|(_, lv)| lv.as_str());
if forbidden_vals.contains("") && actual.is_none() {
return false;
}
if let Some(actual_v) = actual
&& forbidden_vals.contains(actual_v)
{
return false;
}
}
true
}
}
pub fn group_samples_and_reduce(
all_samples: &[Vec<(u64, f64)>],
reducer_type: GroupReducerType,
) -> Vec<(u64, f64)> {
if reducer_type == GroupReducerType::None || all_samples.is_empty() {
return Vec::new();
}
#[derive(Eq, PartialEq)]
struct HeapItem {
ts: u64,
vec_idx: usize,
sample_idx: usize,
}
impl Ord for HeapItem {
fn cmp(&self, other: &Self) -> Ordering {
other.ts.cmp(&self.ts) }
}
impl PartialOrd for HeapItem {
fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
Some(self.cmp(other))
}
}
let mut heap = BinaryHeap::new();
for (i, vec) in all_samples.iter().enumerate() {
if !vec.is_empty() {
heap.push(HeapItem {
ts: vec[0].0,
vec_idx: i,
sample_idx: 0,
});
}
}
let mut result = Vec::new();
let mut current_ts = None;
let mut current_values = Vec::new();
let reduce = |values: &[f64]| -> f64 {
if values.is_empty() {
return 0.0;
}
let count = values.len() as f64;
let sum: f64 = values.iter().sum();
match reducer_type {
GroupReducerType::None => 0.0,
GroupReducerType::Sum => sum,
GroupReducerType::Avg => sum / count,
GroupReducerType::Min => values.iter().copied().fold(f64::INFINITY, f64::min),
GroupReducerType::Max => values.iter().copied().fold(f64::NEG_INFINITY, f64::max),
GroupReducerType::Range => {
let min = values.iter().copied().fold(f64::INFINITY, f64::min);
let max = values.iter().copied().fold(f64::NEG_INFINITY, f64::max);
max - min
}
GroupReducerType::Count => count,
GroupReducerType::VarP | GroupReducerType::StdP => {
let mean = sum / count;
let var = values.iter().map(|&x| (x - mean).powi(2)).sum::<f64>() / count;
if reducer_type == GroupReducerType::StdP {
var.sqrt()
} else {
var
}
}
GroupReducerType::VarS | GroupReducerType::StdS => {
if count <= 1.0 {
0.0
} else {
let mean = sum / count;
let var =
values.iter().map(|&x| (x - mean).powi(2)).sum::<f64>() / (count - 1.0);
if reducer_type == GroupReducerType::StdS {
var.sqrt()
} else {
var
}
}
}
GroupReducerType::Twa => sum / count,
}
};
while let Some(top) = heap.pop() {
let val = all_samples[top.vec_idx][top.sample_idx].1;
match current_ts {
Some(ts) if ts == top.ts => {
current_values.push(val);
}
Some(ts) => {
result.push((ts, reduce(¤t_values)));
current_values.clear();
current_values.push(val);
current_ts = Some(top.ts);
}
None => {
current_ts = Some(top.ts);
current_values.push(val);
}
}
let next_idx = top.sample_idx + 1;
if next_idx < all_samples[top.vec_idx].len() {
heap.push(HeapItem {
ts: all_samples[top.vec_idx][next_idx].0,
vec_idx: top.vec_idx,
sample_idx: next_idx,
});
}
}
if let Some(ts) = current_ts {
result.push((ts, reduce(¤t_values)));
}
result
}
impl WeDb {
pub fn ts_create<K: AsRef<[u8]>, LK: AsRef<[u8]>, LV: AsRef<[u8]>>(
&self,
key: K,
retention_ms: u64,
duplicate_policy: DuplicatePolicy,
labels: &[(LK, LV)],
) -> Result<()> {
let label_vec = labels
.iter()
.map(|(lk, lv)| {
(
String::from_utf8_lossy(lk.as_ref()).to_string(),
String::from_utf8_lossy(lv.as_ref()).to_string(),
)
})
.collect();
self.ts_create_opt(
key,
&TSCreateOption {
retention_time: retention_ms,
chunk_size: TimeSeriesMeta::DEFAULT_CHUNK_SIZE,
chunk_type: ChunkType::Uncompressed,
duplicate_policy,
source_key: String::new(),
labels: label_vec,
},
)
}
pub fn ts_create_opt<K: AsRef<[u8]>>(&self, key: K, opt: &TSCreateOption) -> Result<()> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.ts_meta(k_str);
if self.meta_ks.contains_key(meta_k.as_bytes())? {
return Err(Error::invalid_data("ERR TSDB: key already exists"));
}
let meta = TimeSeriesMeta::with_options(TimeSeriesMetaOptions {
retention_time: opt.retention_time,
chunk_size: opt.chunk_size,
chunk_type: opt.chunk_type,
duplicate_policy: opt.duplicate_policy,
source_key: opt.source_key.clone(),
labels: opt.labels.clone(),
expire_at: 0,
version: 0,
});
self.meta_ks.insert(meta_k.as_bytes(), meta.encode())?;
Ok(())
}
pub fn ts_alter<K: AsRef<[u8]>>(
&self,
key: K,
retention_ms: Option<u64>,
chunk_size: Option<u64>,
duplicate_policy: Option<DuplicatePolicy>,
labels: Option<Vec<(String, String)>>,
) -> Result<()> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.ts_meta(k_str);
let mut meta = match self.meta_ks.get(meta_k.as_bytes())? {
Some(b) => TimeSeriesMeta::decode(&b)
.ok_or_else(|| Error::invalid_data("ERR TSDB: corrupted metadata"))?,
None => return Err(Error::invalid_data("ERR TSDB: the key does not exist")),
};
if let Some(r) = retention_ms {
meta.retention_time = r;
}
if let Some(c) = chunk_size {
meta.chunk_size = if c == 0 {
TimeSeriesMeta::DEFAULT_CHUNK_SIZE
} else {
c
};
}
if let Some(p) = duplicate_policy {
meta.duplicate_policy = p;
}
if let Some(l) = labels {
meta.labels = l;
}
self.meta_ks.insert(meta_k.as_bytes(), meta.encode())?;
Ok(())
}
pub fn ts_add<K: AsRef<[u8]>>(&self, key: K, timestamp: u64, value: f64) -> Result<u64> {
self.ts_add_opt(key, timestamp, value, None, None)
}
pub fn ts_add_opt<K: AsRef<[u8]>>(
&self,
key: K,
timestamp: u64,
value: f64,
on_duplicate: Option<DuplicatePolicy>,
create_opt: Option<&TSCreateOption>,
) -> Result<u64> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.ts_meta(k_str);
let mut meta = match self.meta_ks.get(meta_k.as_bytes())? {
Some(m_bytes) => TimeSeriesMeta::decode(&m_bytes).unwrap_or_else(|| {
TimeSeriesMeta::new(0, 4096, DuplicatePolicy::Block, Vec::new())
}),
None => {
if let Some(opt) = create_opt {
TimeSeriesMeta::with_options(TimeSeriesMetaOptions {
retention_time: opt.retention_time,
chunk_size: opt.chunk_size,
chunk_type: opt.chunk_type,
duplicate_policy: opt.duplicate_policy,
source_key: opt.source_key.clone(),
labels: opt.labels.clone(),
expire_at: 0,
version: 0,
})
} else {
TimeSeriesMeta::new(0, 4096, DuplicatePolicy::Block, Vec::new())
}
}
};
if meta.retention_time > 0 && timestamp < meta.last_time.saturating_sub(meta.retention_time)
{
return Err(Error::invalid_data(
"ERR TSDB: Timestamp is older than retention",
));
}
let policy = on_duplicate.unwrap_or(meta.duplicate_policy);
let sample = TSSample::new(timestamp, value);
let prefix = kc.ts_prefix(k_str);
let mut chunks: Vec<(u64, Vec<u8>, Vec<u8>)> = Vec::new(); for g in self.data_ks.prefix(&prefix) {
let (k, v) = g.into_inner()?;
if k.starts_with(&prefix) {
let sub = &k[prefix.len()..];
if let Ok(ts_str) = str::from_utf8(sub)
&& let Ok(chunk_first_ts) = u64::from_str_radix(ts_str, 16)
{
chunks.push((chunk_first_ts, k.to_vec(), v.to_vec()));
}
}
}
let mut batch = self.db.batch();
let mut final_value = value;
if chunks.is_empty() {
let encoded_chunk = TSChunk::encode_with_type(&[sample], meta.chunk_type);
let item_k = kc.ts_item(k_str, timestamp);
batch.insert(&self.data_ks, item_k.as_bytes(), encoded_chunk);
meta.total_samples = 1;
meta.first_time = timestamp;
meta.last_time = timestamp;
} else {
let target_idx = chunks
.iter()
.rposition(|(chunk_first_ts, _, _)| *chunk_first_ts <= timestamp)
.unwrap_or(0);
let (chunk_first_ts, old_k, old_data) = &chunks[target_idx];
let mut samples = TSChunk::decode_samples(old_data)?;
if let Ok(idx) = samples.binary_search_by_key(×tamp, |s| s.ts) {
match policy.merge_value(samples[idx].v, value) {
Some(merged) => {
final_value = merged;
samples[idx].v = merged;
}
None => {
return Err(Error::invalid_data(
"ERR TSDB: Error at upsert, update is not supported when DUPLICATE_POLICY is set to BLOCK mode",
));
}
}
let new_chunk = TSChunk::encode_with_type(&samples, meta.chunk_type);
batch.insert(&self.data_ks, old_k, new_chunk);
} else {
let is_latest_chunk = target_idx == chunks.len() - 1;
let last_sample_ts = samples.last().map(|s| s.ts).unwrap_or(*chunk_first_ts);
if is_latest_chunk
&& samples.len() >= meta.chunk_size as usize
&& timestamp > last_sample_ts
{
let new_chunk = TSChunk::encode_with_type(&[sample], meta.chunk_type);
let new_item_k = kc.ts_item(k_str, timestamp);
batch.insert(&self.data_ks, new_item_k.as_bytes(), new_chunk);
} else {
let new_chunks = TSChunk::upsert_and_split(
old_data,
&[sample],
policy,
meta.chunk_size as usize,
meta.chunk_type,
)?;
batch.remove(&self.data_ks, old_k);
for chunk_data in new_chunks {
if let Some(first_ts) = TSChunk::get_first_timestamp(&chunk_data) {
let new_item_k = kc.ts_item(k_str, first_ts);
batch.insert(&self.data_ks, new_item_k.as_bytes(), chunk_data);
}
}
}
meta.total_samples += 1;
}
if meta.first_time == 0 || timestamp < meta.first_time {
meta.first_time = timestamp;
}
if timestamp > meta.last_time {
meta.last_time = timestamp;
}
}
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
batch.commit()?;
self.trigger_downstream_upsert(k_str, timestamp, final_value)?;
Ok(timestamp)
}
fn trigger_downstream_upsert(&self, src_key: &str, ts: u64, val: f64) -> Result<()> {
let kc = KeyComposer::new("default");
let ds_prefix = format!("_ts_ds:{src_key}:");
let ds_prefix_k = kc.raw_key(&ds_prefix);
for g in self.meta_ks.prefix(ds_prefix_k.as_bytes()) {
let (k, v) = g.into_inner()?;
if let Ok(k_str) = str::from_utf8(&k)
&& k_str.starts_with(&ds_prefix)
&& let Some(dst_key) = k_str.strip_prefix(&ds_prefix)
&& let Some(mut ds_meta) = TSDownStreamMeta::decode(&v)
{
let bkt_left = ds_meta.aggregator.calculate_aligned_bucket_left(ts);
let agg = &ds_meta.aggregator;
if agg.agg_type.is_incremental() {
let (inc_val, inc_policy) = match agg.agg_type {
AggregationType::Sum => (val, DuplicatePolicy::Sum),
AggregationType::Count => (1.0, DuplicatePolicy::Sum),
AggregationType::Min => (val, DuplicatePolicy::Min),
AggregationType::Max => (val, DuplicatePolicy::Max),
_ => (val, DuplicatePolicy::Last),
};
let _ = self.ts_add_opt(dst_key, bkt_left, inc_val, Some(inc_policy), None);
} else {
let bkt_right = ds_meta.aggregator.calculate_aligned_bucket_right(bkt_left);
let end_bound = if bkt_right > 0 { bkt_right - 1 } else { 0 };
let bucket_samples = self.ts_range(src_key, bkt_left, end_bound)?;
if !bucket_samples.is_empty() {
let agg_val = ds_meta.aggregator.aggregate_samples(&bucket_samples);
let _ = self.ts_add_opt(
dst_key,
bkt_left,
agg_val,
Some(DuplicatePolicy::Last),
None,
);
}
}
ds_meta.latest_bucket_idx = ds_meta.latest_bucket_idx.max(bkt_left);
self.meta_ks.insert(&*k, ds_meta.encode())?;
}
}
Ok(())
}
pub fn ts_madd<K: AsRef<[u8]>>(&self, items: &[(K, u64, f64)]) -> Result<Vec<Result<u64>>> {
let mut results = Vec::with_capacity(items.len());
for (k, ts, v) in items {
results.push(self.ts_add(k, *ts, *v));
}
Ok(results)
}
pub fn ts_incrby<K: AsRef<[u8]>>(
&self,
key: K,
value: f64,
timestamp: Option<u64>,
) -> Result<u64> {
let ts = timestamp.unwrap_or_else(|| ts_::sec() * 1000);
let latest = self.ts_get(&key)?;
if let Some((latest_ts, old_v)) = latest {
if ts < latest_ts {
return Err(Error::invalid_data(
"ERR TSDB: timestamp must be equal to or higher than the maximum existing timestamp",
));
}
self.ts_add(key, ts, old_v + value)
} else {
self.ts_add(key, ts, value)
}
}
pub fn ts_decrby<K: AsRef<[u8]>>(
&self,
key: K,
value: f64,
timestamp: Option<u64>,
) -> Result<u64> {
self.ts_incrby(key, -value, timestamp)
}
pub fn ts_get<K: AsRef<[u8]>>(&self, key: K) -> Result<Option<(u64, f64)>> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.ts_meta(k_str);
let meta = match self.meta_ks.get(meta_k.as_bytes())? {
Some(m_bytes) => TimeSeriesMeta::decode(&m_bytes),
None => None,
};
if let Some(m) = meta
&& m.total_samples > 0
{
let prefix = kc.ts_prefix(k_str);
let mut last_sample = None;
for g in self.data_ks.prefix(&prefix) {
let (k, v) = g.into_inner()?;
if k.starts_with(&prefix)
&& let Ok(samples) = TSChunk::decode_samples(&v)
&& let Some(s) = samples.last()
{
last_sample = Some((s.ts, s.v));
}
}
if last_sample.is_some() {
return Ok(last_sample);
}
}
Ok(None)
}
pub fn ts_range<K: AsRef<[u8]>>(
&self,
key: K,
from_ts: u64,
to_ts: u64,
) -> Result<Vec<(u64, f64)>> {
if from_ts > to_ts {
return Ok(Vec::new());
}
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let prefix = kc.ts_prefix(k_str);
let meta_k = kc.ts_meta(k_str);
let meta = self
.meta_ks
.get(meta_k.as_bytes())?
.and_then(|b| TimeSeriesMeta::decode(&b));
let retention_bound = if let Some(ref m) = meta
&& m.retention_time > 0
{
m.last_time.saturating_sub(m.retention_time)
} else {
0
};
let start_ts = from_ts.max(retention_bound);
let end_ts = to_ts;
let mut samples = Vec::new();
for g in self.data_ks.prefix(&prefix) {
let (k, v) = g.into_inner()?;
if k.starts_with(&prefix) {
let sub = &k[prefix.len()..];
if let Ok(ts_str) = str::from_utf8(sub)
&& let Ok(chunk_first_ts) = u64::from_str_radix(ts_str, 16)
{
if chunk_first_ts > end_ts {
break;
}
if let Ok(decoded) = TSChunk::decode_samples(&v) {
for s in decoded {
if s.ts >= start_ts && s.ts <= end_ts {
samples.push((s.ts, s.v));
}
}
}
}
}
}
samples.sort_by_key(|s| s.0);
Ok(samples)
}
pub fn ts_revrange<K: AsRef<[u8]>>(
&self,
key: K,
from_ts: u64,
to_ts: u64,
) -> Result<Vec<(u64, f64)>> {
let mut range = self.ts_range(key, from_ts, to_ts)?;
range.reverse();
Ok(range)
}
pub fn ts_range_opt<K: AsRef<[u8]>>(
&self,
key: K,
opt: &TSRangeOption,
) -> Result<Vec<(u64, f64)>> {
let raw = self.ts_range(key, opt.start_ts, opt.end_ts)?;
let mut filtered: Vec<(u64, f64)> = raw
.into_iter()
.filter(|&(ts, v)| {
if !opt.filter_by_ts.is_empty() && !opt.filter_by_ts.contains(&ts) {
return false;
}
if let Some((min_v, max_v)) = opt.filter_by_value
&& (v < min_v || v > max_v)
{
return false;
}
true
})
.collect();
if let Some(ref agg) = opt.aggregator {
filtered = agg.split_and_aggregate_opt(
&filtered,
opt.count_limit,
opt.is_return_empty,
opt.bucket_timestamp_type,
);
} else if let Some(limit) = opt.count_limit {
filtered.truncate(limit);
}
Ok(filtered)
}
pub fn ts_del<K: AsRef<[u8]>>(&self, key: K, from_ts: u64, to_ts: u64) -> Result<usize> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let prefix = kc.ts_prefix(k_str);
let meta_k = kc.ts_meta(k_str);
let meta = match self.meta_ks.get(meta_k.as_bytes())? {
Some(b) => TimeSeriesMeta::decode(&b)
.ok_or_else(|| Error::invalid_data("ERR TSDB: corrupted metadata"))?,
None => return Ok(0),
};
let ds_prefix = format!("_ts_ds:{k_str}:");
let ds_prefix_k = kc.raw_key(&ds_prefix);
let has_downstream = self.meta_ks.prefix(ds_prefix_k.as_bytes()).next().is_some();
if has_downstream
&& meta.retention_time > 0
&& from_ts < meta.last_time.saturating_sub(meta.retention_time)
{
return Err(Error::invalid_data(
"ERR TSDB: When a series has compactions, deleting samples or compaction buckets beyond the series retention period is not possible",
));
}
let mut total_deleted = 0;
let mut batch = self.db.batch();
for g in self.data_ks.prefix(&prefix) {
let (k, v) = g.into_inner()?;
if k.starts_with(&prefix) {
let sub = &k[prefix.len()..];
if let Ok(ts_str) = str::from_utf8(sub)
&& let Ok(chunk_ts) = u64::from_str_radix(ts_str, 16)
{
if chunk_ts > to_ts {
break;
}
let (new_chunk_data, deleted) =
TSChunk::remove_samples_between(&v, from_ts, to_ts, meta.chunk_type)?;
if deleted > 0 {
total_deleted += deleted;
let count = TSChunk::get_count(&new_chunk_data);
if count == 0 {
batch.remove(&self.data_ks, &*k);
} else {
let new_first_ts =
TSChunk::get_first_timestamp(&new_chunk_data).unwrap_or(chunk_ts);
let new_key = kc.ts_item(k_str, new_first_ts);
if new_key.as_bytes() != &*k {
batch.remove(&self.data_ks, &*k);
}
batch.insert(&self.data_ks, new_key.as_bytes(), new_chunk_data);
}
}
}
}
}
if total_deleted == 0 {
return Ok(0);
}
let mut updated_meta = meta;
updated_meta.total_samples = updated_meta
.total_samples
.saturating_sub(total_deleted as u64);
if updated_meta.total_samples == 0 {
updated_meta.first_time = 0;
updated_meta.last_time = 0;
} else if let Ok(remaining) = self.ts_range(k_str, 0, u64::MAX) {
if let Some(first_s) = remaining.first() {
updated_meta.first_time = first_s.0;
}
if let Some(last_s) = remaining.last() {
updated_meta.last_time = last_s.0;
}
}
batch.insert(&self.meta_ks, meta_k.as_bytes(), updated_meta.encode());
batch.commit()?;
self.cascade_downstream_del(k_str, from_ts, to_ts)?;
Ok(total_deleted)
}
fn cascade_downstream_del(&self, src_key: &str, from_ts: u64, to_ts: u64) -> Result<()> {
let kc = KeyComposer::new("default");
let ds_prefix = format!("_ts_ds:{src_key}:");
let ds_prefix_k = kc.raw_key(&ds_prefix);
for g in self.meta_ks.prefix(ds_prefix_k.as_bytes()) {
let (k, v) = g.into_inner()?;
if let Ok(k_str) = str::from_utf8(&k)
&& k_str.starts_with(&ds_prefix)
&& let Some(dst_key) = k_str.strip_prefix(&ds_prefix)
&& let Some(mut ds_meta) = TSDownStreamMeta::decode(&v)
{
let start_bkt = ds_meta.aggregator.calculate_aligned_bucket_left(from_ts);
let end_bkt = ds_meta.aggregator.calculate_aligned_bucket_left(to_ts);
let bkt_right = ds_meta.aggregator.calculate_aligned_bucket_right(start_bkt);
let bkt_samples = self.ts_range(src_key, start_bkt, bkt_right.saturating_sub(1))?;
if bkt_samples.is_empty() {
let _ = self.ts_del(dst_key, start_bkt, start_bkt);
} else {
let val = ds_meta.aggregator.aggregate_samples(&bkt_samples);
let _ =
self.ts_add_opt(dst_key, start_bkt, val, Some(DuplicatePolicy::Last), None);
}
if end_bkt > start_bkt + ds_meta.aggregator.bucket_duration {
let _ = self.ts_del(
dst_key,
start_bkt + ds_meta.aggregator.bucket_duration,
end_bkt.saturating_sub(1),
);
}
if end_bkt > start_bkt {
let bkt_right_end = ds_meta.aggregator.calculate_aligned_bucket_right(end_bkt);
let bkt_samples_end =
self.ts_range(src_key, end_bkt, bkt_right_end.saturating_sub(1))?;
if bkt_samples_end.is_empty() {
let _ = self.ts_del(dst_key, end_bkt, end_bkt);
} else {
let val = ds_meta.aggregator.aggregate_samples(&bkt_samples_end);
let _ = self.ts_add_opt(
dst_key,
end_bkt,
val,
Some(DuplicatePolicy::Last),
None,
);
}
}
if let Ok(Some((last_ts, _))) = self.ts_get(src_key) {
ds_meta.latest_bucket_idx =
ds_meta.aggregator.calculate_aligned_bucket_left(last_ts);
} else {
ds_meta.latest_bucket_idx = 0;
}
self.meta_ks.insert(&*k, ds_meta.encode())?;
}
}
Ok(())
}
pub fn ts_mget(&self, opt: &TSMGetOption) -> Result<Vec<TSMGetResult>> {
let kc = KeyComposer::new("default");
let filter = TimeSeriesLabelFilter::parse(&opt.filters);
let prefix = kc.ts_meta_prefix();
let mut results = Vec::new();
for g in self.meta_ks.prefix(&prefix) {
let (k, v) = g.into_inner()?;
if k.starts_with(&prefix)
&& let Some(meta) = TimeSeriesMeta::decode(&v)
&& filter.matches(&meta.labels)
{
let name = str::from_utf8(&k[prefix.len()..]).unwrap_or("").to_string();
let latest_sample = self.ts_get(&name)?;
let labels = if opt.with_labels {
if opt.selected_labels.is_empty() {
meta.labels
} else {
meta.labels
.into_iter()
.filter(|(lk, _)| opt.selected_labels.contains(lk))
.collect()
}
} else {
Vec::new()
};
results.push(TSMGetResult {
name,
labels,
sample: latest_sample,
});
}
}
Ok(results)
}
pub fn ts_mrange(&self, opt: &TSMRangeOption) -> Result<Vec<TSMRangeResult>> {
let mget_res = self.ts_mget(&opt.mget)?;
let mut series_results = Vec::new();
for item in mget_res {
let samples = self.ts_range_opt(&item.name, &opt.range)?;
series_results.push((item.name, item.labels, samples));
}
type GroupedSeriesEntry = (Vec<(String, String)>, Vec<Vec<(u64, f64)>>, Vec<String>);
if let Some(ref group_label) = opt.group_by_label
&& opt.reducer != GroupReducerType::None
{
let kc = KeyComposer::new("default");
let mut groups: RapidHashMap<String, GroupedSeriesEntry> = RapidHashMap::default();
for (name, labels, samples) in series_results {
let group_val = labels
.iter()
.find(|(k, _)| k == group_label)
.map(|(_, v)| v.clone())
.or_else(|| {
let meta_k = kc.ts_meta(&name);
self.meta_ks
.get(meta_k.as_bytes())
.ok()
.flatten()
.and_then(|b| TimeSeriesMeta::decode(&b))
.and_then(|m| {
m.labels
.iter()
.find(|(k, _)| k == group_label)
.map(|(_, v)| v.clone())
})
})
.unwrap_or_default();
let entry = groups.entry(group_val.clone()).or_insert_with(|| {
(
vec![(group_label.clone(), group_val)],
Vec::new(),
Vec::new(),
)
});
entry.1.push(samples);
entry.2.push(name);
}
let mut final_res = Vec::new();
for (group_val, (mut labels, all_samples, source_keys)) in groups {
let reduced = group_samples_and_reduce(&all_samples, opt.reducer);
if opt.mget.with_labels {
labels.push(("__reducer__".to_string(), opt.reducer.as_str().to_string()));
labels.push(("__source__".to_string(), source_keys.join(",")));
}
final_res.push(TSMRangeResult {
name: format!("{group_label}={group_val}"),
labels,
samples: reduced,
source_keys,
});
}
return Ok(final_res);
}
let mut final_res = Vec::with_capacity(series_results.len());
for (name, labels, samples) in series_results {
final_res.push(TSMRangeResult {
source_keys: vec![name.clone()],
name,
labels,
samples,
});
}
Ok(final_res)
}
pub fn ts_mrevrange(&self, opt: &TSMRangeOption) -> Result<Vec<TSMRangeResult>> {
let mut res = self.ts_mrange(opt)?;
for r in &mut res {
r.samples.reverse();
}
Ok(res)
}
pub fn ts_info<K: AsRef<[u8]>>(&self, key: K) -> Result<TSInfoResult> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.ts_meta(k_str);
let meta = match self.meta_ks.get(meta_k.as_bytes())? {
Some(b) => TimeSeriesMeta::decode(&b)
.ok_or_else(|| Error::invalid_data("ERR TSDB: corrupted metadata"))?,
None => return Err(Error::invalid_data("ERR TSDB: the key does not exist")),
};
let ds_prefix = format!("_ts_ds:{k_str}:");
let ds_prefix_k = kc.raw_key(&ds_prefix);
let mut downstream_rules = Vec::new();
for g in self.meta_ks.prefix(ds_prefix_k.as_bytes()) {
let (k, v) = g.into_inner()?;
if let Ok(rk_str) = str::from_utf8(&k)
&& rk_str.starts_with(&ds_prefix)
&& let Some(dst_key) = rk_str.strip_prefix(&ds_prefix)
&& let Some(ds_meta) = TSDownStreamMeta::decode(&v)
{
downstream_rules.push((dst_key.to_string(), ds_meta.aggregator));
}
}
let memory_usage = 128 + meta.total_samples * 16;
let (first_timestamp, last_timestamp) = if meta.total_samples > 0 {
let samples = self.ts_range(k_str, 0, u64::MAX)?;
(
samples.first().map(|s| s.0).unwrap_or(meta.first_time),
samples.last().map(|s| s.0).unwrap_or(meta.last_time),
)
} else {
(0, 0)
};
Ok(TSInfoResult {
total_samples: meta.total_samples,
memory_usage,
first_timestamp,
last_timestamp,
retention_time: meta.retention_time,
chunk_size: meta.chunk_size,
chunk_type: meta.chunk_type,
duplicate_policy: meta.duplicate_policy,
labels: meta.labels,
source_key: meta.source_key,
downstream_rules,
})
}
pub fn ts_queryindex(&self, filters: &[String]) -> Result<Vec<String>> {
let kc = KeyComposer::new("default");
let filter = TimeSeriesLabelFilter::parse(filters);
let prefix = kc.ts_meta_prefix();
let mut keys = Vec::new();
for g in self.meta_ks.prefix(&prefix) {
let (k, v) = g.into_inner()?;
if k.starts_with(&prefix)
&& let Some(meta) = TimeSeriesMeta::decode(&v)
&& filter.matches(&meta.labels)
{
let name = str::from_utf8(&k[prefix.len()..]).unwrap_or("").to_string();
keys.push(name);
}
}
keys.sort();
Ok(keys)
}
pub fn ts_createrule<SK: AsRef<[u8]>, DK: AsRef<[u8]>>(
&self,
src_key: SK,
dst_key: DK,
aggregator: AggregationType,
bucket_duration: u64,
alignment: Option<u64>,
) -> Result<()> {
let kc = KeyComposer::new("default");
let src_str = str::from_utf8(src_key.as_ref()).unwrap_or("");
let dst_str = str::from_utf8(dst_key.as_ref()).unwrap_or("");
if src_str == dst_str {
return Err(Error::invalid_data(
"ERR TSDB: the source key and destination key should be different",
));
}
let src_meta_k = kc.ts_meta(src_str);
let dst_meta_k = kc.ts_meta(dst_str);
let src_meta = match self.meta_ks.get(src_meta_k.as_bytes())? {
Some(b) => TimeSeriesMeta::decode(&b)
.ok_or_else(|| Error::invalid_data("ERR TSDB: corrupted source metadata"))?,
None => return Err(Error::invalid_data("ERR TSDB: the key is not a TSDB key")),
};
let mut dst_meta = match self.meta_ks.get(dst_meta_k.as_bytes())? {
Some(b) => TimeSeriesMeta::decode(&b)
.ok_or_else(|| Error::invalid_data("ERR TSDB: corrupted dest metadata"))?,
None => return Err(Error::invalid_data("ERR TSDB: the key is not a TSDB key")),
};
if !src_meta.source_key.is_empty() {
return Err(Error::invalid_data(
"ERR TSDB: the source key already has a source rule",
));
}
if !dst_meta.source_key.is_empty() && dst_meta.source_key != src_str {
return Err(Error::invalid_data(
"ERR TSDB: the destination key already has a src rule",
));
}
let dst_ds_prefix = format!("_ts_ds:{dst_str}:");
let dst_ds_prefix_k = kc.raw_key(&dst_ds_prefix);
if self
.meta_ks
.prefix(dst_ds_prefix_k.as_bytes())
.next()
.is_some()
{
return Err(Error::invalid_data(
"ERR TSDB: the destination key already has a dst rule",
));
}
let ds_meta_k = kc.raw_key(&format!("_ts_ds:{src_str}:{dst_str}"));
let align = alignment.unwrap_or(0);
let ds = TSDownStreamMeta::new(Aggregator::new(aggregator, bucket_duration, align));
let mut batch = self.db.batch();
batch.insert(&self.meta_ks, ds_meta_k.as_bytes(), ds.encode());
if dst_meta.source_key != src_str {
dst_meta.source_key = src_str.to_string();
batch.insert(&self.meta_ks, dst_meta_k.as_bytes(), dst_meta.encode());
}
batch.commit()?;
Ok(())
}
pub fn ts_deleterule<SK: AsRef<[u8]>, DK: AsRef<[u8]>>(
&self,
src_key: SK,
dst_key: DK,
) -> Result<()> {
let kc = KeyComposer::new("default");
let src_str = str::from_utf8(src_key.as_ref()).unwrap_or("");
let dst_str = str::from_utf8(dst_key.as_ref()).unwrap_or("");
let src_meta_k = kc.ts_meta(src_str);
let dst_meta_k = kc.ts_meta(dst_str);
if !self.meta_ks.contains_key(src_meta_k.as_bytes())?
|| !self.meta_ks.contains_key(dst_meta_k.as_bytes())?
{
return Err(Error::invalid_data("ERR TSDB: the key is not a TSDB key"));
}
let ds_meta_k = kc.raw_key(&format!("_ts_ds:{src_str}:{dst_str}"));
if !self.meta_ks.contains_key(ds_meta_k.as_bytes())? {
return Err(Error::invalid_data(
"ERR TSDB: compaction rule does not exist",
));
}
let mut batch = self.db.batch();
batch.remove(&self.meta_ks, ds_meta_k.as_bytes());
if let Some(b) = self.meta_ks.get(dst_meta_k.as_bytes())?
&& let Some(mut dst_meta) = TimeSeriesMeta::decode(&b)
&& dst_meta.source_key == src_str
{
dst_meta.source_key.clear();
batch.insert(&self.meta_ks, dst_meta_k.as_bytes(), dst_meta.encode());
}
batch.commit()?;
Ok(())
}
}