use std::str;
use rapidhash::RapidHashMap;
use crate::{
api::timeseries::{
TimeSeriesLabelFilter,
chunk::TSChunk,
conf::{
AggregationType, Aggregator, DuplicatePolicy, GroupReducerType, TSCreateOption,
TSDownStreamMeta, TSInfoResult, TSMGetOption, TSMGetResult, TSMRangeOption, TSMRangeResult,
TSRangeOption,
},
gorilla::TSSample,
group_samples_and_reduce, key,
meta::{TimeSeriesMeta, TimeSeriesMetaOptions},
traits::TimeSeries,
},
error::{Error, Result},
key::get_meta_checked,
meta::current_now_ms,
traits::DbLike,
};
fn trigger_downstream_upsert(db: &impl DbLike, src_key: &[u8], ts: u64, val: f64) -> Result<()> {
let kc = db.kc();
let meta_ks = db.meta();
let ds_prefix_k = key::downstream_prefix(&kc, src_key);
for g in meta_ks.prefix(&ds_prefix_k) {
let (k, v) = g.into_inner()?;
if let Some(dst_key) = k.strip_prefix(ds_prefix_k.as_slice())
&& 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 _ = db.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 = db.ts_range(src_key, bkt_left, end_bound)?;
if !bucket_samples.is_empty() {
let agg_val = agg.aggregate_samples(&bucket_samples);
let _ = db.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);
meta_ks.insert(&*k, ds_meta.encode())?;
}
}
Ok(())
}
fn cascade_downstream_del(
db: &impl DbLike,
src_key: &[u8],
from_ts: u64,
to_ts: u64,
) -> Result<()> {
let kc = db.kc();
let meta_ks = db.meta();
let ds_prefix_k = key::downstream_prefix(&kc, src_key);
for g in meta_ks.prefix(&ds_prefix_k) {
let (k, v) = g.into_inner()?;
if let Some(dst_key) = k.strip_prefix(ds_prefix_k.as_slice())
&& let Some(ds_meta) = TSDownStreamMeta::decode(&v)
{
let bkt_left = ds_meta.aggregator.calculate_aligned_bucket_left(from_ts);
let bkt_right = ds_meta.aggregator.calculate_aligned_bucket_right(to_ts);
let end_bound = if bkt_right > 0 { bkt_right - 1 } else { 0 };
let _ = db.ts_del(dst_key, bkt_left, end_bound);
}
}
Ok(())
}
impl<T: DbLike> TimeSeries for T {
fn ts_create_opt<K: AsRef<[u8]>>(&self, key: K, opt: &TSCreateOption) -> Result<()> {
let key_bytes = key.as_ref();
let meta_k = key::meta(&self.kc(), key_bytes);
let now_ms = current_now_ms();
if get_meta_checked::<TimeSeriesMeta>(self, key_bytes, &meta_k, now_ms)?.is_some() {
return Err(Error::invalid_data("ERR TSDB: key already exists"));
}
let meta_ks = self.meta();
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,
});
meta_ks.insert(&meta_k, meta.encode())?;
Ok(())
}
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 key_bytes = key.as_ref();
let meta_k = key::meta(&self.kc(), key_bytes);
let meta_ks = self.meta();
let now_ms = current_now_ms();
let mut meta = match get_meta_checked::<TimeSeriesMeta>(self, key_bytes, &meta_k, now_ms)? {
Some(b) => b,
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;
}
meta_ks.insert(&meta_k, meta.encode())?;
Ok(())
}
#[inline]
fn ts_add<K: AsRef<[u8]>>(&self, key: K, timestamp: u64, value: f64) -> Result<u64> {
self.ts_add_opt(key, timestamp, value, None, None)
}
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 key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = key::meta(&kc, key_bytes);
let meta_ks = self.meta();
let data_ks = self.data();
let now_ms = current_now_ms();
let mut meta = match get_meta_checked::<TimeSeriesMeta>(self, key_bytes, &meta_k, now_ms)? {
Some(m) => m,
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 = key::prefix(&kc, key_bytes);
let mut batch = self.batch();
let mut final_value = value;
if meta.total_samples == 0 {
let encoded_chunk = TSChunk::encode_with_type(&[sample], meta.chunk_type);
let item_k = key::item(&kc, key_bytes, timestamp);
batch.insert(data_ks, &item_k, encoded_chunk);
meta.total_samples = 1;
meta.first_time = timestamp;
meta.last_time = timestamp;
} else {
let target_item_k = key::item(&kc, key_bytes, timestamp);
let target_chunk = if let Some(g) = data_ks
.range(prefix.as_slice()..=target_item_k.as_slice())
.next_back()
{
let (k, v) = g.into_inner()?;
if k.starts_with(&prefix) {
Some((k, v))
} else {
None
}
} else {
None
};
let (old_k, old_data) = match target_chunk {
Some(pair) => pair,
None => {
let first_g = data_ks
.prefix(&prefix)
.next()
.ok_or_else(|| Error::invalid_data("ERR TSDB: corrupted timeseries data index"))?;
first_g.into_inner()?
}
};
let chunk_first_ts = if let Ok(b8) = old_k[prefix.len()..].try_into() {
u64::from_be_bytes(b8)
} else {
0
};
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(data_ks, &*old_k, new_chunk);
} else {
let is_latest_chunk = timestamp >= meta.last_time;
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 = key::item(&kc, key_bytes, timestamp);
batch.insert(data_ks, &new_item_k, new_chunk);
} else {
let new_chunks = TSChunk::upsert_and_split(
&old_data,
&[sample],
policy,
meta.chunk_size as usize,
meta.chunk_type,
)?;
batch.remove(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 = key::item(&kc, key_bytes, first_ts);
batch.insert(data_ks, &new_item_k, 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(meta_ks, &meta_k, meta.encode());
batch.commit()?;
trigger_downstream_upsert(self, key_bytes, timestamp, final_value)?;
Ok(timestamp)
}
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)
}
#[inline]
fn ts_incrby<K: AsRef<[u8]>>(&self, key: K, value: f64, timestamp: Option<u64>) -> Result<u64> {
self.ts_incrby_opt(key, value, timestamp, None)
}
fn ts_incrby_opt<K: AsRef<[u8]>>(
&self,
key: K,
value: f64,
timestamp: Option<u64>,
create_opt: Option<&TSCreateOption>,
) -> Result<u64> {
let ts = timestamp.unwrap_or_else(current_now_ms);
let latest = self.ts_get(key.as_ref())?;
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_opt(
key,
ts,
old_v + value,
Some(DuplicatePolicy::Last),
create_opt,
)
} else {
self.ts_add_opt(key, ts, value, Some(DuplicatePolicy::Last), create_opt)
}
}
#[inline]
fn ts_decrby<K: AsRef<[u8]>>(&self, key: K, value: f64, timestamp: Option<u64>) -> Result<u64> {
self.ts_decrby_opt(key, value, timestamp, None)
}
#[inline]
fn ts_decrby_opt<K: AsRef<[u8]>>(
&self,
key: K,
value: f64,
timestamp: Option<u64>,
create_opt: Option<&TSCreateOption>,
) -> Result<u64> {
self.ts_incrby_opt(key, -value, timestamp, create_opt)
}
fn ts_get<K: AsRef<[u8]>>(&self, key: K) -> Result<Option<(u64, f64)>> {
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = key::meta(&kc, key_bytes);
let data_ks = self.data();
let now_ms = current_now_ms();
let meta = get_meta_checked::<TimeSeriesMeta>(self, key_bytes, &meta_k, now_ms)?;
if let Some(m) = meta
&& m.total_samples > 0
{
let prefix = key::prefix(&kc, key_bytes);
if let Some(g) = data_ks.prefix(&prefix).next_back() {
let (k, v) = g.into_inner()?;
if k.starts_with(&prefix)
&& let Ok(Some(sample)) = TSChunk::get_latest_sample(&v)
{
return Ok(Some(sample));
}
}
}
Ok(None)
}
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 key_bytes = key.as_ref();
let kc = self.kc();
let prefix = key::prefix(&kc, key_bytes);
let meta_k = key::meta(&kc, key_bytes);
let data_ks = self.data();
let now_ms = current_now_ms();
let meta = get_meta_checked::<TimeSeriesMeta>(self, key_bytes, &meta_k, now_ms)?;
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 cap = meta
.as_ref()
.map(|m| (m.total_samples as usize).min(4096))
.unwrap_or(0);
let mut samples = Vec::with_capacity(cap);
let mut decoded_buf = Vec::with_capacity(1024);
let start_item_k = key::item(&kc, key_bytes, start_ts);
let start_range_k = if let Some(g) = data_ks
.range(prefix.as_slice()..=start_item_k.as_slice())
.next_back()
{
let (k, _) = g.into_inner()?;
if k.starts_with(&prefix) {
k.to_vec()
} else {
prefix.clone()
}
} else {
prefix.clone()
};
for g in data_ks.range(start_range_k..) {
let (k, v) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
let sub = &k[prefix.len()..];
if let Ok(b8) = sub.try_into() {
let chunk_first_ts = u64::from_be_bytes(b8);
if chunk_first_ts > end_ts {
break;
}
if let Some(chunk_last_ts) = TSChunk::get_last_timestamp(&v)
&& chunk_last_ts < start_ts
{
continue;
}
decoded_buf.clear();
if TSChunk::decode_samples_into(&v, &mut decoded_buf).is_ok() && !decoded_buf.is_empty() {
let first_sample_ts = decoded_buf[0].ts;
let last_sample_ts = decoded_buf[decoded_buf.len() - 1].ts;
if first_sample_ts >= start_ts && last_sample_ts <= end_ts {
samples.extend(decoded_buf.iter().map(|s| (s.ts, s.v)));
} else {
let start_idx = decoded_buf.partition_point(|s| s.ts < start_ts);
let end_idx = decoded_buf.partition_point(|s| s.ts <= end_ts);
if start_idx < end_idx {
samples.extend(decoded_buf[start_idx..end_idx].iter().map(|s| (s.ts, s.v)));
}
}
}
}
}
if !samples.windows(2).all(|w| w[0].0 <= w[1].0) {
samples.sort_by_key(|s| s.0);
}
Ok(samples)
}
#[inline]
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)
}
fn ts_range_opt<K: AsRef<[u8]>>(&self, key: K, opt: &TSRangeOption) -> Result<Vec<(u64, f64)>> {
let mut filtered = self.ts_range(key, opt.start_ts, opt.end_ts)?;
if !opt.filter_by_ts.is_empty() || opt.filter_by_value.is_some() {
filtered.retain(|&(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
});
}
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)
}
fn ts_del<K: AsRef<[u8]>>(&self, key: K, from_ts: u64, to_ts: u64) -> Result<usize> {
let key_bytes = key.as_ref();
let kc = self.kc();
let prefix = key::prefix(&kc, key_bytes);
let meta_k = key::meta(&kc, key_bytes);
let meta_ks = self.meta();
let data_ks = self.data();
let now_ms = current_now_ms();
let meta = match get_meta_checked::<TimeSeriesMeta>(self, key_bytes, &meta_k, now_ms)? {
Some(b) => b,
None => return Ok(0),
};
let ds_prefix_k = key::downstream_prefix(&kc, key_bytes);
let retention_bound = if meta.retention_time > 0 && meta.retention_time < meta.last_time {
meta.last_time - meta.retention_time
} else {
0
};
if retention_bound > 0 {
for g in meta_ks.prefix(&ds_prefix_k) {
let (k, v) = g.into_inner()?;
if k.starts_with(&ds_prefix_k)
&& let Some(ds_meta) = TSDownStreamMeta::decode(&v)
&& ds_meta.aggregator.calculate_aligned_bucket_left(from_ts) < retention_bound
{
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.batch();
let mut first_remaining_ts = None;
let mut last_remaining_ts = None;
for g in data_ks.prefix(&prefix) {
let (k, v) = g.into_inner()?;
if k.starts_with(&prefix) {
let sub = &k[prefix.len()..];
if let Ok(b8) = sub.try_into() {
let chunk_ts = u64::from_be_bytes(b8);
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(data_ks, &*k);
} else {
let new_first_ts = TSChunk::get_first_timestamp(&new_chunk_data).unwrap_or(chunk_ts);
let new_last_ts =
TSChunk::get_last_timestamp(&new_chunk_data).unwrap_or(new_first_ts);
let new_key = key::item(&kc, key_bytes, new_first_ts);
if new_key[..] != *k {
batch.remove(data_ks, &*k);
}
batch.insert(data_ks, &new_key, new_chunk_data);
if first_remaining_ts.is_none() {
first_remaining_ts = Some(new_first_ts);
}
last_remaining_ts = Some(new_last_ts);
}
} else {
let count = TSChunk::get_count(&v);
if count > 0 {
let chunk_first = TSChunk::get_first_timestamp(&v).unwrap_or(chunk_ts);
let chunk_last = TSChunk::get_last_timestamp(&v).unwrap_or(chunk_first);
if first_remaining_ts.is_none() {
first_remaining_ts = Some(chunk_first);
}
last_remaining_ts = Some(chunk_last);
}
}
}
}
}
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 Some(fts) = first_remaining_ts {
updated_meta.first_time = fts;
}
if let Some(lts) = last_remaining_ts {
updated_meta.last_time = lts;
}
}
batch.insert(meta_ks, &meta_k, updated_meta.encode());
batch.commit()?;
cascade_downstream_del(self, key_bytes, from_ts, to_ts)?;
Ok(total_deleted)
}
fn ts_mget(&self, opt: &TSMGetOption) -> Result<Vec<TSMGetResult>> {
let kc = self.kc();
let meta_ks = self.meta();
let filter = TimeSeriesLabelFilter::parse(&opt.filters);
let prefix = key::meta_prefix(&kc);
let mut results = Vec::new();
for g in 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(&k[prefix.len()..])?;
let labels = if opt.with_labels {
meta.labels
} else if !opt.selected_labels.is_empty() {
opt
.selected_labels
.iter()
.map(|sel_k| {
let val = meta
.labels
.iter()
.find(|(k, _)| k == sel_k)
.map(|(_, v)| v.as_str())
.unwrap_or("");
(sel_k.clone(), val.to_string())
})
.collect()
} else {
Vec::new()
};
results.push(TSMGetResult {
name,
labels,
sample: latest_sample,
});
}
}
Ok(results)
}
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 mut groups: RapidHashMap<String, GroupedSeriesEntry> = RapidHashMap::default();
let kc = self.kc();
let meta_ks = self.meta();
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 = key::meta(&kc, name.as_bytes());
meta_ks
.get(&meta_k)
.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)
}
#[inline]
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)
}
fn ts_info<K: AsRef<[u8]>>(&self, key: K) -> Result<TSInfoResult> {
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = key::meta(&kc, key_bytes);
let meta_ks = self.meta();
let now_ms = current_now_ms();
let meta = match get_meta_checked::<TimeSeriesMeta>(self, key_bytes, &meta_k, now_ms)? {
Some(b) => b,
None => return Err(Error::invalid_data("ERR TSDB: the key does not exist")),
};
let ds_prefix_k = key::downstream_prefix(&kc, key_bytes);
let mut downstream_rules = Vec::new();
for g in meta_ks.prefix(&ds_prefix_k) {
let (k, v) = g.into_inner()?;
if let Some(dst_key) = k.strip_prefix(ds_prefix_k.as_slice())
&& let Some(ds_meta) = TSDownStreamMeta::decode(&v)
{
downstream_rules.push((
String::from_utf8_lossy(dst_key).into_owned(),
ds_meta.aggregator,
));
}
}
let memory_usage = 128 + meta.total_samples * 16;
let (first_timestamp, last_timestamp) = if meta.total_samples > 0 {
(meta.first_time, 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,
})
}
fn ts_queryindex(&self, filters: &[String]) -> Result<Vec<String>> {
let filter = TimeSeriesLabelFilter::parse(filters);
let prefix = key::meta_prefix(&self.kc());
let meta_ks = self.meta();
let mut keys = Vec::new();
for g in 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)
}
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 src_bytes = src_key.as_ref();
let dst_bytes = dst_key.as_ref();
if src_bytes == dst_bytes {
return Err(Error::invalid_data(
"ERR TSDB: the source key and destination key should be different",
));
}
let kc = self.kc();
let meta_ks = self.meta();
let src_meta_k = key::meta(&kc, src_bytes);
let dst_meta_k = key::meta(&kc, dst_bytes);
let src_meta = match meta_ks.get(&src_meta_k)? {
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 meta_ks.get(&dst_meta_k)? {
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.as_bytes() != src_bytes {
return Err(Error::invalid_data(
"ERR TSDB: the destination key already has a src rule",
));
}
let dst_ds_prefix_k = key::downstream_prefix(&kc, dst_bytes);
if meta_ks.prefix(&dst_ds_prefix_k).next().is_some() {
return Err(Error::invalid_data(
"ERR TSDB: the destination key already has a dst rule",
));
}
let ds_meta_k = key::downstream_meta(&kc, src_bytes, dst_bytes);
let align = alignment.unwrap_or(0);
let ds = TSDownStreamMeta::new(Aggregator::new(aggregator, bucket_duration, align));
let mut batch = self.batch();
batch.insert(meta_ks, &ds_meta_k, ds.encode());
if dst_meta.source_key.as_bytes() != src_bytes {
dst_meta.source_key = String::from_utf8_lossy(src_bytes).into_owned();
batch.insert(meta_ks, &dst_meta_k, dst_meta.encode());
}
batch.commit()?;
Ok(())
}
fn ts_deleterule<SK: AsRef<[u8]>, DK: AsRef<[u8]>>(
&self,
src_key: SK,
dst_key: DK,
) -> Result<()> {
let src_bytes = src_key.as_ref();
let dst_bytes = dst_key.as_ref();
let kc = self.kc();
let meta_ks = self.meta();
let src_meta_k = key::meta(&kc, src_bytes);
let dst_meta_k = key::meta(&kc, dst_bytes);
if !meta_ks.contains_key(&src_meta_k)? || !meta_ks.contains_key(&dst_meta_k)? {
return Err(Error::invalid_data("ERR TSDB: the key is not a TSDB key"));
}
let ds_meta_k = key::downstream_meta(&kc, src_bytes, dst_bytes);
if !meta_ks.contains_key(&ds_meta_k)? {
return Err(Error::invalid_data(
"ERR TSDB: compaction rule does not exist",
));
}
let mut batch = self.batch();
batch.remove(meta_ks, &ds_meta_k);
if let Some(b) = meta_ks.get(&dst_meta_k)?
&& let Some(mut dst_meta) = TimeSeriesMeta::decode(&b)
&& dst_meta.source_key.as_bytes() == src_bytes
{
dst_meta.source_key.clear();
batch.insert(meta_ks, &dst_meta_k, dst_meta.encode());
}
batch.commit()?;
Ok(())
}
}