use yo_common::num::{DOUBLE_MAX, parse_f64, parse_i64, write_dragonbox};
use yo_common::{Code, Error, Result};
use yo_kv::{Db, Foreign, KeyCursor, Keyspace, Kind};
use yo_series::{
Agg, Buckets, Encoding, Policy, Query, Refused, Rows, Sample, Series, Stamp, Unread,
bucket_start,
};
use yo_series::Rule as Compaction;
use super::args::{self, Args};
use super::table::Spec;
use crate::reply::Out;
const EXISTS: &[u8] = b"TSDB: key already exists";
const MISSING: &[u8] = b"TSDB: the key does not exist";
const NOT_A_SERIES: &[u8] = b"TSDB: the key is not a TSDB key";
const BAD_LABELS: &[u8] = b"TSDB: Couldn't parse LABELS";
const BAD_RETENTION: &[u8] = b"TSDB: Couldn't parse RETENTION";
const BAD_CHUNK: &[u8] = b"TSDB: Couldn't parse CHUNK_SIZE";
const CHUNK_RANGE: &[u8] =
b"TSDB: CHUNK_SIZE value must be a multiple of 8 in the range [48 .. 1048576]";
const BAD_ENCODING: &[u8] = b"TSDB: unknown ENCODING parameter";
const BAD_POLICY: &[u8] = b"TSDB: Couldn't parse DUPLICATE_POLICY";
const UNKNOWN_POLICY: &[u8] = b"TSDB: Unknown DUPLICATE_POLICY";
const BAD_IGNORE: &[u8] = b"TSDB: Couldn't parse IGNORE";
const NEGATIVE_IGNORE: &[u8] = b"TSDB: IGNORE arguments cannot be negative";
const BAD_TIMESTAMP: &[u8] = b"TSDB: invalid timestamp";
const NEGATIVE_TIMESTAMP: &[u8] = b"TSDB: invalid timestamp, must be a nonnegative integer";
const BAD_VALUE: &[u8] = b"TSDB: invalid value";
const BAD_INCREMENT: &[u8] = b"TSDB: invalid increase/decrease value";
const NAN_INCREMENT: &[u8] = b"TSDB: cannot increment/decrement NaN value";
const TOO_OLD: &[u8] = b"TSDB: Timestamp is older than retention";
const UPSERT: &[u8] = b"TSDB: Error at upsert, update is not supported when DUPLICATE_POLICY is set to BLOCK mode, or either current or new value is NaN and DUPLICATE_POLICY is MAX/MIN/SUM";
const BAD_FROM: &[u8] = b"TSDB: wrong fromTimestamp";
const BAD_TO: &[u8] = b"TSDB: wrong toTimestamp";
const COUNT_MISSING: &[u8] = b"TSDB: COUNT argument is missing";
const BAD_COUNT: &[u8] = b"TSDB: Couldn't parse COUNT";
const COUNT_RANGE: &[u8] = b"TSDB: Invalid COUNT value";
const BAD_AGGREGATION: &[u8] = b"TSDB: Couldn't parse AGGREGATION";
const EMPTY_AGG: &[u8] = b"TSDB: Empty aggregation type in list";
const TOO_MANY_AGGS: &[u8] = b"TSDB: Too many aggregation types";
const UNKNOWN_AGG: &[u8] = b"TSDB: Unknown aggregation type";
const BAD_BUCKET: &[u8] = b"TSDB: bucketDuration must be greater than zero";
const EMPTY_PLACE: &[u8] = b"TSDB: EMPTY flag should be the 3rd or 5th flag after AGGREGATION flag";
const BUCKET_TS_PLACE: &[u8] =
b"TSDB: BUCKETTIMESTAMP flag should be the 3rd or 4th flag after AGGREGATION flag";
const BAD_BUCKET_TS: &[u8] = b"TSDB: unknown BUCKETTIMESTAMP parameter";
const BAD_ALIGN: &[u8] = b"TSDB: unknown ALIGN parameter";
const ALIGN_NO_AGG: &[u8] = b"TSDB: ALIGN parameter can only be used with AGGREGATION";
const ALIGN_START: &[u8] = b"TSDB: start alignment can only be used with explicit start timestamp";
const ALIGN_END: &[u8] = b"TSDB: end alignment can only be used with explicit end timestamp";
const FILTER_VALUE_MISSING: &[u8] = b"TSDB: FILTER_BY_VALUE one or more arguments are missing";
const BAD_MIN: &[u8] = b"TSDB: Couldn't parse MIN";
const BAD_MAX: &[u8] = b"TSDB: Couldn't parse MAX";
const FILTER_TS_MISSING: &[u8] = b"TSDB: FILTER_BY_TS one or more arguments are missing";
const TOO_WIDE: &[u8] = b"TSDB: the requested range holds too many empty buckets";
const BAD_FILTER: &[u8] = b"TSDB: failed parsing labels";
const NO_MATCHER: &[u8] = b"TSDB: please provide at least one matcher";
const NO_EXPRESSIONS: &[u8] = b"TSDB: FILTER given with no filter expressions";
const BAD_SUBTYPE: &[u8] = b"TSDB: unknown subtype, must be one of LABELS|VALUES";
const EXPECTED_FILTER: &[u8] = b"TSDB: unknown argument, expected FILTER";
const BOTH_LABELS: &[u8] = b"TSDB: cannot accept WITHLABELS and SELECT_LABELS together";
const NO_SELECTED: &[u8] = b"TSDB: SELECT_LABELS should have at least 1 parameter";
const RULE_EXISTS: &[u8] = b"TSDB: the destination key already has a src rule";
const SOURCE_IS_DEST: &[u8] = b"TSDB: the source key already has a source rule";
const DEST_IS_SOURCE: &[u8] = b"TSDB: the destination key already has a dst rule";
const SAME_KEY: &[u8] = b"TSDB: the source key and destination key should be different";
const NO_RULE: &[u8] = b"TSDB: compaction rule does not exist";
const BAD_ALIGN_STAMP: &[u8] = b"TSDB: Couldn't parse alignTimestamp";
const THIRD_WORD: &[u8] = b"TSDB: wrong 3rd argument";
const BAD_NUMKEYS: &[u8] = b"TSDB: numkeys must be a positive integer";
const AGG_NUMKEYS: &[u8] = b"TSDB: the number of AGGREGATION arguments must be equal to numkeys";
const BAD_READ_AT: &[u8] = b"TSDB: invalid timestamp";
const MISSING_FILTER: &[u8] = b"TSDB: missing FILTER argument";
const NO_FILTER_LABELS: &[u8] = b"TSDB: missing labels for filter argument";
const GROUPBY_ORDER: &[u8] = b"TSDB: GROUPBY should always come after filter";
const BAD_REDUCER: &[u8] = b"TSDB: Invalid reducer type";
const GROUPBY_COLUMNS: &[u8] =
b"TSDB: GROUPBY is not allowed when multiple aggregators are specified";
const BARE_RETENTION: &[u8] = b"TSDB: Couldn't parse RETENTION";
const BARE_BACKWARDS: &[u8] =
b"TSDB: timestamp must be equal to or higher than the maximum existing timestamp";
const WRONG_KIND: &str = "WRONGTYPE Operation against a key holding the wrong kind of value";
const CHUNK_MIN: i64 = 48;
const CHUNK_MAX: i64 = 1_048_576;
const MAX_AGGS: usize = 16;
const MAX_FILTER_TS: usize = 128;
const SPAN_AT: usize = 2;
#[derive(Debug)]
pub(super) struct TsBody {
s: Series,
}
impl Foreign for TsBody {
fn type_name(&self) -> &'static str {
"TSDB-TYPE"
}
fn encoding(&self) -> &'static str {
"raw"
}
fn memory_bytes(&self) -> usize {
self.s.memory_bytes()
}
fn is_empty(&self) -> bool {
false
}
}
pub(super) fn execute(db: &Db, spec: &Spec, args: Args<'_>, out: &mut Out) -> Result<()> {
match spec.name {
"ts.create" => create(db, &args, out),
"ts.alter" => alter(db, &args, out),
"ts.add" => add(db, &args, out),
"ts.madd" => madd(db, &args, out),
"ts.incrby" => incr(db, &args, out, true),
"ts.decrby" => incr(db, &args, out, false),
"ts.del" => del(db, &args, out),
"ts.get" => get(db, &args, out),
"ts.range" => range(db, &args, out, false),
"ts.revrange" => range(db, &args, out, true),
"ts.nrange" => nrange(db, &args, out, false),
"ts.nrevrange" => nrange(db, &args, out, true),
"ts.read" => tail(db, &args, out),
"ts.queryindex" => queryindex(db, &args, out),
"ts.querylabels" => querylabels(db, &args, out),
"ts.mget" => mget(db, &args, out),
"ts.createrule" => createrule(db, &args, out),
"ts.deleterule" => deleterule(db, &args, out),
"ts.mrange" => mrange(db, &args, out, false),
"ts.mrevrange" => mrange(db, &args, out, true),
"ts.info" => info(db, &args, out),
other => unreachable!("{other} is not a time series command"),
}
}
fn create(db: &Db, args: &Args<'_>, out: &mut Out) -> Result<()> {
let opts = match options(args) {
Ok(opts) => opts,
Err(bad) => return said(bad, "ts.create", out),
};
let key = args.get(1);
if db.hold(key).kind_of(key).is_some() {
return say(out, EXISTS);
}
let mut s = Series::new();
apply(&mut s, opts);
db.hold(key).put_foreign(key, Box::new(TsBody { s }));
out.ok();
Ok(())
}
fn alter(db: &Db, args: &Args<'_>, out: &mut Out) -> Result<()> {
let mut opts = match options(args) {
Ok(opts) => opts,
Err(bad) => return said(bad, "ts.alter", out),
};
opts.encoding = None;
let key = args.get(1);
let mut stripe = db.hold(key);
let Some(body) = write(&mut stripe, key)? else {
return say(out, MISSING);
};
apply(&mut body.s, opts);
out.ok();
Ok(())
}
fn add(db: &Db, args: &Args<'_>, out: &mut Out) -> Result<()> {
let key = args.get(1);
let Some(value) = number(args.get(3)) else {
return say(out, BAD_VALUE);
};
let at = match moment(db, args.get(2)) {
Ok(at) => at,
Err(msg) => return say(out, msg),
};
let over = if db.hold(key).kind_of(key).is_none() {
let opts = match options(args) {
Ok(opts) => opts,
Err(bad) => return said(bad, "ts.add", out),
};
let mut s = Series::new();
apply(&mut s, opts);
db.hold(key).put_foreign(key, Box::new(TsBody { s }));
None
} else {
if write(&mut db.hold(key), key).is_err() {
return say(out, NOT_A_SERIES);
}
match find(args, b"ON_DUPLICATE") {
None => None,
Some(at) => match policy_at(args, at) {
Ok(policy) => Some(policy),
Err(bad) => return said(bad, "ts.add", out),
},
}
};
let stored = {
let mut stripe = db.hold(key);
let body = write(&mut stripe, key)?.expect("the series is there by now");
let before = body.s.last();
store(body, at, value, over, out).then_some(before)
};
if let Some(before) = stored {
feed(db, key, at, before)?;
}
Ok(())
}
fn madd(db: &Db, args: &Args<'_>, out: &mut Out) -> Result<()> {
if args.len() < 4 || !(args.len() - 1).is_multiple_of(3) {
return Err(args::wrong_arity("ts.madd"));
}
out.array((args.len() - 1) / 3);
for i in (1..args.len()).step_by(3) {
let Some(value) = number(args.get(i + 2)) else {
out.error_line(b"ERR ", BAD_VALUE);
continue;
};
let at = match moment(db, args.get(i + 1)) {
Ok(at) => at,
Err(msg) => {
out.error_line(b"ERR ", msg);
continue;
}
};
let key = args.get(i);
let stored = {
let mut stripe = db.hold(key);
match write(&mut stripe, key) {
Ok(Some(body)) => {
let before = body.s.last();
store(body, at, value, None, out).then_some(before)
}
Ok(None) | Err(_) => {
out.error_line(b"ERR ", NOT_A_SERIES);
None
}
}
};
if let Some(before) = stored {
feed(db, key, at, before)?;
}
}
Ok(())
}
fn incr(db: &Db, args: &Args<'_>, out: &mut Out, up: bool) -> Result<()> {
let key = args.get(1);
let name = if up { "ts.incrby" } else { "ts.decrby" };
let exists = db.hold(key).kind_of(key).is_some();
if exists {
write(&mut db.hold(key), key)?;
}
let Some(by) = parse_f64(args.get(2)).filter(|n| !n.is_nan()) else {
return say(out, BAD_INCREMENT);
};
let labels_at = find_from(args, 3, b"LABELS");
let stamp_at =
find_from(args, 3, b"TIMESTAMP").filter(|&at| labels_at.is_none_or(|labels| at < labels));
let at = match stamp_at {
None => now(db),
Some(at) => match args.opt(at + 1) {
None => return say(out, BAD_TIMESTAMP),
Some(b"*") => now(db),
Some(word) => match parse_i64(word) {
Some(at) => at,
None => return say(out, BAD_TIMESTAMP),
},
},
};
if !exists {
let opts = match options(args) {
Ok(opts) => opts,
Err(bad) => return said(bad, name, out),
};
let mut s = Series::new();
apply(&mut s, opts);
db.hold(key).put_foreign(key, Box::new(TsBody { s }));
}
let stored = {
let mut stripe = db.hold(key);
let body = write(&mut stripe, key)?.expect("the series is there by now");
let last = body.s.last_sample();
if last.is_some_and(|s| at < s.at) {
out.error(BARE_BACKWARDS);
return Ok(());
}
let base = last.map_or(0.0, |s| s.value);
if base.is_nan() {
return say(out, NAN_INCREMENT);
}
let value = if up { base + by } else { base - by };
let before = body.s.last();
store(body, at, value, Some(Policy::Last), out).then_some(before)
};
if let Some(before) = stored {
feed(db, key, at, before)?;
}
Ok(())
}
fn del(db: &Db, args: &Args<'_>, out: &mut Out) -> Result<()> {
let Some(from) = span(args.get(2), b"-", 0) else {
return say(out, BAD_FROM);
};
let Some(to) = span(args.get(3), b"+", i64::MAX) else {
return say(out, BAD_TO);
};
let key = args.get(1);
let gone = {
let mut stripe = db.hold(key);
let Some(body) = write(&mut stripe, key)? else {
return say(out, MISSING);
};
body.s.delete(from, to)
};
undo(db, key, from, to)?;
out.uint(gone as u64);
Ok(())
}
fn get(db: &Db, args: &Args<'_>, out: &mut Out) -> Result<()> {
if args.len() > 3 {
return Err(args::wrong_arity("ts.get"));
}
let latest = args.opt(2).is_some_and(|word| args::is(word, b"LATEST"));
let open = if latest {
open_bucket(db, args.get(1))?
} else {
None
};
let key = args.get(1);
let mut stripe = db.hold(key);
let Some(body) = read(&mut stripe, key)? else {
return say(out, MISSING);
};
if args.len() == 3 && !latest {
return say(out, THIRD_WORD);
}
match open.or_else(|| body.s.last_sample()) {
None => out.array(0),
Some(sample) => {
out.array(2);
out.int(sample.at);
value(out, sample.value);
}
}
Ok(())
}
fn range(db: &Db, args: &Args<'_>, out: &mut Out, reverse: bool) -> Result<()> {
let name = if reverse { "ts.revrange" } else { "ts.range" };
let key = args.get(1);
if read(&mut db.hold(key), key)?.is_none() {
return say(out, MISSING);
}
let mut query = match reading(args, reverse, SPAN_AT) {
Ok(query) => query,
Err(bad) => return said(bad, name, out),
};
if find_from(args, SPAN_AT + 1, b"LATEST").is_some() {
query.latest = open_bucket(db, key)?;
}
let mut stripe = db.hold(key);
let body = read(&mut stripe, key)?.expect("the series was there a moment ago");
let rows = match body.s.read(&query) {
Ok(rows) => rows,
Err(Unread::TooWide) => return say(out, TOO_WIDE),
};
spread(out, &rows);
Ok(())
}
fn nrange(db: &Db, args: &Args<'_>, out: &mut Out, reverse: bool) -> Result<()> {
let name = if reverse { "ts.nrevrange" } else { "ts.nrange" };
let Some(keys) = parse_i64(args.get(1))
.filter(|&n| n > 0)
.and_then(|n| usize::try_from(n).ok())
else {
return say(out, BAD_NUMKEYS);
};
let span = 2 + keys;
if span + 1 >= args.len() {
return Err(args::wrong_arity(name));
}
if keys > 1
&& let Some(at) = find_from(args, span + 1, b"AGGREGATION")
&& let Err(bad) = agg_lists(args, at, keys)
{
return said(bad, name, out);
}
let (query, lists) = match reading_keys(args, reverse, span, keys) {
Ok(read) => read,
Err(bad) => return said(bad, name, out),
};
for i in 0..keys {
let key = args.get(2 + i);
if read(&mut db.hold(key), key)?.is_none() {
return say(out, MISSING);
}
}
let latest = find_from(args, span + 1, b"LATEST").is_some();
let mut taken = Vec::with_capacity(keys);
for (i, list) in lists.iter().enumerate() {
let key = args.get(2 + i).to_vec();
let open = if latest { open_bucket(db, &key)? } else { None };
let mut one = Query {
reverse: false,
count: None,
latest: open,
..query.clone()
};
if let Some(buckets) = one.buckets.as_mut() {
buckets.aggs.clone_from(list);
}
let mut stripe = db.hold(&key);
let body = read(&mut stripe, &key)?.expect("the series was there a moment ago");
match body.s.read(&one) {
Ok(rows) => taken.push(rows),
Err(Unread::TooWide) => return say(out, TOO_WIDE),
}
}
let mut rows = join(&taken);
if reverse {
rows.flip();
}
if let Some(n) = query.count {
rows.keep(n);
}
columns(out, &rows);
Ok(())
}
fn join(taken: &[Rows]) -> Rows {
let width = taken.iter().map(|rows| rows.width).sum();
let mut rows = Rows {
width,
..Rows::default()
};
let mut at = vec![0usize; taken.len()];
while let Some(now) = taken
.iter()
.zip(&at)
.filter_map(|(one, &i)| one.stamps.get(i).copied())
.min()
{
rows.stamps.push(now);
for (k, one) in taken.iter().enumerate() {
if one.stamps.get(at[k]) == Some(&now) {
rows.values.extend_from_slice(one.row(at[k]));
at[k] += 1;
} else {
rows.values
.extend(core::iter::repeat_n(f64::NAN, one.width));
}
}
}
rows
}
fn tail(db: &Db, args: &Args<'_>, out: &mut Out) -> Result<()> {
if args.len() != 3 {
return Err(args::wrong_arity("ts.read"));
}
let word = args.get(2);
let newest = word == b"+";
let from = if word == b"-" || newest {
0
} else {
match parse_i64(word).filter(|&n| n >= 0) {
Some(at) => at,
None => return said(Bad::Bare(BAD_READ_AT), "ts.read", out),
}
};
let key = args.get(1);
let mut stripe = db.hold(key);
let Some(body) = bare_read(&mut stripe, key)? else {
out.array(0);
return Ok(());
};
let from = match (newest, body.s.last_sample()) {
(true, None) => {
out.array(0);
return Ok(());
}
(true, Some(last)) => last.at,
(false, _) => from,
};
let query = Query {
from,
to: i64::MAX,
..Query::default()
};
let rows = body
.s
.read(&query)
.expect("a read with no buckets fills no gaps");
spread(out, &rows);
Ok(())
}
fn reading(args: &Args<'_>, reverse: bool, at: usize) -> core::result::Result<Query, Bad> {
reading_keys(args, reverse, at, 1).map(|(query, _)| query)
}
fn reading_keys(
args: &Args<'_>,
reverse: bool,
at: usize,
keys: usize,
) -> core::result::Result<(Query, Vec<Vec<Agg>>), Bad> {
let opts = at + 1;
let open_start = args.get(at) == b"-";
let Some(from) = span(args.get(at), b"-", 0) else {
return Err(Bad::Said(BAD_FROM));
};
let open_end = args.get(at + 1) == b"+";
let Some(to) = span(args.get(at + 1), b"+", i64::MAX) else {
return Err(Bad::Said(BAD_TO));
};
let count = count_of(args, opts, keys)?;
let read = buckets_of(args, opts, keys)?;
let lists = read
.as_ref()
.map_or_else(|| vec![Vec::new(); keys], |(_, lists)| lists.clone());
let mut buckets = read.map(|(buckets, _)| buckets);
if let Some(align) = align_of(
args,
opts,
buckets.is_some(),
open_start,
open_end,
from,
to,
)? && let Some(buckets) = buckets.as_mut()
{
buckets.align = align;
}
Ok((
Query {
from,
to,
reverse,
count,
by_ts: by_ts(args, opts)?,
by_value: by_value(args, opts)?,
buckets,
latest: None,
},
lists,
))
}
fn count_of(args: &Args<'_>, opts: usize, keys: usize) -> core::result::Result<Option<usize>, Bad> {
let Some(mut at) = find_from(args, opts, b"COUNT") else {
return Ok(None);
};
while find_from(args, opts, b"AGGREGATION").is_some_and(|agg| at > agg && at <= agg + keys)
|| find_from(args, opts, b"REDUCE") == Some(at - 1)
{
match find_from(args, at + 1, b"COUNT") {
Some(next) => at = next,
None => return Ok(None),
}
}
if at + 1 == args.len() {
return Err(Bad::Said(COUNT_MISSING));
}
let Some(n) = parse_i64(args.get(at + 1)) else {
return Err(Bad::Said(BAD_COUNT));
};
if n < 1 {
return Err(Bad::Said(COUNT_RANGE));
}
Ok(Some(usize::try_from(n).unwrap_or(usize::MAX)))
}
type Bucketing = (Buckets, Vec<Vec<Agg>>);
fn buckets_of(
args: &Args<'_>,
opts: usize,
keys: usize,
) -> core::result::Result<Option<Bucketing>, Bad> {
let Some(at) = find_from(args, opts, b"AGGREGATION") else {
return Ok(None);
};
let width = at + 1 + keys;
if width >= args.len() {
return Err(Bad::Said(BAD_AGGREGATION));
}
let Some(delta) = parse_i64(args.get(width)) else {
return Err(Bad::Said(BAD_AGGREGATION));
};
let lists = agg_lists(args, at, keys)?;
if delta <= 0 {
return Err(Bad::Said(BAD_BUCKET));
}
let mut buckets = Buckets {
aggs: lists[0].clone(),
delta,
align: 0,
empty: false,
stamp: Stamp::Start,
};
if let Some(flag) = find_from(args, opts, b"EMPTY") {
if flag != width + 1 && flag != width + 3 {
return Err(Bad::Said(EMPTY_PLACE));
}
buckets.empty = true;
}
if let Some(flag) = find_from(args, opts, b"BUCKETTIMESTAMP") {
if flag != width + 1 && flag != width + 2 {
return Err(Bad::Said(BUCKET_TS_PLACE));
}
if flag + 1 >= args.len() {
return Err(Bad::Arity);
}
let word = args.get(flag + 1);
buckets.stamp = if args::is(word, b"start") || word == b"-" {
Stamp::Start
} else if args::is(word, b"end") || word == b"+" {
Stamp::End
} else if args::is(word, b"mid") || word == b"~" {
Stamp::Mid
} else {
return Err(Bad::Said(BAD_BUCKET_TS));
};
}
Ok(Some((buckets, lists)))
}
fn agg_lists(args: &Args<'_>, at: usize, keys: usize) -> core::result::Result<Vec<Vec<Agg>>, Bad> {
if keys == 1 {
return Ok(vec![reductions(args.get(at + 1))?]);
}
let mut lists = Vec::with_capacity(keys);
for i in 0..keys {
let Some(word) = args.opt(at + 1 + i) else {
return Err(Bad::Said(AGG_NUMKEYS));
};
if parse_i64(word).is_some() {
return Err(Bad::Said(AGG_NUMKEYS));
}
lists.push(reductions(word)?);
}
if args
.opt(at + 1 + keys)
.is_some_and(|word| reductions(word).is_ok())
{
return Err(Bad::Said(AGG_NUMKEYS));
}
Ok(lists)
}
fn reductions(spec: &[u8]) -> core::result::Result<Vec<Agg>, Bad> {
let mut aggs = Vec::new();
for word in spec.split(|&b| b == b',') {
if word.is_empty() {
return Err(Bad::Said(EMPTY_AGG));
}
if aggs.len() >= MAX_AGGS {
return Err(Bad::Said(TOO_MANY_AGGS));
}
let Some(agg) = Agg::parse(word) else {
return Err(Bad::Said(UNKNOWN_AGG));
};
aggs.push(agg);
}
Ok(aggs)
}
fn align_of(
args: &Args<'_>,
opts: usize,
bucketed: bool,
open_start: bool,
open_end: bool,
from: i64,
to: i64,
) -> core::result::Result<Option<i64>, Bad> {
let Some(at) = find_from(args, opts, b"ALIGN") else {
return Ok(None);
};
if at + 1 >= args.len() {
return Err(Bad::Arity);
}
let word = args.get(at + 1);
let start = args::is(word, b"start") || word == b"-";
let end = args::is(word, b"end") || word == b"+";
let align = if start {
from
} else if end {
to
} else {
match parse_i64(word).filter(|&n| n >= 0) {
Some(n) => n,
None => return Err(Bad::Said(BAD_ALIGN)),
}
};
if !bucketed {
return Err(Bad::Said(ALIGN_NO_AGG));
}
if start && open_start {
return Err(Bad::Said(ALIGN_START));
}
if end && open_end {
return Err(Bad::Said(ALIGN_END));
}
Ok(Some(align))
}
fn by_value(args: &Args<'_>, opts: usize) -> core::result::Result<Option<(f64, f64)>, Bad> {
let Some(at) = find_from(args, opts, b"FILTER_BY_VALUE") else {
return Ok(None);
};
if at + 2 >= args.len() {
return Err(Bad::Said(FILTER_VALUE_MISSING));
}
let Some(min) = parse_f64(args.get(at + 1)) else {
return Err(Bad::Said(BAD_MIN));
};
let Some(max) = parse_f64(args.get(at + 2)) else {
return Err(Bad::Said(BAD_MAX));
};
Ok(Some((min, max)))
}
fn by_ts(args: &Args<'_>, opts: usize) -> core::result::Result<Option<Vec<i64>>, Bad> {
let Some(at) = find_from(args, opts, b"FILTER_BY_TS") else {
return Ok(None);
};
if at + 1 == args.len() {
return Err(Bad::Said(FILTER_TS_MISSING));
}
let mut list = Vec::new();
let mut i = at + 1;
while i < args.len() && list.len() < MAX_FILTER_TS {
match parse_i64(args.get(i)).filter(|&n| n >= 0) {
Some(n) => list.push(n),
None => break,
}
i += 1;
}
if list.is_empty() {
return Err(Bad::Said(FILTER_TS_MISSING));
}
list.sort_unstable();
list.dedup();
Ok(Some(list))
}
#[derive(Debug)]
enum Rule {
Is(Vec<u8>, Vec<u8>),
IsNot(Vec<u8>, Vec<u8>),
In(Vec<u8>, Vec<Vec<u8>>),
NotIn(Vec<u8>, Vec<Vec<u8>>),
Present(Vec<u8>),
Absent(Vec<u8>),
}
impl Rule {
fn positive(&self) -> bool {
matches!(self, Rule::Is(..) | Rule::In(..))
}
fn holds(&self, labels: &[(Vec<u8>, Vec<u8>)]) -> bool {
let any = |name: &Vec<u8>, take: &dyn Fn(&Vec<u8>) -> bool| {
labels.iter().any(|(n, v)| n == name && take(v))
};
match self {
Rule::Is(name, want) => any(name, &|v| v == want),
Rule::IsNot(name, want) => !any(name, &|v| v == want),
Rule::In(name, list) => any(name, &|v| list.contains(v)),
Rule::NotIn(name, list) => !any(name, &|v| list.contains(v)),
Rule::Present(name) => any(name, &|_| true),
Rule::Absent(name) => !any(name, &|_| true),
}
}
}
fn rule(arg: &[u8]) -> core::result::Result<Rule, Bad> {
let negated = arg.windows(2).any(|w| w == b"!=");
let separator: &[u8] = if negated {
b"!="
} else if arg.contains(&b'=') {
b"="
} else {
return Err(Bad::Said(BAD_FILTER));
};
if let Some(open) = arg.iter().position(|&b| b == b'(')
&& open > 0
&& arg[open - 1] == b'='
{
if arg.last() != Some(&b')') {
return Err(Bad::Said(BAD_FILTER));
}
let mut at = 0;
let Some(name) = word(&arg[..open], separator, &mut at) else {
return Err(Bad::Said(BAD_FILTER));
};
let inside = &arg[open + 1..arg.len() - 1];
let mut list = Vec::new();
if !inside.is_empty() {
for part in inside.split(|&b| b == b',') {
if part.is_empty() {
return Err(Bad::Said(BAD_FILTER));
}
list.push(part.to_vec());
}
}
let name = name.to_vec();
return Ok(if negated {
Rule::NotIn(name, list)
} else {
Rule::In(name, list)
});
}
let mut at = 0;
let Some(name) = word(arg, separator, &mut at) else {
return Err(Bad::Said(BAD_FILTER));
};
let name = name.to_vec();
match word(arg, separator, &mut at) {
Some(value) => {
let value = value.to_vec();
Ok(if negated {
Rule::IsNot(name, value)
} else {
Rule::Is(name, value)
})
}
None if negated => Ok(Rule::Present(name)),
None => Ok(Rule::Absent(name)),
}
}
fn word<'a>(arg: &'a [u8], separator: &[u8], at: &mut usize) -> Option<&'a [u8]> {
while *at < arg.len() && separator.contains(&arg[*at]) {
*at += 1;
}
let start = *at;
while *at < arg.len() && !separator.contains(&arg[*at]) {
*at += 1;
}
(start < *at).then(|| &arg[start..*at])
}
fn rules(args: &Args<'_>, from: usize) -> core::result::Result<Vec<Rule>, Bad> {
rules_in(args, from, args.len())
}
fn rules_in(args: &Args<'_>, from: usize, to: usize) -> core::result::Result<Vec<Rule>, Bad> {
let mut list = Vec::with_capacity(to.saturating_sub(from));
for i in from..to {
list.push(rule(args.get(i))?);
}
if !list.iter().any(Rule::positive) {
return Err(Bad::Said(NO_MATCHER));
}
Ok(list)
}
fn chosen(db: &Db, rules: &[Rule]) -> Vec<Vec<u8>> {
let mut names = Vec::new();
db.scan(KeyCursor::START, usize::MAX, Some(Kind::Foreign), |key| {
names.push(key.to_vec());
});
names.retain(|name| {
db.hold(name)
.foreign(name)
.ok()
.flatten()
.and_then(<dyn Foreign>::downcast_ref::<TsBody>)
.is_some_and(|body| rules.iter().all(|r| r.holds(body.s.labels())))
});
names.sort_unstable();
names
}
fn queryindex(db: &Db, args: &Args<'_>, out: &mut Out) -> Result<()> {
let rules = match rules(args, 1) {
Ok(rules) => rules,
Err(bad) => return said(bad, "ts.queryindex", out),
};
let names = chosen(db, &rules);
out.set(names.len());
for name in &names {
out.bulk(name);
}
Ok(())
}
fn querylabels(db: &Db, args: &Args<'_>, out: &mut Out) -> Result<()> {
let values = if args::is(args.get(1), b"LABELS") {
false
} else if args::is(args.get(1), b"VALUES") {
true
} else {
return say(out, BAD_SUBTYPE);
};
let at = if values { 3 } else { 2 };
if args.len() < at {
return Err(args::wrong_arity("ts.querylabels"));
}
let mut rules = Vec::new();
if args.len() > at {
if !args::is(args.get(at), b"FILTER") {
return say(out, EXPECTED_FILTER);
}
if args.len() == at + 1 {
return say(out, NO_EXPRESSIONS);
}
rules = match self::rules(args, at + 1) {
Ok(rules) => rules,
Err(bad) => return said(bad, "ts.querylabels", out),
};
}
let wanted = if values {
args.get(2).to_vec()
} else {
Vec::new()
};
let mut found: Vec<Vec<u8>> = Vec::new();
for name in chosen(db, &rules) {
let mut stripe = db.hold(&name);
let Some(body) = read(&mut stripe, &name)? else {
continue;
};
if values {
let smallest = body
.s
.labels()
.iter()
.filter(|(label, _)| *label == wanted)
.map(|(_, value)| value)
.min();
if let Some(value) = smallest {
found.push(value.clone());
}
} else {
for (label, _) in body.s.labels() {
found.push(label.clone());
}
}
}
found.sort_unstable();
found.dedup();
out.set(found.len());
for word in &found {
out.bulk(word);
}
Ok(())
}
#[derive(Default)]
enum Wearing {
#[default]
None,
All,
These(Vec<Vec<u8>>),
}
fn mget(db: &Db, args: &Args<'_>, out: &mut Out) -> Result<()> {
let Some(filter) = find(args, b"FILTER") else {
return Err(args::wrong_arity("ts.mget"));
};
let wearing = match wearing(args) {
Ok(wearing) => wearing,
Err(bad) => return said(bad, "ts.mget", out),
};
let rules = match rules(args, filter + 1) {
Ok(rules) => rules,
Err(bad) => return said(bad, "ts.mget", out),
};
let names = chosen(db, &rules);
let latest = latest(args);
if out.proto().is_resp3() {
out.map(names.len());
} else {
out.array(names.len());
}
for name in &names {
let open = if latest { open_bucket(db, name)? } else { None };
let mut stripe = db.hold(name);
let Some(body) = read(&mut stripe, name)? else {
continue;
};
let last = open.or_else(|| body.s.last_sample());
let labels = body.s.labels().to_vec();
if out.proto().is_resp3() {
out.bulk(name);
out.array(2);
} else {
out.array(3);
out.bulk(name);
}
wearing.write(out, &labels);
match last {
None => out.array(0),
Some(sample) => {
out.array(2);
out.int(sample.at);
value(out, sample.value);
}
}
}
Ok(())
}
const ENDS_LABELS: &[&[u8]] = &[
b"FILTER",
b"FILTER_BY_TS",
b"FILTER_BY_VALUE",
b"COUNT",
b"AGGREGATION",
b"ALIGN",
b"WITHLABELS",
b"REDUCE",
b"GROUPBY",
];
fn latest(args: &Args<'_>) -> bool {
let stop = find(args, b"FILTER").unwrap_or_else(|| args.len());
find(args, b"LATEST").is_some_and(|at| at < stop)
}
fn wearing(args: &Args<'_>) -> core::result::Result<Wearing, Bad> {
let all = find(args, b"WITHLABELS").is_some();
let some = find(args, b"SELECTED_LABELS");
if all && some.is_some() {
return Err(Bad::Said(BOTH_LABELS));
}
if let Some(at) = some {
let mut list: Vec<Vec<u8>> = Vec::new();
for i in at + 1..args.len() {
if ENDS_LABELS.iter().any(|word| args::is(args.get(i), word)) {
break;
}
list.push(args.get(i).to_vec());
}
if list.is_empty() {
return Err(Bad::Said(NO_SELECTED));
}
return Ok(Wearing::These(list));
}
Ok(if all { Wearing::All } else { Wearing::None })
}
impl Wearing {
fn write(&self, out: &mut Out, labels: &[(Vec<u8>, Vec<u8>)]) {
let pairs: Vec<(&[u8], Option<&[u8]>)> = match self {
Wearing::None => Vec::new(),
Wearing::All => labels
.iter()
.map(|(n, v)| (n.as_slice(), Some(v.as_slice())))
.collect(),
Wearing::These(wanted) => wanted
.iter()
.map(|name| {
let found = labels.iter().find(|(n, _)| n == name);
(name.as_slice(), found.map(|(_, v)| v.as_slice()))
})
.collect(),
};
let resp3 = out.proto().is_resp3();
if resp3 {
out.map(pairs.len());
} else {
out.array(pairs.len());
}
for (name, value) in pairs {
if !resp3 {
out.array(2);
}
out.bulk(name);
match value {
Some(value) => out.bulk(value),
None => out.nil(),
}
}
}
}
type Took = (Vec<u8>, Vec<(Vec<u8>, Vec<u8>)>, Rows);
type Group = (Vec<u8>, Vec<Vec<u8>>, Vec<Rows>);
const REDUCERS: &[Agg] = &[
Agg::Avg,
Agg::Sum,
Agg::Min,
Agg::Max,
Agg::Range,
Agg::Count,
Agg::StdP,
Agg::StdS,
Agg::VarP,
Agg::VarS,
];
fn mrange(db: &Db, args: &Args<'_>, out: &mut Out, reverse: bool) -> Result<()> {
let name = if reverse { "ts.mrevrange" } else { "ts.mrange" };
let query = match reading(args, reverse, SPAN_AT - 1) {
Ok(query) => query,
Err(bad) => return said(bad, name, out),
};
let Some(filter) = find(args, b"FILTER") else {
return say(out, MISSING_FILTER);
};
let wearing = match wearing(args) {
Ok(wearing) => wearing,
Err(bad) => return said(bad, name, out),
};
let group = find(args, b"GROUPBY");
if group.is_some_and(|at| at < filter) {
return say(out, GROUPBY_ORDER);
}
let end = group.unwrap_or_else(|| args.len());
if filter + 1 == end {
return say(out, NO_FILTER_LABELS);
}
let rules = match rules_in(args, filter + 1, end) {
Ok(rules) => rules,
Err(bad) => return said(bad, name, out),
};
if group.is_some_and(|at| at + 4 != args.len()) {
return Err(args::wrong_arity(name));
}
let reducer = match group {
None => None,
Some(_) => {
let word = find(args, b"REDUCE")
.filter(|at| at + 1 < args.len())
.map_or(&[][..], |at| args.get(at + 1));
match Agg::parse(word).filter(|agg| REDUCERS.contains(agg)) {
Some(agg) => Some(agg),
None => return say(out, BAD_REDUCER),
}
}
};
if reducer.is_some() && query.buckets.as_ref().is_some_and(|b| b.aggs.len() > 1) {
return say(out, GROUPBY_COLUMNS);
}
let names = chosen(db, &rules);
let latest = latest(args);
let mut query = query;
let mut taken: Vec<Took> = Vec::with_capacity(names.len());
for key in names {
query.latest = if latest { open_bucket(db, &key)? } else { None };
let mut stripe = db.hold(&key);
let Some(body) = read(&mut stripe, &key)? else {
continue;
};
let labels = body.s.labels().to_vec();
let rows = match body.s.read(&query) {
Ok(rows) => rows,
Err(Unread::TooWide) => return say(out, TOO_WIDE),
};
taken.push((key, labels, rows));
}
match (group, reducer) {
(Some(at), Some(agg)) => {
grouped(out, args.get(at + 1), agg, &wearing, &query, taken);
}
_ => plainly(out, &wearing, &query, &taken),
}
Ok(())
}
fn plainly(out: &mut Out, wearing: &Wearing, query: &Query, taken: &[Took]) {
let resp3 = out.proto().is_resp3();
if resp3 {
out.map(taken.len());
} else {
out.array(taken.len());
}
for (key, labels, rows) in taken {
if resp3 {
out.bulk(key);
out.array(3);
wearing.write(out, labels);
named(out, b"aggregators", query);
} else {
out.array(3);
out.bulk(key);
wearing.write(out, labels);
}
spread(out, rows);
}
}
fn named(out: &mut Out, under: &[u8], query: &Query) {
let aggs = query.buckets.as_ref().map_or(&[][..], |b| &b.aggs);
out.map(1);
out.bulk(under);
out.array(aggs.len());
for agg in aggs {
out.bulk(agg.name().as_bytes());
}
}
fn grouped(
out: &mut Out,
label: &[u8],
agg: Agg,
wearing: &Wearing,
query: &Query,
taken: Vec<Took>,
) {
let mut groups: Vec<Group> = Vec::new();
for (key, labels, rows) in taken {
let Some((_, value)) = labels.iter().find(|(n, _)| n == label) else {
continue;
};
let at = match groups.iter().position(|(v, _, _)| v == value) {
Some(at) => at,
None => {
groups.push((value.clone(), Vec::new(), Vec::new()));
groups.len() - 1
}
};
groups[at].1.push(key);
groups[at].2.push(rows);
}
groups.sort_by(|a, b| a.0.cmp(&b.0));
let resp3 = out.proto().is_resp3();
if resp3 {
out.map(groups.len());
} else {
out.array(groups.len());
}
for (value, sources, rows) in &groups {
let mut title = label.to_vec();
title.push(b'=');
title.extend_from_slice(value);
let mut labels = vec![(label.to_vec(), value.clone())];
if !resp3 && matches!(wearing, Wearing::All) {
labels.push((b"__reducer__".to_vec(), agg.name().as_bytes().to_vec()));
labels.push((b"__source__".to_vec(), sources.join(&b","[..])));
}
if resp3 {
out.bulk(&title);
out.array(4);
wearing.write(out, &labels);
out.map(1);
out.bulk(b"reducers");
out.array(1);
out.bulk(agg.name().as_bytes());
out.map(1);
out.bulk(b"sources");
out.array(sources.len());
for key in sources {
out.bulk(key);
}
} else {
out.array(3);
out.bulk(&title);
wearing.write(out, &labels);
}
spread(out, &folded(rows, agg, query));
}
}
fn folded(taken: &[Rows], agg: Agg, query: &Query) -> Rows {
let mut all: Vec<(i64, f64)> = Vec::new();
for rows in taken {
for i in 0..rows.len() {
all.push((rows.stamps[i], rows.row(i)[0]));
}
}
all.sort_by_key(|&(at, _)| at);
let mut stamps = Vec::new();
let mut values = Vec::new();
let mut at = 0;
while at < all.len() {
let moment = all[at].0;
let mut past = at;
while past < all.len() && all[past].0 == moment {
past += 1;
}
let readings: Vec<f64> = all[at..past].iter().map(|&(_, v)| v).collect();
stamps.push(moment);
values.push(yo_series::group(agg, &readings));
at = past;
}
if query.reverse {
stamps.reverse();
values.reverse();
}
if let Some(n) = query.count {
stamps.truncate(n);
values.truncate(n);
}
Rows {
stamps,
values,
width: 1,
}
}
fn spread(out: &mut Out, rows: &Rows) {
out.array(rows.len());
for i in 0..rows.len() {
let row = rows.row(i);
out.array(1 + row.len());
out.int(rows.stamps[i]);
for &d in row {
value(out, d);
}
}
}
fn columns(out: &mut Out, rows: &Rows) {
out.array(rows.len());
for i in 0..rows.len() {
let row = rows.row(i);
out.array(2);
out.int(rows.stamps[i]);
out.array(row.len());
for &d in row {
value(out, d);
}
}
}
fn createrule(db: &Db, args: &Args<'_>, out: &mut Out) -> Result<()> {
if !matches!(args.len(), 6 | 7) || !args::is(args.get(3), b"AGGREGATION") {
return Err(args::wrong_arity("ts.createrule"));
}
let Some(delta) = parse_i64(args.get(5)) else {
return say(out, BAD_AGGREGATION);
};
let Some(agg) = Agg::parse(args.get(4)) else {
return say(out, UNKNOWN_AGG);
};
if delta <= 0 {
return say(out, BAD_BUCKET);
}
let align = match args.opt(6) {
None => 0,
Some(word) => match parse_i64(word).filter(|&n| n >= 0) {
Some(n) => n,
None => return say(out, BAD_ALIGN_STAMP),
},
};
let source = args.get(1).to_vec();
let dest = args.get(2).to_vec();
if source == dest {
return say(out, SAME_KEY);
}
prune(db, &source)?;
prune(db, &dest)?;
{
let mut stripe = db.hold(&source);
let Some(body) = read(&mut stripe, &source)? else {
return say(out, MISSING);
};
if body.s.source().is_some() {
return say(out, SOURCE_IS_DEST);
}
}
{
let mut stripe = db.hold(&dest);
let Some(body) = write(&mut stripe, &dest)? else {
return say(out, MISSING);
};
if !body.s.rules().is_empty() {
return say(out, DEST_IS_SOURCE);
}
if body.s.source().is_some() {
return say(out, RULE_EXISTS);
}
body.s.set_source(Some(source.clone()));
}
let mut stripe = db.hold(&source);
let body = write(&mut stripe, &source)?.expect("the source was there a moment ago");
body.s.add_rule(Compaction {
dest,
delta,
agg,
align,
open: None,
start: 0,
});
out.ok();
Ok(())
}
fn deleterule(db: &Db, args: &Args<'_>, out: &mut Out) -> Result<()> {
if args.len() != 3 {
return Err(args::wrong_arity("ts.deleterule"));
}
let source = args.get(1).to_vec();
let dest = args.get(2).to_vec();
prune(db, &source)?;
{
let mut stripe = db.hold(&source);
let Some(body) = write(&mut stripe, &source)? else {
return say(out, MISSING);
};
if !body.s.drop_rule(&dest) {
return say(out, NO_RULE);
}
}
let mut stripe = db.hold(&dest);
if let Some(body) = write(&mut stripe, &dest)? {
body.s.set_source(None);
}
out.ok();
Ok(())
}
fn prune(db: &Db, key: &[u8]) -> Result<()> {
let (source, dests) = {
let mut stripe = db.hold(key);
let Ok(Some(body)) = read(&mut stripe, key) else {
return Ok(());
};
let source = body.s.source().map(<[u8]>::to_vec);
let dests: Vec<Vec<u8>> = body
.s
.rules()
.iter()
.map(|rule| rule.dest.clone())
.collect();
(source, dests)
};
let mut kept = Vec::with_capacity(dests.len());
for dest in dests {
if points_at(db, &dest, key)? {
kept.push(dest);
}
}
let orphan = match &source {
None => false,
Some(source) => !feeds(db, source, key)?,
};
let mut stripe = db.hold(key);
let Some(body) = write(&mut stripe, key)? else {
return Ok(());
};
body.s.keep_rules(&kept);
if orphan {
body.s.set_source(None);
}
Ok(())
}
fn points_at(db: &Db, key: &[u8], source: &[u8]) -> Result<bool> {
let mut stripe = db.hold(key);
let body = match read(&mut stripe, key) {
Ok(body) => body,
Err(_) => return Ok(false),
};
Ok(body.is_some_and(|body| body.s.source() == Some(source)))
}
fn feeds(db: &Db, key: &[u8], dest: &[u8]) -> Result<bool> {
let mut stripe = db.hold(key);
let body = match read(&mut stripe, key) {
Ok(body) => body,
Err(_) => return Ok(false),
};
Ok(body.is_some_and(|body| body.s.rules().iter().any(|rule| rule.dest == dest)))
}
fn fold(s: &Series, rule: &Compaction, at: i64, from: i64) -> Option<f64> {
let wide = rule.agg == Agg::Twa;
let room = if wide { 2 } else { 1 };
let query = Query {
from: if wide { at } else { from },
to: at.saturating_add(rule.delta * room) - 1,
buckets: Some(Buckets {
aggs: vec![rule.agg],
delta: rule.delta,
align: rule.align,
empty: false,
stamp: Stamp::Start,
}),
..Query::default()
};
let rows = s.read(&query).ok()?;
let want = at.max(0);
(0..rows.len())
.find(|&i| rows.stamps[i] == want)
.map(|i| rows.row(i)[0])
}
fn feed(db: &Db, key: &[u8], at: i64, before: Option<i64>) -> Result<()> {
for rule in rules_of(db, key)? {
let (written, moved) = {
let mut stripe = db.hold(key);
let Some(body) = read(&mut stripe, key)? else {
return Ok(());
};
let landed = bucket_start(at, rule.delta, rule.align);
let mut written = None;
let mut moved = None;
if before.is_none_or(|last| at >= last) {
match rule.open {
Some(open) if landed > open => {
written = fold(&body.s, &rule, open, rule.start).map(|value| (open, value));
moved = Some((Some(landed), at));
}
Some(_) => {}
None => moved = Some((Some(landed), at)),
}
} else {
let newest = bucket_start(before.unwrap_or(at), rule.delta, rule.align);
if landed < newest {
written = fold(&body.s, &rule, landed, landed).map(|value| (landed, value));
} else if rule.open == Some(landed) {
moved = Some((Some(landed), landed));
}
}
(written, moved)
};
let mut stripe = db.hold(&rule.dest);
if let Some((at, value)) = written
&& let Some(dest) = write(&mut stripe, &rule.dest)?
{
let at = at.max(0);
let _ = dest.s.add(Sample::new(at, value), Some(Policy::Last));
}
drop(stripe);
let mut stripe = db.hold(key);
if let Some((open, start)) = moved
&& let Some(body) = write(&mut stripe, key)?
&& let Some(mine) = body.s.rule_mut(&rule.dest)
{
mine.open = open;
mine.start = start;
}
}
Ok(())
}
fn undo(db: &Db, key: &[u8], from: i64, to: i64) -> Result<()> {
for rule in rules_of(db, key)? {
let open = {
let mut stripe = db.hold(key);
let Some(body) = read(&mut stripe, key)? else {
return Ok(());
};
body.s
.last()
.map(|last| bucket_start(last, rule.delta, rule.align))
};
let head = bucket_start(from.max(0), rule.delta, rule.align);
let tail = to.saturating_add(rule.delta);
let stamps: Vec<i64> = {
let mut stripe = db.hold(&rule.dest);
let Some(dest) = read(&mut stripe, &rule.dest)? else {
continue;
};
dest.s.range(head, tail).map(|sample| sample.at).collect()
};
let rows = {
let mut stripe = db.hold(key);
let body = read(&mut stripe, key)?.expect("the source was there a moment ago");
let mut rows = Vec::with_capacity(stamps.len());
for at in stamps {
let value = match open {
Some(open) if at < open => fold(&body.s, &rule, at, at),
_ => None,
};
rows.push((at, value));
}
rows
};
{
let mut stripe = db.hold(&rule.dest);
for (at, value) in rows {
if let Some(dest) = write(&mut stripe, &rule.dest)? {
match value {
Some(value) => {
let _ = dest.s.add(Sample::new(at, value), Some(Policy::Last));
}
None => {
dest.s.delete(at, at);
}
}
}
}
if let Some(dest) = write(&mut stripe, &rule.dest)? {
match open {
Some(open) => dest.s.delete(open, i64::MAX),
None => dest.s.delete(0, i64::MAX),
};
}
}
let mut stripe = db.hold(key);
if let Some(body) = write(&mut stripe, key)?
&& let Some(mine) = body.s.rule_mut(&rule.dest)
{
mine.open = open;
mine.start = open.unwrap_or(0);
}
}
Ok(())
}
fn rules_of(db: &Db, key: &[u8]) -> Result<Vec<Compaction>> {
prune(db, key)?;
Ok(read(&mut db.hold(key), key)?.map_or_else(Vec::new, |body| body.s.rules().to_vec()))
}
fn open_bucket(db: &Db, key: &[u8]) -> Result<Option<Sample>> {
prune(db, key)?;
let source = {
let mut stripe = db.hold(key);
let Some(body) = read(&mut stripe, key)? else {
return Ok(None);
};
let Some(source) = body.s.source().map(<[u8]>::to_vec) else {
return Ok(None);
};
source
};
let mut stripe = db.hold(&source);
let Some(body) = read(&mut stripe, &source)? else {
return Ok(None);
};
let Some(rule) = body.s.rules().iter().find(|rule| rule.dest == key).cloned() else {
return Ok(None);
};
let Some(open) = rule.open else {
return Ok(None);
};
Ok(fold(&body.s, &rule, open, rule.start).map(|value| Sample::new(open.max(0), value)))
}
fn info(db: &Db, args: &Args<'_>, out: &mut Out) -> Result<()> {
let resp3 = out.proto().is_resp3();
let key = args.get(1);
prune(db, key)?;
let mut stripe = db.hold(key);
let Some(body) = read(&mut stripe, key)? else {
return say(out, MISSING);
};
let s = &body.s;
let (ignore_time, ignore_value) = s.ignore();
out.map(14);
out.simple(b"totalSamples");
out.uint(s.len() as u64);
out.simple(b"memoryUsage");
out.uint(s.memory_bytes() as u64);
out.simple(b"firstTimestamp");
out.int(s.first().unwrap_or(0));
out.simple(b"lastTimestamp");
out.int(s.last().unwrap_or(0));
out.simple(b"retentionTime");
out.int(s.retention());
out.simple(b"chunkCount");
out.uint(s.chunk_count() as u64);
out.simple(b"chunkSize");
out.uint(s.chunk_bytes() as u64);
out.simple(b"chunkType");
out.simple(s.encoding().name().as_bytes());
out.simple(b"duplicatePolicy");
out.simple(s.policy().unwrap_or(Policy::Block).name().as_bytes());
out.simple(b"labels");
if resp3 {
out.map(s.labels().len());
for (name, value) in s.labels() {
out.bulk(name);
out.bulk(value);
}
} else {
out.array(s.labels().len());
for (name, value) in s.labels() {
out.array(2);
out.bulk(name);
out.bulk(value);
}
}
out.simple(b"sourceKey");
match s.source() {
Some(key) => out.bulk(key),
None => out.nil(),
}
out.simple(b"rules");
if resp3 {
out.map(s.rules().len());
for rule in s.rules() {
out.bulk(&rule.dest);
out.array(3);
out.int(rule.delta);
out.simple(rule.agg.name().to_uppercase().as_bytes());
out.int(rule.align);
}
} else {
out.array(s.rules().len());
for rule in s.rules() {
out.array(4);
out.bulk(&rule.dest);
out.int(rule.delta);
out.simple(rule.agg.name().to_uppercase().as_bytes());
out.int(rule.align);
}
}
out.simple(b"ignoreMaxTimeDiff");
out.int(ignore_time);
out.simple(b"ignoreMaxValDiff");
out.double(ignore_value);
Ok(())
}
fn store(body: &mut TsBody, at: i64, value: f64, over: Option<Policy>, out: &mut Out) -> bool {
match body.s.add(Sample::new(at, value), over) {
Ok(when) => {
out.int(when);
when == at
}
Err(Refused::Old) => {
out.error_line(b"ERR ", TOO_OLD);
false
}
Err(Refused::Duplicate) => {
out.error_line(b"ERR ", UPSERT);
false
}
}
}
fn say(out: &mut Out, msg: &[u8]) -> Result<()> {
out.error_line(b"ERR ", msg);
Ok(())
}
#[derive(Debug, Default)]
struct Options {
retention: Option<i64>,
chunk_bytes: Option<usize>,
encoding: Option<Encoding>,
policy: Option<Policy>,
labels: Option<Vec<(Vec<u8>, Vec<u8>)>>,
ignore: Option<(i64, f64)>,
}
#[derive(Debug)]
enum Bad {
Said(&'static [u8]),
Bare(&'static [u8]),
Arity,
}
fn said(bad: Bad, name: &'static str, out: &mut Out) -> Result<()> {
match bad {
Bad::Said(msg) => say(out, msg),
Bad::Bare(msg) => {
out.error(msg);
Ok(())
}
Bad::Arity => Err(args::wrong_arity(name)),
}
}
fn options(args: &Args<'_>) -> core::result::Result<Options, Bad> {
let mut opts = Options::default();
if let Some(at) = find(args, b"LABELS") {
let first = at + 1;
let mut labels = Vec::new();
for i in 0..args.len().saturating_sub(first) / 2 {
let name = args.get(first + i * 2);
let value = args.get(first + i * 2 + 1);
let ok = !name.is_empty()
&& !value.is_empty()
&& !value.iter().any(|b| matches!(b, b'(' | b')' | b','));
if !ok {
return Err(Bad::Said(BAD_LABELS));
}
labels.push((name.to_vec(), value.to_vec()));
}
opts.labels = Some(labels);
}
if let Some(at) = find(args, b"RETENTION") {
let Some(n) = args.opt(at + 1).and_then(parse_i64) else {
return Err(Bad::Said(BAD_RETENTION));
};
if n < 0 {
return Err(Bad::Bare(BARE_RETENTION));
}
opts.retention = Some(n);
}
if let Some(at) = find(args, b"CHUNK_SIZE") {
let Some(n) = args.opt(at + 1).and_then(parse_i64) else {
return Err(Bad::Said(BAD_CHUNK));
};
if !(CHUNK_MIN..=CHUNK_MAX).contains(&n) || !(n as usize).is_multiple_of(8) {
return Err(Bad::Said(CHUNK_RANGE));
}
opts.chunk_bytes = Some(n as usize);
}
if let Some(at) = find(args, b"ENCODING") {
let Some(word) = args.opt(at + 1) else {
return Err(Bad::Arity);
};
opts.encoding = Some(if args::is(word, b"uncompressed") {
Encoding::Uncompressed
} else if args::is(word, b"compressed") {
Encoding::Compressed
} else {
return Err(Bad::Said(BAD_ENCODING));
});
}
if let Some(at) = find(args, b"DUPLICATE_POLICY") {
opts.policy = Some(policy_at(args, at)?);
}
if let Some(at) = find(args, b"IGNORE") {
let time = args.opt(at + 1).and_then(parse_i64);
let value = args.opt(at + 2).and_then(parse_f64);
let (Some(time), Some(value)) = (time, value) else {
return Err(Bad::Said(BAD_IGNORE));
};
if time < 0 || value < 0.0 {
return Err(Bad::Said(NEGATIVE_IGNORE));
}
opts.ignore = Some((time, value));
}
Ok(opts)
}
fn policy_at(args: &Args<'_>, at: usize) -> core::result::Result<Policy, Bad> {
let Some(word) = args.opt(at + 1) else {
return Err(Bad::Said(BAD_POLICY));
};
Policy::parse(word).ok_or(Bad::Said(UNKNOWN_POLICY))
}
fn apply(s: &mut Series, opts: Options) {
if let Some(n) = opts.retention {
s.set_retention(n);
}
if let Some(n) = opts.chunk_bytes {
s.set_chunk_bytes(n);
}
if let Some(encoding) = opts.encoding {
s.set_encoding(encoding);
}
if let Some(policy) = opts.policy {
s.set_policy(policy);
}
if let Some(labels) = opts.labels {
s.set_labels(labels);
}
if let Some((time, value)) = opts.ignore {
s.set_ignore(time, value);
}
}
fn find(args: &Args<'_>, word: &[u8]) -> Option<usize> {
find_from(args, 1, word)
}
fn find_from(args: &Args<'_>, from: usize, word: &[u8]) -> Option<usize> {
(from..args.len()).find(|&i| args::is(args.get(i), word))
}
fn span(arg: &[u8], open: &[u8], edge: i64) -> Option<i64> {
if arg == open {
return Some(edge);
}
parse_i64(arg).filter(|&n| n >= 0)
}
fn moment(db: &Db, arg: &[u8]) -> core::result::Result<i64, &'static [u8]> {
if arg == b"*" {
return Ok(now(db));
}
let Some(at) = parse_i64(arg) else {
return Err(BAD_TIMESTAMP);
};
if at < 0 {
return Err(NEGATIVE_TIMESTAMP);
}
Ok(at)
}
fn now(db: &Db) -> i64 {
db.now_ms() as i64
}
fn number(arg: &[u8]) -> Option<f64> {
if nan_word(arg) {
return Some(f64::NAN);
}
let mut at = usize::from(arg.first() == Some(&b'-'));
if !digits(arg, &mut at) {
return None;
}
if arg.get(at) == Some(&b'.') {
at += 1;
if !digits(arg, &mut at) {
return None;
}
}
if matches!(arg.get(at), Some(b'e' | b'E')) {
at += 1;
if matches!(arg.get(at), Some(b'+' | b'-')) {
at += 1;
}
if !digits(arg, &mut at) {
return None;
}
}
if at != arg.len() {
return None;
}
parse_f64(arg).filter(|value| value.is_finite())
}
fn digits(arg: &[u8], at: &mut usize) -> bool {
let start = *at;
while arg.get(*at).is_some_and(u8::is_ascii_digit) {
*at += 1;
}
*at > start
}
fn nan_word(arg: &[u8]) -> bool {
match arg.len() {
3 => arg.eq_ignore_ascii_case(b"nan"),
4 => arg.eq_ignore_ascii_case(b"-nan") || arg.eq_ignore_ascii_case(b"+nan"),
_ => false,
}
}
fn value(out: &mut Out, d: f64) {
if out.proto().is_resp3() {
out.double(d);
return;
}
let mut buf = [0u8; DOUBLE_MAX];
out.simple(write_dragonbox(&mut buf, d));
}
fn wrong_kind() -> Error {
Error::new(Code::Invalid, WRONG_KIND)
}
fn write<'k>(stripe: &'k mut Keyspace, key: &[u8]) -> Result<Option<&'k mut TsBody>> {
match stripe.foreign_mut(key) {
Ok(Some(body)) => match body.downcast_mut::<TsBody>() {
Some(body) => Ok(Some(body)),
None => Err(wrong_kind()),
},
Ok(None) => Ok(None),
Err(e) if e.code() == Code::WrongType => Err(wrong_kind()),
Err(e) => Err(e),
}
}
fn bare_read<'k>(stripe: &'k mut Keyspace, key: &[u8]) -> Result<Option<&'k TsBody>> {
match stripe.foreign(key)? {
Some(body) => match body.downcast_ref::<TsBody>() {
Some(body) => Ok(Some(body)),
None => Err(yo_kv::keyspace::wrong_type()),
},
None => Ok(None),
}
}
fn read<'k>(stripe: &'k mut Keyspace, key: &[u8]) -> Result<Option<&'k TsBody>> {
match stripe.foreign(key) {
Ok(Some(body)) => match body.downcast_ref::<TsBody>() {
Some(body) => Ok(Some(body)),
None => Err(wrong_kind()),
},
Ok(None) => Ok(None),
Err(e) if e.code() == Code::WrongType => Err(wrong_kind()),
Err(e) => Err(e),
}
}