use crate::{
client::{PreparedCommand, prepare_command},
resp::{Response, Value, cmd, serialize_flag},
};
use serde::{
Deserialize, Serialize,
de::{self, value::SeqAccessDeserializer},
};
use smallvec::SmallVec;
use std::{collections::HashMap, fmt};
pub trait TimeSeriesCommands<'a>: Sized {
#[must_use]
fn ts_add(
self,
key: impl Serialize,
timestamp: TsTimestamp,
value: f64,
options: TsAddOptions,
) -> PreparedCommand<'a, Self, u64> {
prepare_command(
self,
cmd("TS.ADD")
.key(key)
.arg(timestamp)
.arg(value)
.arg(options),
)
}
#[must_use]
fn ts_alter(
self,
key: impl Serialize,
options: TsCreateOptions,
) -> PreparedCommand<'a, Self, ()> {
prepare_command(self, cmd("TS.ALTER").key(key).arg(options))
}
#[must_use]
fn ts_create(
self,
key: impl Serialize,
options: TsCreateOptions,
) -> PreparedCommand<'a, Self, ()> {
prepare_command(self, cmd("TS.CREATE").key(key).arg(options))
}
#[must_use]
fn ts_createrule(
self,
src_key: impl Serialize,
dst_key: impl Serialize,
aggregator: TsAggregationType,
bucket_duration: u64,
options: TsCreateRuleOptions,
) -> PreparedCommand<'a, Self, ()> {
prepare_command(
self,
cmd("TS.CREATERULE")
.key(src_key)
.key(dst_key)
.arg("AGGREGATION")
.arg(aggregator)
.arg(bucket_duration)
.arg(options),
)
}
#[must_use]
fn ts_decrby(
self,
key: impl Serialize,
value: f64,
options: TsIncrByDecrByOptions,
) -> PreparedCommand<'a, Self, u64> {
prepare_command(self, cmd("TS.DECRBY").key(key).arg(value).arg(options))
}
#[must_use]
fn ts_del(
self,
key: impl Serialize,
from_timestamp: u64,
to_timestamp: u64,
) -> PreparedCommand<'a, Self, usize> {
prepare_command(
self,
cmd("TS.DEL").key(key).arg(from_timestamp).arg(to_timestamp),
)
}
#[must_use]
fn ts_deleterule(
self,
src_key: impl Serialize,
dst_key: impl Serialize,
) -> PreparedCommand<'a, Self, ()> {
prepare_command(self, cmd("TS.DELETERULE").key(src_key).key(dst_key))
}
#[must_use]
fn ts_get(
self,
key: impl Serialize,
options: TsGetOptions,
) -> PreparedCommand<'a, Self, TsGetResult> {
prepare_command(self, cmd("TS.GET").key(key).arg(options).readonly())
}
#[must_use]
fn ts_incrby(
self,
key: impl Serialize,
value: f64,
options: TsIncrByDecrByOptions,
) -> PreparedCommand<'a, Self, u64> {
prepare_command(self, cmd("TS.INCRBY").key(key).arg(value).arg(options))
}
#[must_use]
fn ts_info(self, key: impl Serialize, debug: bool) -> PreparedCommand<'a, Self, TsInfoResult> {
prepare_command(
self,
cmd("TS.INFO").key(key).arg_if(debug, "DEBUG").readonly(),
)
}
#[must_use]
fn ts_madd<R: Response>(self, items: impl Serialize) -> PreparedCommand<'a, Self, R> {
prepare_command(
self,
cmd("TS.MADD")
.key_with_step(items, 3)
.cluster_info(None, None, 3),
)
}
#[must_use]
fn ts_mget<R: Response>(
self,
options: TsMGetOptions,
filters: impl Serialize,
) -> PreparedCommand<'a, Self, R> {
prepare_command(
self,
cmd("TS.MGET")
.arg(options)
.arg("FILTER")
.arg(filters)
.readonly(),
)
}
#[must_use]
fn ts_mrange<R: Response>(
self,
from_timestamp: impl Serialize,
to_timestamp: impl Serialize,
options: TsMRangeOptions,
filters: impl Serialize,
groupby_options: TsGroupByOptions,
) -> PreparedCommand<'a, Self, R> {
prepare_command(
self,
cmd("TS.MRANGE")
.arg(from_timestamp)
.arg(to_timestamp)
.arg(options)
.arg("FILTER")
.arg(filters)
.arg(groupby_options)
.readonly(),
)
}
#[must_use]
fn ts_mrevrange<R: Response>(
self,
from_timestamp: impl Serialize,
to_timestamp: impl Serialize,
options: TsMRangeOptions,
filters: impl Serialize,
groupby_options: TsGroupByOptions,
) -> PreparedCommand<'a, Self, R> {
prepare_command(
self,
cmd("TS.MREVRANGE")
.arg(from_timestamp)
.arg(to_timestamp)
.arg(options)
.arg("FILTER")
.arg(filters)
.arg(groupby_options)
.readonly(),
)
}
#[must_use]
fn ts_queryindex<R: Response>(self, filters: impl Serialize) -> PreparedCommand<'a, Self, R> {
prepare_command(self, cmd("TS.QUERYINDEX").arg(filters).readonly())
}
#[must_use]
fn ts_range<R: Response>(
self,
key: impl Serialize,
from_timestamp: impl Serialize,
to_timestamp: impl Serialize,
options: TsRangeOptions,
) -> PreparedCommand<'a, Self, R> {
prepare_command(
self,
cmd("TS.RANGE")
.key(key)
.arg(from_timestamp)
.arg(to_timestamp)
.arg(options)
.readonly(),
)
}
#[must_use]
fn ts_revrange<R: Response>(
self,
key: impl Serialize,
from_timestamp: impl Serialize,
to_timestamp: impl Serialize,
options: TsRangeOptions,
) -> PreparedCommand<'a, Self, R> {
prepare_command(
self,
cmd("TS.REVRANGE")
.key(key)
.arg(from_timestamp)
.arg(to_timestamp)
.arg(options)
.readonly(),
)
}
}
#[derive(Default, Serialize)]
#[serde(rename_all = "UPPERCASE")]
pub struct TsAddOptions<'a> {
#[serde(skip_serializing_if = "Option::is_none")]
retention: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
encoding: Option<TsEncoding>,
#[serde(skip_serializing_if = "Option::is_none")]
chunk_size: Option<u32>,
#[serde(skip_serializing_if = "Option::is_none")]
on_duplicate: Option<TsDuplicatePolicy>,
#[serde(skip_serializing_if = "SmallVec::is_empty")]
labels: SmallVec<[(&'a str, &'a str); 10]>,
}
impl<'a> TsAddOptions<'a> {
#[must_use]
pub fn retention(mut self, retention_period: u64) -> Self {
self.retention = Some(retention_period);
self
}
#[must_use]
pub fn encoding(mut self, encoding: TsEncoding) -> Self {
self.encoding = Some(encoding);
self
}
#[must_use]
pub fn chunk_size(mut self, chunk_size: u32) -> Self {
self.chunk_size = Some(chunk_size);
self
}
#[must_use]
pub fn on_duplicate(mut self, policy: TsDuplicatePolicy) -> Self {
self.on_duplicate = Some(policy);
self
}
#[must_use]
pub fn labels(mut self, labels: impl IntoIterator<Item = (&'a str, &'a str)>) -> Self {
self.labels.extend(labels);
self
}
}
#[derive(Serialize)]
#[serde(rename_all = "UPPERCASE")]
#[non_exhaustive]
pub enum TsEncoding {
Compressed,
Uncompressed,
}
#[derive(Debug, Deserialize, Serialize)]
#[serde(rename_all = "lowercase")]
#[non_exhaustive]
pub enum TsDuplicatePolicy {
Block,
First,
Last,
Min,
Max,
Sum,
}
#[derive(Default, Serialize)]
#[serde(rename_all = "UPPERCASE")]
pub struct TsCreateOptions<'a> {
#[serde(skip_serializing_if = "Option::is_none")]
retention: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
encoding: Option<TsEncoding>,
#[serde(skip_serializing_if = "Option::is_none")]
chunk_size: Option<u32>,
#[serde(skip_serializing_if = "Option::is_none")]
duplicate_policy: Option<TsDuplicatePolicy>,
#[serde(skip_serializing_if = "SmallVec::is_empty")]
labels: SmallVec<[(&'a str, &'a str); 10]>,
}
impl<'a> TsCreateOptions<'a> {
#[must_use]
pub fn retention(mut self, retention_period: u64) -> Self {
self.retention = Some(retention_period);
self
}
#[must_use]
pub fn encoding(mut self, encoding: TsEncoding) -> Self {
self.encoding = Some(encoding);
self
}
#[must_use]
pub fn chunk_size(mut self, chunk_size: u32) -> Self {
self.chunk_size = Some(chunk_size);
self
}
#[must_use]
pub fn duplicate_policy(mut self, policy: TsDuplicatePolicy) -> Self {
self.duplicate_policy = Some(policy);
self
}
#[must_use]
pub fn labels(mut self, labels: impl IntoIterator<Item = (&'a str, &'a str)>) -> Self {
self.labels.extend(labels);
self
}
}
#[derive(Debug, Serialize, Deserialize)]
#[serde(rename_all = "UPPERCASE")]
#[non_exhaustive]
pub enum TsAggregationType {
Avg,
Sum,
Min,
Max,
Range,
Count,
CountNan,
CountAll,
First,
Last,
#[serde(rename = "STD.P")]
StdP,
#[serde(rename = "STD.S")]
StdS,
#[serde(rename = "VAR.P")]
VarP,
#[serde(rename = "VAR.S")]
VarS,
Twa,
}
#[derive(Default, Serialize)]
#[serde(rename_all = "UPPERCASE")]
pub struct TsCreateRuleOptions {
#[serde(rename = "", skip_serializing_if = "Option::is_none")]
align_timestamp: Option<u64>,
}
impl TsCreateRuleOptions {
#[must_use]
pub fn align_timestamp(mut self, align_timestamp: u64) -> Self {
self.align_timestamp = Some(align_timestamp);
self
}
}
#[derive(Default, Serialize)]
#[serde(rename_all = "UPPERCASE")]
pub struct TsIncrByDecrByOptions<'a> {
#[serde(skip_serializing_if = "Option::is_none")]
timestamp: Option<TsTimestamp>,
#[serde(skip_serializing_if = "Option::is_none")]
retention: Option<u64>,
#[serde(
skip_serializing_if = "std::ops::Not::not",
serialize_with = "serialize_flag"
)]
uncompressed: bool,
#[serde(skip_serializing_if = "Option::is_none")]
chunk_size: Option<u32>,
#[serde(skip_serializing_if = "SmallVec::is_empty")]
labels: SmallVec<[(&'a str, &'a str); 10]>,
}
impl<'a> TsIncrByDecrByOptions<'a> {
#[must_use]
pub fn timestamp(mut self, timestamp: TsTimestamp) -> Self {
self.timestamp = Some(timestamp);
self
}
#[must_use]
pub fn retention(mut self, retention_period: u64) -> Self {
self.retention = Some(retention_period);
self
}
#[must_use]
pub fn uncompressed(mut self) -> Self {
self.uncompressed = true;
self
}
#[must_use]
pub fn chunk_size(mut self, chunk_size: u32) -> Self {
self.chunk_size = Some(chunk_size);
self
}
#[must_use]
pub fn labels(mut self, label: &'a str, value: &'a str) -> Self {
self.labels.push((label, value));
self
}
}
#[derive(Default, Serialize)]
#[serde(rename_all = "UPPERCASE")]
pub struct TsGetOptions {
#[serde(
skip_serializing_if = "std::ops::Not::not",
serialize_with = "serialize_flag"
)]
latest: bool,
}
impl TsGetOptions {
#[must_use]
pub fn latest(mut self) -> Self {
self.latest = true;
self
}
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase")]
#[non_exhaustive]
pub struct TsInfoResult {
pub key_self_name: Option<String>,
pub total_samples: usize,
pub memory_usage: usize,
pub first_timestamp: u64,
pub last_timestamp: u64,
pub retention_time: u64,
pub chunk_count: usize,
pub chunk_size: usize,
pub chunk_type: String,
pub duplicate_policy: Option<TsDuplicatePolicy>,
pub labels: HashMap<String, String>,
pub source_key: String,
#[serde(deserialize_with = "deserialize_compation_rules")]
pub rules: Vec<TsCompactionRule>,
#[serde(rename = "Chunks")]
pub chunks: Option<Vec<TsInfoChunkResult>>,
#[serde(flatten)]
pub additional_values: HashMap<String, Value>,
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase")]
#[non_exhaustive]
pub struct TsInfoChunkResult {
pub start_timestamp: i64,
pub end_timestamp: i64,
pub samples: usize,
pub size: usize,
pub bytes_per_sample: f64,
}
#[derive(Debug)]
#[non_exhaustive]
pub struct TsCompactionRule {
pub compaction_key: String,
pub bucket_duration: u64,
pub aggregator: TsAggregationType,
pub alignment: u64,
}
fn deserialize_compation_rules<'de, D>(deserializer: D) -> Result<Vec<TsCompactionRule>, D::Error>
where
D: de::Deserializer<'de>,
{
struct Visitor;
impl<'de> de::Visitor<'de> for Visitor {
type Value = Vec<TsCompactionRule>;
fn expecting(&self, formatter: &mut fmt::Formatter) -> fmt::Result {
formatter.write_str("an array of TsCompactionRule")
}
fn visit_map<A>(self, mut map: A) -> Result<Self::Value, A::Error>
where
A: de::MapAccess<'de>,
{
let mut rules = Vec::with_capacity(map.size_hint().unwrap_or_default());
while let Some(compaction_key) = map.next_key()? {
let (bucket_duration, aggregator, alignment) =
map.next_value::<(u64, TsAggregationType, u64)>()?;
rules.push(TsCompactionRule {
compaction_key,
bucket_duration,
aggregator,
alignment,
});
}
Ok(rules)
}
}
deserializer.deserialize_map(Visitor)
}
#[derive(Default, Serialize)]
#[serde(rename_all = "UPPERCASE")]
pub struct TsMGetOptions<'a> {
#[serde(
skip_serializing_if = "std::ops::Not::not",
serialize_with = "serialize_flag"
)]
latest: bool,
#[serde(
skip_serializing_if = "std::ops::Not::not",
serialize_with = "serialize_flag"
)]
withlabels: bool,
#[serde(skip_serializing_if = "SmallVec::is_empty")]
selected_labels: SmallVec<[&'a str; 10]>,
}
impl<'a> TsMGetOptions<'a> {
#[must_use]
pub fn latest(mut self) -> Self {
self.latest = true;
self
}
#[must_use]
pub fn withlabels(mut self) -> Self {
self.withlabels = true;
self
}
#[must_use]
pub fn selected_label(mut self, label: &'a str) -> Self {
self.selected_labels.push(label);
self
}
}
#[derive(Debug, Clone, Copy, PartialEq)]
pub struct TsGetResult(pub Option<(u64, f64)>);
impl std::ops::Deref for TsGetResult {
type Target = Option<(u64, f64)>;
fn deref(&self) -> &Self::Target {
&self.0
}
}
impl From<TsGetResult> for Option<(u64, f64)> {
fn from(result: TsGetResult) -> Self {
result.0
}
}
impl<'de> de::Deserialize<'de> for TsGetResult {
fn deserialize<D>(deserializer: D) -> std::result::Result<Self, D::Error>
where
D: de::Deserializer<'de>,
{
struct TsGetResultVisitor;
impl<'de> de::Visitor<'de> for TsGetResultVisitor {
type Value = TsGetResult;
fn expecting(&self, formatter: &mut fmt::Formatter) -> fmt::Result {
formatter.write_str("TsGetResult")
}
fn visit_seq<A>(self, mut seq: A) -> std::result::Result<Self::Value, A::Error>
where
A: de::SeqAccess<'de>,
{
let Some(timestamp) = seq.next_element::<u64>()? else {
return Ok(TsGetResult(None));
};
let Some(value) = seq.next_element::<f64>()? else {
return Err(de::Error::invalid_length(1, &"a timestamp and a value"));
};
Ok(TsGetResult(Some((timestamp, value))))
}
fn visit_none<E>(self) -> std::result::Result<Self::Value, E>
where
E: de::Error,
{
Ok(TsGetResult(None))
}
fn visit_unit<E>(self) -> std::result::Result<Self::Value, E>
where
E: de::Error,
{
Ok(TsGetResult(None))
}
}
deserializer.deserialize_seq(TsGetResultVisitor)
}
}
#[derive(Debug, Deserialize)]
#[non_exhaustive]
pub struct TsSample {
pub labels: HashMap<String, String>,
pub timestamp_value: (u64, f64),
}
#[derive(Default, Serialize)]
#[serde(rename_all = "UPPERCASE")]
pub struct TsMRangeOptions<'a> {
#[serde(
skip_serializing_if = "std::ops::Not::not",
serialize_with = "serialize_flag"
)]
latest: bool,
#[serde(skip_serializing_if = "SmallVec::is_empty")]
filter_by_ts: SmallVec<[u64; 10]>,
#[serde(skip_serializing_if = "Option::is_none")]
filter_by_value: Option<(f64, f64)>,
#[serde(
skip_serializing_if = "std::ops::Not::not",
serialize_with = "serialize_flag"
)]
withlabels: bool,
#[serde(skip_serializing_if = "SmallVec::is_empty")]
selected_labels: SmallVec<[&'a str; 10]>,
#[serde(skip_serializing_if = "Option::is_none")]
count: Option<u32>,
#[serde(skip_serializing_if = "Option::is_none")]
align: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
aggregation: Option<(TsAggregationType, u64)>,
#[serde(skip_serializing_if = "Option::is_none")]
buckettimestamp: Option<TsBucketTimestamp>,
#[serde(
skip_serializing_if = "std::ops::Not::not",
serialize_with = "serialize_flag"
)]
empty: bool,
}
impl<'a> TsMRangeOptions<'a> {
#[must_use]
pub fn latest(mut self) -> Self {
self.latest = true;
self
}
#[must_use]
pub fn filter_by_ts(mut self, ts: impl IntoIterator<Item = u64>) -> Self {
self.filter_by_ts.extend(ts);
self
}
#[must_use]
pub fn filter_by_value(mut self, min: f64, max: f64) -> Self {
self.filter_by_value = Some((min, max));
self
}
#[must_use]
pub fn withlabels(mut self) -> Self {
self.withlabels = true;
self
}
#[must_use]
pub fn selected_label(mut self, label: &'a str) -> Self {
self.selected_labels.push(label);
self
}
#[must_use]
pub fn count(mut self, count: u32) -> Self {
self.count = Some(count);
self
}
#[must_use]
pub fn align(mut self, align: &'a str) -> Self {
self.align = Some(align);
self
}
#[must_use]
pub fn aggregation(mut self, aggregator: TsAggregationType, bucket_duration: u64) -> Self {
self.aggregation = Some((aggregator, bucket_duration));
self
}
#[must_use]
pub fn bucket_timestamp(mut self, bucket_timestamp: TsBucketTimestamp) -> Self {
self.buckettimestamp = Some(bucket_timestamp);
self
}
#[must_use]
pub fn empty(mut self) -> Self {
self.empty = true;
self
}
}
#[derive(Debug)]
#[non_exhaustive]
pub struct TsRangeSample {
pub labels: Vec<(String, String)>,
pub reducers: Vec<String>,
pub sources: Vec<String>,
pub aggregators: Vec<String>,
pub values: Vec<(u64, f64)>,
}
impl<'de> de::Deserialize<'de> for TsRangeSample {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: de::Deserializer<'de>,
{
enum TsRangeSampleField {
Aggregators(Vec<String>),
Reducers(Vec<String>),
Sources(Vec<String>),
Values(Vec<(u64, f64)>),
}
impl<'de> de::Deserialize<'de> for TsRangeSampleField {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: de::Deserializer<'de>,
{
struct Visitor;
impl<'de> de::Visitor<'de> for Visitor {
type Value = TsRangeSampleField;
fn expecting(&self, formatter: &mut fmt::Formatter) -> fmt::Result {
formatter.write_str("TsRangeSampleField")
}
fn visit_seq<A>(self, seq: A) -> Result<Self::Value, A::Error>
where
A: de::SeqAccess<'de>,
{
Ok(TsRangeSampleField::Values(Vec::<(u64, f64)>::deserialize(
SeqAccessDeserializer::new(seq),
)?))
}
fn visit_map<A>(self, mut map: A) -> Result<Self::Value, A::Error>
where
A: de::MapAccess<'de>,
{
let (Some((field, value)), None) = (
map.next_entry::<&str, Vec<String>>()?,
map.next_entry::<&str, Vec<String>>()?,
) else {
return Err(de::Error::invalid_length(0, &"1 in map"));
};
match field {
"reducers" => Ok(TsRangeSampleField::Reducers(value)),
"sources" => Ok(TsRangeSampleField::Sources(value)),
"aggregators" => Ok(TsRangeSampleField::Aggregators(value)),
_ => Err(de::Error::unknown_field(
field,
&["reducers", "sources", "aggregators"],
)),
}
}
}
deserializer.deserialize_any(Visitor)
}
}
struct Visitor;
impl<'de> de::Visitor<'de> for Visitor {
type Value = TsRangeSample;
fn expecting(&self, formatter: &mut fmt::Formatter) -> fmt::Result {
formatter.write_str("TsRangeSample")
}
fn visit_seq<A>(self, mut seq: A) -> Result<Self::Value, A::Error>
where
A: de::SeqAccess<'de>,
{
let mut sample = TsRangeSample {
labels: Vec::new(),
reducers: Vec::new(),
sources: Vec::new(),
aggregators: Vec::new(),
values: Vec::new(),
};
let Some(labels) = seq.next_element::<Vec<(String, String)>>()? else {
return Err(de::Error::invalid_length(0, &"more elements in sequence"));
};
sample.labels = labels;
while let Some(field) = seq.next_element::<TsRangeSampleField>()? {
match field {
TsRangeSampleField::Aggregators(aggregators) => {
sample.aggregators = aggregators
}
TsRangeSampleField::Reducers(reducers) => sample.reducers = reducers,
TsRangeSampleField::Sources(sources) => sample.sources = sources,
TsRangeSampleField::Values(values) => sample.values = values,
}
}
Ok(sample)
}
}
deserializer.deserialize_seq(Visitor)
}
}
#[derive(Default, Serialize)]
#[serde(rename_all = "UPPERCASE")]
pub struct TsGroupByOptions<'a> {
#[serde(skip_serializing_if = "Option::is_none")]
groupby: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
reduce: Option<TsAggregationType>,
}
impl<'a> TsGroupByOptions<'a> {
#[must_use]
pub fn new(label: &'a str, reducer: TsAggregationType) -> Self {
Self {
groupby: Some(label),
reduce: Some(reducer),
}
}
}
#[derive(Serialize)]
#[non_exhaustive]
pub enum TsBucketTimestamp {
#[serde(rename = "-")]
Low,
#[serde(rename = "+")]
High,
#[serde(rename = "~")]
Mid,
}
#[derive(Default, Serialize)]
#[serde(rename_all = "UPPERCASE")]
pub struct TsRangeOptions<'a> {
#[serde(
skip_serializing_if = "std::ops::Not::not",
serialize_with = "serialize_flag"
)]
latest: bool,
#[serde(skip_serializing_if = "SmallVec::is_empty")]
filter_by_ts: SmallVec<[u64; 10]>,
#[serde(skip_serializing_if = "Option::is_none")]
filter_by_value: Option<(f64, f64)>,
#[serde(skip_serializing_if = "Option::is_none")]
count: Option<u32>,
#[serde(skip_serializing_if = "Option::is_none")]
align: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
aggregation: Option<(TsAggregationType, u64)>,
#[serde(skip_serializing_if = "Option::is_none")]
buckettimestamp: Option<TsBucketTimestamp>,
#[serde(
skip_serializing_if = "std::ops::Not::not",
serialize_with = "serialize_flag"
)]
empty: bool,
}
impl<'a> TsRangeOptions<'a> {
#[must_use]
pub fn latest(mut self) -> Self {
self.latest = true;
self
}
#[must_use]
pub fn filter_by_ts(mut self, ts: impl IntoIterator<Item = u64>) -> Self {
self.filter_by_ts.extend(ts);
self
}
#[must_use]
pub fn filter_by_value(mut self, min: f64, max: f64) -> Self {
self.filter_by_value = Some((min, max));
self
}
#[must_use]
pub fn count(mut self, count: u32) -> Self {
self.count = Some(count);
self
}
#[must_use]
pub fn align(mut self, align: &'a str) -> Self {
self.align = Some(align);
self
}
#[must_use]
pub fn aggregation(mut self, aggregator: TsAggregationType, bucket_duration: u64) -> Self {
self.aggregation = Some((aggregator, bucket_duration));
self
}
#[must_use]
pub fn bucket_timestamp(mut self, bucket_timestamp: TsBucketTimestamp) -> Self {
self.buckettimestamp = Some(bucket_timestamp);
self
}
#[must_use]
pub fn empty(mut self) -> Self {
self.empty = true;
self
}
}
#[non_exhaustive]
pub enum TsTimestamp {
Value(u64),
ServerClock,
}
impl Serialize for TsTimestamp {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
match self {
TsTimestamp::Value(ts) => serializer.serialize_u64(*ts),
TsTimestamp::ServerClock => serializer.serialize_str("*"),
}
}
}