1use crate::types::Atom;
8use serde::{Deserialize, Serialize};
9use std::fmt;
10
11#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
17pub enum QueryAst {
18 Select(SelectQuery),
20 Find(FindQuery),
22 Aggregate(AggregateQuery),
24 TemporalJoin(TemporalJoinQuery),
26 Stream(StreamQuery),
28 Explain(Box<QueryAst>),
30 Macro(MacroDefinition),
32 Lineage(LineageQuery),
34}
35
36#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
42pub struct SelectQuery {
43 pub projections: Vec<Projection>,
45 pub from: Option<KeyPattern>,
47 pub where_clause: Option<WhereClause>,
49 pub order_by: Option<OrderBy>,
51 pub limit: Option<u64>,
53 pub offset: Option<u64>,
55}
56
57#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
63pub struct FindQuery {
64 pub filter: FilterDocument,
66 pub projection: Option<ProjectionDocument>,
68 pub sort: Option<SortDocument>,
70 pub limit: Option<u64>,
72 pub skip: Option<u64>,
74}
75
76#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
78pub struct FilterDocument {
79 pub raw: serde_json::Value,
81}
82
83#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
85pub struct ProjectionDocument {
86 pub raw: serde_json::Value,
88}
89
90#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
92pub struct SortDocument {
93 pub raw: serde_json::Value,
95}
96
97#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
103pub struct AggregateQuery {
104 pub aggregations: Vec<AggregateFunction>,
106 pub from: Option<KeyPattern>,
108 pub where_clause: Option<WhereClause>,
110 pub group_by: Option<GroupBy>,
112 pub having: Option<WhereClause>,
114 pub order_by: Option<OrderBy>,
116 pub limit: Option<u64>,
118}
119
120#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
122pub enum AggregateFunction {
123 Count,
125 Sum,
127 Avg,
129 Min,
131 Max,
133 First,
135 Last,
137}
138
139#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
141pub enum GroupBy {
142 Key,
144 TimeBucket(TimeBucket),
146}
147
148#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
150pub enum TimeBucket {
151 Minute,
153 Hour,
155 Day,
157 Week,
159 Month,
161}
162
163#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
171pub struct TemporalJoinQuery {
172 pub left: KeyPattern,
174 pub right: KeyPattern,
176 pub join_type: String,
178 pub within_micros: Option<u64>,
180}
181
182#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
186pub struct StreamQuery {
187 pub name: String,
189 pub body: Box<QueryAst>,
191}
192
193#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
197pub struct MacroDefinition {
198 pub name: String,
200 pub params: Vec<String>,
202 pub body: String,
204}
205
206#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
210pub struct LineageQuery {
211 pub key: String,
213}
214
215#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
221pub enum KeyPattern {
222 Exact(String),
224 Prefix(String),
226 Glob(String),
228 Regex(String),
230 Union(Vec<KeyPattern>),
232}
233
234#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
236pub enum ComparisonOp {
237 Eq,
239 Ne,
241 Gt,
243 Gte,
245 Lt,
247 Lte,
249 In,
251 Nin,
253 Like,
255 Regex,
257}
258
259#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
261pub enum BooleanOp {
262 And(Vec<Condition>),
264 Or(Vec<Condition>),
266 Not(Box<Condition>),
268}
269
270#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
272pub enum OrderField {
273 Key,
275 Value,
277 Timestamp,
279}
280
281#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
283pub enum Direction {
284 Asc,
286 Desc,
288}
289
290#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
292pub struct OrderBy {
293 pub field: OrderField,
295 pub direction: Direction,
297}
298
299#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Default)]
301pub struct TimeRange {
302 pub start: Option<u64>,
304 pub end: Option<u64>,
306}
307
308#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
310pub enum ValueFilter {
311 Single(Atom),
313 List(Vec<Atom>),
315}
316
317#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
319pub struct WhereClause {
320 pub root: Condition,
322}
323
324#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
326pub enum Condition {
327 Comparison {
329 field: OrderField,
331 op: ComparisonOp,
333 rhs: ValueFilter,
335 },
336 Key(KeyPattern),
338 TimeRange(TimeRange),
340 Anomaly(AnomalyMethod),
342 Pattern(TimeSeriesPattern),
344 Freshness(FreshnessCondition),
346 Similarity(SimilarityCondition),
348 Boolean(BooleanOp),
350}
351
352#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq)]
356pub enum AnomalyMethod {
357 ZScore(f64),
359 Iqr(f64),
361 MovingAverage {
363 window: u64,
365 threshold: f64,
367 },
368}
369
370#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
372pub enum TimeSeriesPattern {
373 Spike {
375 threshold: f64,
377 max_duration_micros: Option<u64>,
379 },
380 Dip {
382 threshold: f64,
384 max_duration_micros: Option<u64>,
386 },
387 Rising {
389 min_duration_micros: u64,
391 },
392 Falling {
394 min_duration_micros: u64,
396 },
397 Plateau {
399 min_duration_micros: u64,
401 tolerance: f64,
403 },
404}
405
406#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
408pub enum FreshnessCondition {
409 Compare {
411 op: ComparisonOp,
413 value: f64,
415 },
416 Stale,
418 Fresh,
420}
421
422#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
424pub struct SimilarityCondition {
425 pub query_vector: Vec<f32>,
427 pub k: usize,
429 pub index_hint: Option<String>,
431}
432
433#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
437pub enum Projection {
438 All,
440 Key,
442 Value,
444 Timestamp,
446 Function(FunctionCall),
448 Predict(PredictExpr),
450 Aliased {
452 inner: Box<Projection>,
454 alias: String,
456 },
457}
458
459#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
461pub struct FunctionCall {
462 pub name: String,
464 pub args: Vec<String>,
466}
467
468#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
470pub struct PredictExpr {
471 pub method: String,
473 pub horizon: u64,
475 pub interval: Option<String>,
477 pub args: Vec<(String, String)>,
479}
480
481impl fmt::Display for QueryAst {
486 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
487 match self {
488 QueryAst::Select(q) => write!(f, "{}", q),
489 QueryAst::Find(q) => write!(f, "FIND {:?}", q.filter.raw),
490 QueryAst::Aggregate(q) => write!(f, "{}", q),
491 QueryAst::TemporalJoin(q) => {
492 write!(f, "{} TEMPORAL JOIN {} {}", q.left, q.right, q.join_type)
493 }
494 QueryAst::Stream(q) => write!(f, "CREATE STREAM {} AS {}", q.name, q.body),
495 QueryAst::Explain(inner) => write!(f, "EXPLAIN {}", inner),
496 QueryAst::Macro(m) => {
497 write!(
498 f,
499 "CREATE MACRO {}({}) AS {}",
500 m.name,
501 m.params.join(", "),
502 m.body
503 )
504 }
505 QueryAst::Lineage(l) => write!(f, "LINEAGE({})", l.key),
506 }
507 }
508}
509
510impl fmt::Display for SelectQuery {
511 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
512 write!(f, "SELECT ")?;
513 if self.projections.is_empty() {
514 write!(f, "*")?;
515 } else {
516 for (i, p) in self.projections.iter().enumerate() {
517 if i > 0 {
518 write!(f, ", ")?;
519 }
520 write!(f, "{}", p)?;
521 }
522 }
523 if let Some(from) = &self.from {
524 write!(f, " FROM {}", from)?;
525 }
526 if let Some(w) = &self.where_clause {
527 write!(f, " WHERE {}", w.root)?;
528 }
529 if let Some(order) = &self.order_by {
530 write!(f, " ORDER BY {:?} {:?}", order.field, order.direction)?;
531 }
532 if let Some(lim) = self.limit {
533 write!(f, " LIMIT {}", lim)?;
534 }
535 if let Some(off) = self.offset {
536 write!(f, " OFFSET {}", off)?;
537 }
538 Ok(())
539 }
540}
541
542impl fmt::Display for AggregateQuery {
543 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
544 write!(f, "SELECT ")?;
545 for (i, a) in self.aggregations.iter().enumerate() {
546 if i > 0 {
547 write!(f, ", ")?;
548 }
549 write!(f, "{:?}", a)?;
550 }
551 if let Some(from) = &self.from {
552 write!(f, " FROM {}", from)?;
553 }
554 if let Some(w) = &self.where_clause {
555 write!(f, " WHERE {}", w.root)?;
556 }
557 if let Some(g) = &self.group_by {
558 write!(f, " GROUP BY {:?}", g)?;
559 }
560 Ok(())
561 }
562}
563
564impl fmt::Display for KeyPattern {
565 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
566 match self {
567 KeyPattern::Exact(s) => write!(f, "\"{}\"", s),
568 KeyPattern::Prefix(s) => write!(f, "\"{}*\"", s),
569 KeyPattern::Glob(s) => write!(f, "\"{}\"", s),
570 KeyPattern::Regex(s) => write!(f, "REGEX(\"{}\")", s),
571 KeyPattern::Union(parts) => {
572 for (i, p) in parts.iter().enumerate() {
573 if i > 0 {
574 write!(f, " OR ")?;
575 }
576 write!(f, "{}", p)?;
577 }
578 Ok(())
579 }
580 }
581 }
582}
583
584impl fmt::Display for Condition {
585 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
586 match self {
587 Condition::Comparison { field, op, rhs } => {
588 write!(f, "{:?} {:?} {:?}", field, op, rhs)
589 }
590 Condition::Key(kp) => write!(f, "key MATCHES {}", kp),
591 Condition::TimeRange(tr) => {
592 write!(f, "timestamp BETWEEN {:?} AND {:?}", tr.start, tr.end)
593 }
594 Condition::Anomaly(m) => write!(f, "ANOMALY(value, {:?})", m),
595 Condition::Pattern(p) => write!(f, "MATCHES_PATTERN(value, {:?})", p),
596 Condition::Freshness(fc) => write!(f, "{:?}", fc),
597 Condition::Similarity(s) => write!(f, "SIMILAR_TO(<vec>, {})", s.k),
598 Condition::Boolean(BooleanOp::And(cs)) => {
599 for (i, c) in cs.iter().enumerate() {
600 if i > 0 {
601 write!(f, " AND ")?;
602 }
603 write!(f, "({})", c)?;
604 }
605 Ok(())
606 }
607 Condition::Boolean(BooleanOp::Or(cs)) => {
608 for (i, c) in cs.iter().enumerate() {
609 if i > 0 {
610 write!(f, " OR ")?;
611 }
612 write!(f, "({})", c)?;
613 }
614 Ok(())
615 }
616 Condition::Boolean(BooleanOp::Not(c)) => write!(f, "NOT ({})", c),
617 }
618 }
619}
620
621impl fmt::Display for Projection {
622 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
623 match self {
624 Projection::All => write!(f, "*"),
625 Projection::Key => write!(f, "key"),
626 Projection::Value => write!(f, "value"),
627 Projection::Timestamp => write!(f, "timestamp"),
628 Projection::Function(fc) => write!(f, "{}({})", fc.name, fc.args.join(", ")),
629 Projection::Predict(p) => write!(f, "PREDICT({}, horizon={})", p.method, p.horizon),
630 Projection::Aliased { inner, alias } => write!(f, "{} AS {}", inner, alias),
631 }
632 }
633}
634
635#[cfg(test)]
640mod tests {
641 use super::*;
642
643 #[test]
644 fn serde_roundtrip_simple_select() {
645 let ast = QueryAst::Select(SelectQuery {
646 projections: vec![Projection::All],
647 from: Some(KeyPattern::Glob("sensor/*".into())),
648 where_clause: None,
649 order_by: Some(OrderBy {
650 field: OrderField::Timestamp,
651 direction: Direction::Asc,
652 }),
653 limit: Some(100),
654 offset: None,
655 });
656
657 let bytes = bincode::serialize(&ast).unwrap();
658 let decoded: QueryAst = bincode::deserialize(&bytes).unwrap();
659 assert_eq!(ast, decoded);
660 }
661
662 #[test]
663 fn display_pretty_prints_select() {
664 let ast = QueryAst::Select(SelectQuery {
665 projections: vec![Projection::Key, Projection::Value],
666 from: Some(KeyPattern::Prefix("sensor/".into())),
667 where_clause: None,
668 order_by: None,
669 limit: Some(10),
670 offset: None,
671 });
672 let s = format!("{}", ast);
673 assert!(s.contains("SELECT"));
674 assert!(s.contains("FROM"));
675 assert!(s.contains("LIMIT 10"));
676 }
677}