1use arrow_schema::{DataType, Schema as ArrowSchema};
16use catalog::table_source::DfTableSourceProvider;
17use common_function::utils::escape_like_pattern;
18use datafusion::datasource::DefaultTableSource;
19use datafusion::execution::SessionState;
20use datafusion_common::{DFSchema, ScalarValue, TableReference};
21use datafusion_expr::utils::{conjunction, disjunction};
22use datafusion_expr::{
23 BinaryExpr, Expr, ExprSchemable, LogicalPlan, LogicalPlanBuilder, Operator, col, lit, not,
24};
25use datatypes::schema::Schema;
26use log_query::{AggFunc, BinaryOperator, EqualValue, LogExpr, LogQuery, TimeFilter};
27use snafu::{OptionExt, ResultExt};
28use table::table::adapter::DfTableProviderAdapter;
29
30use crate::log_query::error::{
31 CatalogSnafu, DataFusionPlanningSnafu, Result, TimeIndexNotFoundSnafu, UnexpectedLogExprSnafu,
32 UnimplementedSnafu, UnknownAggregateFunctionSnafu, UnknownScalarFunctionSnafu,
33 UnknownTableSnafu,
34};
35
36pub struct LogQueryPlanner {
37 table_provider: DfTableSourceProvider,
38 session_state: SessionState,
39}
40
41impl LogQueryPlanner {
42 pub fn new(table_provider: DfTableSourceProvider, session_state: SessionState) -> Self {
43 Self {
44 table_provider,
45 session_state,
46 }
47 }
48
49 pub async fn query_to_plan(&mut self, query: LogQuery) -> Result<LogicalPlan> {
50 let table_ref: TableReference = query.table.table_ref().into();
52 let table_source = self
53 .table_provider
54 .resolve_table(table_ref.clone())
55 .await
56 .context(CatalogSnafu)?;
57 let schema = table_source
58 .downcast_ref::<DefaultTableSource>()
59 .context(UnknownTableSnafu)?
60 .table_provider
61 .downcast_ref::<DfTableProviderAdapter>()
62 .context(UnknownTableSnafu)?
63 .table()
64 .schema();
65
66 let mut plan_builder = LogicalPlanBuilder::scan(table_ref, table_source, None)
68 .context(DataFusionPlanningSnafu)?;
69 let df_schema = plan_builder.schema().clone();
70
71 let mut filters = Vec::new();
73
74 filters.push(self.build_time_filter(&query.time_filter, &schema)?);
76
77 if let Some(filters_expr) = self.build_filters(&query.filters, df_schema.as_arrow())? {
78 filters.push(filters_expr);
79 }
80
81 if !filters.is_empty() {
83 let filter_expr = filters.into_iter().reduce(|a, b| a.and(b)).unwrap();
84 plan_builder = plan_builder
85 .filter(filter_expr)
86 .context(DataFusionPlanningSnafu)?;
87 }
88
89 if !query.columns.is_empty() {
91 let projected_columns = query.columns.iter().map(col).collect::<Vec<_>>();
92 plan_builder = plan_builder
93 .project(projected_columns)
94 .context(DataFusionPlanningSnafu)?;
95 }
96
97 for expr in &query.exprs {
99 plan_builder = self.process_log_expr(plan_builder, expr)?;
100 }
101
102 if query.limit.skip.is_some() || query.limit.fetch.is_some() {
104 plan_builder = plan_builder
105 .limit(query.limit.skip.unwrap_or(0), query.limit.fetch)
106 .context(DataFusionPlanningSnafu)?;
107 }
108
109 let plan = plan_builder.build().context(DataFusionPlanningSnafu)?;
111
112 Ok(plan)
113 }
114
115 fn build_time_filter(&self, time_filter: &TimeFilter, schema: &Schema) -> Result<Expr> {
116 let timestamp_col = schema
117 .timestamp_column()
118 .with_context(|| TimeIndexNotFoundSnafu {})?
119 .name
120 .clone();
121
122 let start_time = ScalarValue::Utf8(time_filter.start.clone());
123 let expr = col(timestamp_col.clone())
124 .gt_eq(lit(start_time))
125 .and(match &time_filter.end {
126 Some(end) => col(timestamp_col).lt(lit(ScalarValue::Utf8(Some(end.clone())))),
127 None => col(timestamp_col).lt_eq(lit(ScalarValue::Utf8(Some(
128 "9999-12-31T23:59:59Z".to_string(),
129 )))),
130 });
131
132 Ok(expr)
133 }
134 fn build_filters(
136 &self,
137 filters: &log_query::Filters,
138 schema: &ArrowSchema,
139 ) -> Result<Option<Expr>> {
140 match filters {
141 log_query::Filters::And(filters) => {
142 let exprs = filters
143 .iter()
144 .filter_map(|filter| self.build_filters(filter, schema).transpose())
145 .collect::<std::result::Result<Vec<_>, _>>()?;
146 if exprs.is_empty() {
147 Ok(None)
148 } else {
149 Ok(conjunction(exprs))
150 }
151 }
152 log_query::Filters::Or(filters) => {
153 let exprs = filters
154 .iter()
155 .filter_map(|filter| self.build_filters(filter, schema).transpose())
156 .collect::<std::result::Result<Vec<_>, _>>()?;
157 if exprs.is_empty() {
158 Ok(None)
159 } else {
160 Ok(disjunction(exprs))
161 }
162 }
163 log_query::Filters::Not(filter) => {
164 if let Some(expr) = self.build_filters(filter, schema)? {
165 Ok(Some(not(expr)))
166 } else {
167 Ok(None)
168 }
169 }
170 log_query::Filters::Single(column_filters) => {
171 self.build_column_filter(column_filters, schema)
173 }
174 }
175 }
176
177 fn build_column_filter(
179 &self,
180 column_filter: &log_query::ColumnFilters,
181 schema: &ArrowSchema,
182 ) -> Result<Option<Expr>> {
183 let df_schema = DFSchema::try_from(schema.clone()).context(DataFusionPlanningSnafu)?;
185 let col_expr = self.log_expr_to_df_expr(&column_filter.expr, &df_schema)?;
186
187 let filter_exprs = column_filter
188 .filters
189 .iter()
190 .filter_map(|filter| {
191 self.build_content_filter_with_expr(col_expr.clone(), filter, &df_schema)
192 .transpose()
193 })
194 .collect::<std::result::Result<Vec<_>, _>>()?;
195
196 if filter_exprs.is_empty() {
197 return Ok(Some(col_expr.is_true()));
198 }
199
200 Ok(conjunction(filter_exprs))
202 }
203
204 #[allow(clippy::only_used_in_recursion)]
206 fn build_content_filter_with_expr(
207 &self,
208 col_expr: Expr,
209 filter: &log_query::ContentFilter,
210 schema: &DFSchema,
211 ) -> Result<Option<Expr>> {
212 match filter {
213 log_query::ContentFilter::Exact(value) => Ok(Some(
214 col_expr.like(lit(ScalarValue::Utf8(Some(escape_like_pattern(value))))),
215 )),
216 log_query::ContentFilter::Prefix(value) => Ok(Some(col_expr.like(lit(
217 ScalarValue::Utf8(Some(format!("{}%", escape_like_pattern(value)))),
218 )))),
219 log_query::ContentFilter::Postfix(value) => Ok(Some(col_expr.like(lit(
220 ScalarValue::Utf8(Some(format!("%{}", escape_like_pattern(value)))),
221 )))),
222 log_query::ContentFilter::Contains(value) => Ok(Some(col_expr.like(lit(
223 ScalarValue::Utf8(Some(format!("%{}%", escape_like_pattern(value)))),
224 )))),
225 log_query::ContentFilter::Regex(_pattern) => Err(UnimplementedSnafu {
226 feature: "regex filter",
227 }
228 .build()),
229 log_query::ContentFilter::Exist => Ok(Some(col_expr.is_not_null())),
230 log_query::ContentFilter::Between {
231 start,
232 end,
233 start_inclusive,
234 end_inclusive,
235 } => {
236 let start_literal = self.create_inferred_literal(start, &col_expr, schema);
237 let end_literal = self.create_inferred_literal(end, &col_expr, schema);
238
239 let left = if *start_inclusive {
240 col_expr.clone().gt_eq(start_literal)
241 } else {
242 col_expr.clone().gt(start_literal)
243 };
244 let right = if *end_inclusive {
245 col_expr.lt_eq(end_literal)
246 } else {
247 col_expr.lt(end_literal)
248 };
249 Ok(Some(left.and(right)))
250 }
251 log_query::ContentFilter::GreatThan { value, inclusive } => {
252 let value_literal = self.create_inferred_literal(value, &col_expr, schema);
253 let comparison_expr = if *inclusive {
254 col_expr.gt_eq(value_literal)
255 } else {
256 col_expr.gt(value_literal)
257 };
258 Ok(Some(comparison_expr))
259 }
260 log_query::ContentFilter::LessThan { value, inclusive } => {
261 let value_literal = self.create_inferred_literal(value, &col_expr, schema);
262 if *inclusive {
263 Ok(Some(col_expr.lt_eq(value_literal)))
264 } else {
265 Ok(Some(col_expr.lt(value_literal)))
266 }
267 }
268 log_query::ContentFilter::In(values) => {
269 let inferred_values: Vec<_> = values
270 .iter()
271 .map(|v| self.create_inferred_literal(v, &col_expr, schema))
272 .collect();
273 Ok(Some(col_expr.in_list(inferred_values, false)))
274 }
275 log_query::ContentFilter::IsTrue => Ok(Some(col_expr.is_true())),
276 log_query::ContentFilter::IsFalse => Ok(Some(col_expr.is_false())),
277 log_query::ContentFilter::Equal(value) => {
278 let value_literal = Self::create_eq_literal(value.clone());
279 Ok(Some(col_expr.eq(value_literal)))
280 }
281 log_query::ContentFilter::Compound(filters, op) => {
282 let exprs = filters
283 .iter()
284 .filter_map(|filter| {
285 self.build_content_filter_with_expr(col_expr.clone(), filter, schema)
286 .transpose()
287 })
288 .collect::<std::result::Result<Vec<_>, _>>()?;
289
290 if exprs.is_empty() {
291 return Ok(None);
292 }
293
294 match op {
295 log_query::ConjunctionOperator::And => Ok(conjunction(exprs)),
296 log_query::ConjunctionOperator::Or => {
297 Ok(exprs.into_iter().reduce(|a, b| a.or(b)))
299 }
300 }
301 }
302 }
303 }
304
305 fn build_aggr_func(
306 &self,
307 schema: &DFSchema,
308 expr: &[AggFunc],
309 by: &[LogExpr],
310 ) -> Result<(Vec<Expr>, Vec<Expr>)> {
311 let aggr_expr = expr
312 .iter()
313 .map(|agg_func| {
314 let AggFunc {
315 name: fn_name,
316 args,
317 alias,
318 } = agg_func;
319 let aggr_fn = self
320 .session_state
321 .aggregate_functions()
322 .get(fn_name)
323 .with_context(|| UnknownAggregateFunctionSnafu {
324 name: fn_name.clone(),
325 })?;
326 let args = args
327 .iter()
328 .map(|expr| self.log_expr_to_df_expr(expr, schema))
329 .collect::<std::result::Result<Vec<_>, _>>()?;
330 if let Some(alias) = alias {
331 Ok(aggr_fn.call(args).alias(alias))
332 } else {
333 Ok(aggr_fn.call(args))
334 }
335 })
336 .collect::<std::result::Result<Vec<_>, _>>()?;
337
338 let group_exprs = by
339 .iter()
340 .map(|expr| self.log_expr_to_df_expr(expr, schema))
341 .collect::<std::result::Result<Vec<_>, _>>()?;
342
343 Ok((aggr_expr, group_exprs))
344 }
345
346 fn log_expr_to_df_expr(&self, expr: &LogExpr, schema: &DFSchema) -> Result<Expr> {
348 match expr {
349 LogExpr::NamedIdent(name) => Ok(col(name)),
350 LogExpr::PositionalIdent(index) => Ok(col(schema.field(*index).name())),
351 LogExpr::Literal(literal) => Ok(lit(ScalarValue::Utf8(Some(literal.clone())))),
352 LogExpr::BinaryOp { left, op, right } => {
353 self.build_binary_expr(left, op, right, schema)
355 }
356 LogExpr::ScalarFunc { name, args, alias } => {
357 self.build_scalar_func(schema, name, args, alias)
358 }
359 LogExpr::Alias { expr, alias } => {
360 let df_expr = self.log_expr_to_df_expr(expr, schema)?;
361 Ok(df_expr.alias(alias))
362 }
363 LogExpr::AggrFunc { .. } | LogExpr::Filter { .. } | LogExpr::Decompose { .. } => {
364 UnexpectedLogExprSnafu {
365 expr: expr.clone(),
366 expected: "not a typical expression",
367 }
368 .fail()
369 }
370 }
371 }
372
373 fn build_scalar_func(
374 &self,
375 schema: &DFSchema,
376 name: &str,
377 args: &[LogExpr],
378 alias: &Option<String>,
379 ) -> Result<Expr> {
380 let args = args
381 .iter()
382 .map(|expr| self.log_expr_to_df_expr(expr, schema))
383 .collect::<std::result::Result<Vec<_>, _>>()?;
384 let func = self.session_state.scalar_functions().get(name).context(
385 UnknownScalarFunctionSnafu {
386 name: name.to_string(),
387 },
388 )?;
389 let expr = func.call(args);
390
391 if let Some(alias) = alias {
392 Ok(expr.alias(alias))
393 } else {
394 Ok(expr)
395 }
396 }
397
398 fn binary_operator_to_df_operator(op: &BinaryOperator) -> Operator {
400 match op {
401 BinaryOperator::Eq => Operator::Eq,
402 BinaryOperator::Ne => Operator::NotEq,
403 BinaryOperator::Lt => Operator::Lt,
404 BinaryOperator::Le => Operator::LtEq,
405 BinaryOperator::Gt => Operator::Gt,
406 BinaryOperator::Ge => Operator::GtEq,
407 BinaryOperator::Plus => Operator::Plus,
408 BinaryOperator::Minus => Operator::Minus,
409 BinaryOperator::Multiply => Operator::Multiply,
410 BinaryOperator::Divide => Operator::Divide,
411 BinaryOperator::Modulo => Operator::Modulo,
412 BinaryOperator::And => Operator::And,
413 BinaryOperator::Or => Operator::Or,
414 }
415 }
416
417 fn infer_literal_scalar_value(&self, literal: &str, target_type: &DataType) -> ScalarValue {
420 let utf8_literal = ScalarValue::Utf8(Some(literal.to_string()));
421 utf8_literal.cast_to(target_type).unwrap_or(utf8_literal)
422 }
423
424 fn build_binary_expr(
427 &self,
428 left: &LogExpr,
429 op: &BinaryOperator,
430 right: &LogExpr,
431 schema: &DFSchema,
432 ) -> Result<Expr> {
433 let mut left_expr = self.log_expr_to_df_expr(left, schema)?;
435 let mut right_expr = self.log_expr_to_df_expr(right, schema)?;
436
437 match (left, right) {
439 (LogExpr::Literal(_), LogExpr::Literal(_)) => {
440 }
442 (LogExpr::Literal(literal), _) => {
443 if let Ok(right_type) = right_expr.get_type(schema) {
445 let inferred_scalar = self.infer_literal_scalar_value(literal, &right_type);
446 left_expr = lit(inferred_scalar);
447 }
448 }
449 (_, LogExpr::Literal(literal)) => {
450 if let Ok(left_type) = left_expr.get_type(schema) {
452 let inferred_scalar = self.infer_literal_scalar_value(literal, &left_type);
453 right_expr = lit(inferred_scalar);
454 }
455 }
456 _ => {
457 }
459 }
460
461 let df_op = Self::binary_operator_to_df_operator(op);
462 Ok(Expr::BinaryExpr(BinaryExpr {
463 left: Box::new(left_expr),
464 op: df_op,
465 right: Box::new(right_expr),
466 }))
467 }
468
469 fn create_inferred_literal(&self, value: &str, expr: &Expr, schema: &DFSchema) -> Expr {
472 if let Ok(expr_type) = expr.get_type(schema) {
473 lit(self.infer_literal_scalar_value(value, &expr_type))
474 } else {
475 lit(ScalarValue::Utf8(Some(value.to_string())))
476 }
477 }
478
479 fn create_eq_literal(value: EqualValue) -> Expr {
480 match value {
481 EqualValue::String(s) => lit(ScalarValue::Utf8(Some(s))),
482 EqualValue::Float(n) => lit(ScalarValue::Float64(Some(n))),
483 EqualValue::Int(n) => lit(ScalarValue::Int64(Some(n))),
484 EqualValue::Boolean(b) => lit(ScalarValue::Boolean(Some(b))),
485 EqualValue::UInt(n) => lit(ScalarValue::UInt64(Some(n))),
486 }
487 }
488
489 fn process_log_expr(
493 &self,
494 plan_builder: LogicalPlanBuilder,
495 expr: &LogExpr,
496 ) -> Result<LogicalPlanBuilder> {
497 let mut plan_builder = plan_builder;
498
499 match expr {
500 LogExpr::AggrFunc { expr, by } => {
501 let schema = plan_builder.schema();
502 let (aggr_expr, group_exprs) = self.build_aggr_func(schema, expr, by)?;
503
504 plan_builder = plan_builder
505 .aggregate(group_exprs, aggr_expr)
506 .context(DataFusionPlanningSnafu)?;
507 }
508 LogExpr::Filter { filter } => {
509 let schema = plan_builder.schema();
510 if let Some(filter_expr) = self.build_column_filter(filter, schema.as_arrow())? {
511 plan_builder = plan_builder
512 .filter(filter_expr)
513 .context(DataFusionPlanningSnafu)?;
514 }
515 }
516 LogExpr::ScalarFunc { name, args, alias } => {
517 let schema = plan_builder.schema();
518 let expr = self.build_scalar_func(schema, name, args, alias)?;
519 plan_builder = plan_builder
520 .project([expr])
521 .context(DataFusionPlanningSnafu)?;
522 }
523 LogExpr::NamedIdent(_) | LogExpr::PositionalIdent(_) => {
524 }
526 LogExpr::Alias { expr, alias } => {
527 let schema = plan_builder.schema();
528 let df_expr = self.log_expr_to_df_expr(expr, schema)?;
529 let aliased_expr = df_expr.alias(alias);
530 plan_builder = plan_builder
531 .project([aliased_expr.clone()])
532 .context(DataFusionPlanningSnafu)?;
533 }
534 LogExpr::BinaryOp { .. } => {
535 let schema = plan_builder.schema();
536 let binary_expr = self.log_expr_to_df_expr(expr, schema)?;
537
538 plan_builder = plan_builder
539 .project([binary_expr])
540 .context(DataFusionPlanningSnafu)?;
541 }
542 _ => {
543 UnimplementedSnafu {
544 feature: "log expression",
545 }
546 .fail()?;
547 }
548 }
549 Ok(plan_builder)
550 }
551}
552
553#[cfg(test)]
554mod tests {
555 use std::sync::Arc;
556
557 use catalog::RegisterTableRequest;
558 use catalog::memory::MemoryCatalogManager;
559 use common_catalog::consts::DEFAULT_CATALOG_NAME;
560 use common_query::test_util::DummyDecoder;
561 use datafusion::execution::SessionStateBuilder;
562 use datatypes::prelude::ConcreteDataType;
563 use datatypes::schema::{ColumnSchema, SchemaRef};
564 use log_query::{
565 ColumnFilters, ConjunctionOperator, ContentFilter, Context, Filters, Limit, LogExpr,
566 };
567 use session::context::QueryContext;
568 use table::metadata::{TableInfoBuilder, TableMetaBuilder};
569 use table::table_name::TableName;
570 use table::test_util::EmptyTable;
571
572 use super::*;
573
574 fn mock_schema() -> SchemaRef {
575 let columns = vec![
576 ColumnSchema::new(
577 "message".to_string(),
578 ConcreteDataType::string_datatype(),
579 false,
580 ),
581 ColumnSchema::new(
582 "timestamp".to_string(),
583 ConcreteDataType::timestamp_millisecond_datatype(),
584 false,
585 )
586 .with_time_index(true),
587 ColumnSchema::new(
588 "host".to_string(),
589 ConcreteDataType::string_datatype(),
590 true,
591 ),
592 ColumnSchema::new(
593 "is_active".to_string(),
594 ConcreteDataType::boolean_datatype(),
595 true,
596 ),
597 ];
598
599 Arc::new(Schema::new(columns))
600 }
601
602 fn mock_schema_with_typed_columns() -> SchemaRef {
603 let columns = vec![
604 ColumnSchema::new(
605 "message".to_string(),
606 ConcreteDataType::string_datatype(),
607 false,
608 ),
609 ColumnSchema::new(
610 "timestamp".to_string(),
611 ConcreteDataType::timestamp_millisecond_datatype(),
612 false,
613 )
614 .with_time_index(true),
615 ColumnSchema::new(
616 "host".to_string(),
617 ConcreteDataType::string_datatype(),
618 true,
619 ),
620 ColumnSchema::new(
621 "is_active".to_string(),
622 ConcreteDataType::boolean_datatype(),
623 true,
624 ),
625 ColumnSchema::new("age".to_string(), ConcreteDataType::int32_datatype(), true),
627 ColumnSchema::new(
628 "score".to_string(),
629 ConcreteDataType::float64_datatype(),
630 true,
631 ),
632 ColumnSchema::new(
633 "count".to_string(),
634 ConcreteDataType::uint64_datatype(),
635 true,
636 ),
637 ];
638
639 Arc::new(Schema::new(columns))
640 }
641
642 async fn build_test_table_provider(
644 table_name_tuples: &[(String, String)],
645 ) -> DfTableSourceProvider {
646 build_test_table_provider_with_schema(table_name_tuples, mock_schema()).await
647 }
648
649 async fn build_test_table_provider_with_typed_columns(
651 table_name_tuples: &[(String, String)],
652 ) -> DfTableSourceProvider {
653 build_test_table_provider_with_schema(table_name_tuples, mock_schema_with_typed_columns())
654 .await
655 }
656
657 async fn build_test_table_provider_with_schema(
658 table_name_tuples: &[(String, String)],
659 schema: SchemaRef,
660 ) -> DfTableSourceProvider {
661 let catalog_list = MemoryCatalogManager::with_default_setup();
662 for (schema_name, table_name) in table_name_tuples {
663 let table_meta = TableMetaBuilder::empty()
664 .schema(schema.clone())
665 .primary_key_indices(vec![2])
666 .value_indices(vec![0])
667 .next_column_id(1024)
668 .build()
669 .unwrap();
670 let table_info = TableInfoBuilder::default()
671 .name(table_name.clone())
672 .meta(table_meta)
673 .build()
674 .unwrap();
675 let table = EmptyTable::from_table_info(&table_info);
676
677 catalog_list
678 .register_table_sync(RegisterTableRequest {
679 catalog: DEFAULT_CATALOG_NAME.to_string(),
680 schema: schema_name.clone(),
681 table_name: table_name.clone(),
682 table_id: 1024,
683 table,
684 })
685 .unwrap();
686 }
687
688 DfTableSourceProvider::new(
689 catalog_list,
690 false,
691 QueryContext::arc(),
692 DummyDecoder::arc(),
693 false,
694 )
695 }
696
697 #[tokio::test]
698 async fn test_query_to_plan() {
699 let table_provider =
700 build_test_table_provider(&[("public".to_string(), "test_table".to_string())]).await;
701 let session_state = SessionStateBuilder::new().with_default_features().build();
702 let mut planner = LogQueryPlanner::new(table_provider, session_state);
703
704 let log_query = LogQuery {
705 table: TableName::new(DEFAULT_CATALOG_NAME, "public", "test_table"),
706 time_filter: TimeFilter {
707 start: Some("2021-01-01T00:00:00Z".to_string()),
708 end: Some("2021-01-02T00:00:00Z".to_string()),
709 span: None,
710 },
711 filters: Filters::Single(ColumnFilters {
712 expr: Box::new(LogExpr::NamedIdent("message".to_string())),
713 filters: vec![ContentFilter::Contains("error".to_string())],
714 }),
715 limit: Limit {
716 skip: None,
717 fetch: Some(100),
718 },
719 context: Context::None,
720 columns: vec![],
721 exprs: vec![],
722 };
723
724 let plan = planner.query_to_plan(log_query).await.unwrap();
725 let expected = "Limit: skip=0, fetch=100 [message:Utf8, timestamp:Timestamp(ms), host:Utf8;N, is_active:Boolean;N]\
726\n Filter: greptime.public.test_table.timestamp >= Utf8(\"2021-01-01T00:00:00Z\") AND greptime.public.test_table.timestamp < Utf8(\"2021-01-02T00:00:00Z\") AND greptime.public.test_table.message LIKE Utf8(\"%error%\") [message:Utf8, timestamp:Timestamp(ms), host:Utf8;N, is_active:Boolean;N]\
727\n TableScan: greptime.public.test_table [message:Utf8, timestamp:Timestamp(ms), host:Utf8;N, is_active:Boolean;N]";
728
729 assert_eq!(plan.display_indent_schema().to_string(), expected);
730 }
731
732 #[tokio::test]
733 async fn test_build_time_filter() {
734 let table_provider =
735 build_test_table_provider(&[("public".to_string(), "test_table".to_string())]).await;
736 let session_state = SessionStateBuilder::new().with_default_features().build();
737 let planner = LogQueryPlanner::new(table_provider, session_state);
738
739 let time_filter = TimeFilter {
740 start: Some("2021-01-01T00:00:00Z".to_string()),
741 end: Some("2021-01-02T00:00:00Z".to_string()),
742 span: None,
743 };
744
745 let expr = planner
746 .build_time_filter(&time_filter, &mock_schema())
747 .unwrap();
748
749 let expected_expr = col("timestamp")
750 .gt_eq(lit(ScalarValue::Utf8(Some(
751 "2021-01-01T00:00:00Z".to_string(),
752 ))))
753 .and(col("timestamp").lt(lit(ScalarValue::Utf8(Some(
754 "2021-01-02T00:00:00Z".to_string(),
755 )))));
756
757 assert_eq!(format!("{:?}", expr), format!("{:?}", expected_expr));
758 }
759
760 #[tokio::test]
761 async fn test_build_time_filter_without_end() {
762 let table_provider =
763 build_test_table_provider(&[("public".to_string(), "test_table".to_string())]).await;
764 let session_state = SessionStateBuilder::new().with_default_features().build();
765 let planner = LogQueryPlanner::new(table_provider, session_state);
766
767 let time_filter = TimeFilter {
768 start: Some("2021-01-01T00:00:00Z".to_string()),
769 end: None,
770 span: None,
771 };
772
773 let expr = planner
774 .build_time_filter(&time_filter, &mock_schema())
775 .unwrap();
776
777 let expected_expr = col("timestamp")
778 .gt_eq(lit(ScalarValue::Utf8(Some(
779 "2021-01-01T00:00:00Z".to_string(),
780 ))))
781 .and(col("timestamp").lt_eq(lit(ScalarValue::Utf8(Some(
782 "9999-12-31T23:59:59Z".to_string(),
783 )))));
784
785 assert_eq!(format!("{:?}", expr), format!("{:?}", expected_expr));
786 }
787
788 #[tokio::test]
789 async fn test_build_content_filter() {
790 let table_provider =
791 build_test_table_provider(&[("public".to_string(), "test_table".to_string())]).await;
792 let session_state = SessionStateBuilder::new().with_default_features().build();
793 let planner = LogQueryPlanner::new(table_provider, session_state);
794 let schema = mock_schema();
795
796 let column_filter = ColumnFilters {
797 expr: Box::new(LogExpr::NamedIdent("message".to_string())),
798 filters: vec![
799 ContentFilter::Contains("error".to_string()),
800 ContentFilter::Prefix("WARN".to_string()),
801 ],
802 };
803
804 let expr_option = planner
805 .build_column_filter(&column_filter, schema.arrow_schema())
806 .unwrap();
807 assert!(expr_option.is_some());
808
809 let expr = expr_option.unwrap();
810
811 let expected_expr = col("message")
812 .like(lit(ScalarValue::Utf8(Some("%error%".to_string()))))
813 .and(col("message").like(lit(ScalarValue::Utf8(Some("WARN%".to_string())))));
814
815 assert_eq!(format!("{:?}", expr), format!("{:?}", expected_expr));
816 }
817
818 #[tokio::test]
819 async fn test_query_to_plan_with_only_skip() {
820 let table_provider =
821 build_test_table_provider(&[("public".to_string(), "test_table".to_string())]).await;
822 let session_state = SessionStateBuilder::new().with_default_features().build();
823 let mut planner = LogQueryPlanner::new(table_provider, session_state);
824
825 let log_query = LogQuery {
826 table: TableName::new(DEFAULT_CATALOG_NAME, "public", "test_table"),
827 time_filter: TimeFilter {
828 start: Some("2021-01-01T00:00:00Z".to_string()),
829 end: Some("2021-01-02T00:00:00Z".to_string()),
830 span: None,
831 },
832 filters: Filters::Single(ColumnFilters {
833 expr: Box::new(LogExpr::NamedIdent("message".to_string())),
834 filters: vec![ContentFilter::Contains("error".to_string())],
835 }),
836 limit: Limit {
837 skip: Some(10),
838 fetch: None,
839 },
840 context: Context::None,
841 columns: vec![],
842 exprs: vec![],
843 };
844
845 let plan = planner.query_to_plan(log_query).await.unwrap();
846 let expected = "Limit: skip=10, fetch=None [message:Utf8, timestamp:Timestamp(ms), host:Utf8;N, is_active:Boolean;N]\
847\n Filter: greptime.public.test_table.timestamp >= Utf8(\"2021-01-01T00:00:00Z\") AND greptime.public.test_table.timestamp < Utf8(\"2021-01-02T00:00:00Z\") AND greptime.public.test_table.message LIKE Utf8(\"%error%\") [message:Utf8, timestamp:Timestamp(ms), host:Utf8;N, is_active:Boolean;N]\
848\n TableScan: greptime.public.test_table [message:Utf8, timestamp:Timestamp(ms), host:Utf8;N, is_active:Boolean;N]";
849
850 assert_eq!(plan.display_indent_schema().to_string(), expected);
851 }
852
853 #[tokio::test]
854 async fn test_query_to_plan_without_limit() {
855 let table_provider =
856 build_test_table_provider(&[("public".to_string(), "test_table".to_string())]).await;
857 let session_state = SessionStateBuilder::new().with_default_features().build();
858 let mut planner = LogQueryPlanner::new(table_provider, session_state);
859
860 let log_query = LogQuery {
861 table: TableName::new(DEFAULT_CATALOG_NAME, "public", "test_table"),
862 time_filter: TimeFilter {
863 start: Some("2021-01-01T00:00:00Z".to_string()),
864 end: Some("2021-01-02T00:00:00Z".to_string()),
865 span: None,
866 },
867 filters: Filters::Single(ColumnFilters {
868 expr: Box::new(LogExpr::NamedIdent("message".to_string())),
869 filters: vec![ContentFilter::Contains("error".to_string())],
870 }),
871 limit: Limit {
872 skip: None,
873 fetch: None,
874 },
875 context: Context::None,
876 columns: vec![],
877 exprs: vec![],
878 };
879
880 let plan = planner.query_to_plan(log_query).await.unwrap();
881 let expected = "Filter: greptime.public.test_table.timestamp >= Utf8(\"2021-01-01T00:00:00Z\") AND greptime.public.test_table.timestamp < Utf8(\"2021-01-02T00:00:00Z\") AND greptime.public.test_table.message LIKE Utf8(\"%error%\") [message:Utf8, timestamp:Timestamp(ms), host:Utf8;N, is_active:Boolean;N]\
882\n TableScan: greptime.public.test_table [message:Utf8, timestamp:Timestamp(ms), host:Utf8;N, is_active:Boolean;N]";
883
884 assert_eq!(plan.display_indent_schema().to_string(), expected);
885 }
886
887 #[test]
888 fn test_escape_pattern() {
889 assert_eq!(escape_like_pattern("test"), "test");
890 assert_eq!(escape_like_pattern("te%st"), "te\\%st");
891 assert_eq!(escape_like_pattern("te_st"), "te\\_st");
892 assert_eq!(escape_like_pattern("te\\st"), "te\\\\st");
893 }
894
895 #[tokio::test]
896 async fn test_query_to_plan_with_aggr_func() {
897 let table_provider =
898 build_test_table_provider(&[("public".to_string(), "test_table".to_string())]).await;
899 let session_state = SessionStateBuilder::new().with_default_features().build();
900 let mut planner = LogQueryPlanner::new(table_provider, session_state);
901
902 let log_query = LogQuery {
903 table: TableName::new(DEFAULT_CATALOG_NAME, "public", "test_table"),
904 time_filter: TimeFilter {
905 start: Some("2021-01-01T00:00:00Z".to_string()),
906 end: Some("2021-01-02T00:00:00Z".to_string()),
907 span: None,
908 },
909 filters: Default::default(),
910 limit: Limit {
911 skip: None,
912 fetch: Some(100),
913 },
914 context: Context::None,
915 columns: vec![],
916 exprs: vec![LogExpr::AggrFunc {
917 expr: vec![AggFunc::new(
918 "count".to_string(),
919 vec![LogExpr::NamedIdent("message".to_string())],
920 Some("count_result".to_string()),
921 )],
922 by: vec![LogExpr::NamedIdent("host".to_string())],
923 }],
924 };
925
926 let plan = planner.query_to_plan(log_query).await.unwrap();
927 let expected = "Limit: skip=0, fetch=100 [host:Utf8;N, count_result:Int64]\
928\n Aggregate: groupBy=[[greptime.public.test_table.host]], aggr=[[count(greptime.public.test_table.message) AS count_result]] [host:Utf8;N, count_result:Int64]\
929\n Filter: greptime.public.test_table.timestamp >= Utf8(\"2021-01-01T00:00:00Z\") AND greptime.public.test_table.timestamp < Utf8(\"2021-01-02T00:00:00Z\") [message:Utf8, timestamp:Timestamp(ms), host:Utf8;N, is_active:Boolean;N]\
930\n TableScan: greptime.public.test_table [message:Utf8, timestamp:Timestamp(ms), host:Utf8;N, is_active:Boolean;N]";
931
932 assert_eq!(plan.display_indent_schema().to_string(), expected);
933 }
934
935 #[tokio::test]
936 async fn test_query_to_plan_with_scalar_func() {
937 let table_provider =
938 build_test_table_provider(&[("public".to_string(), "test_table".to_string())]).await;
939 let session_state = SessionStateBuilder::new().with_default_features().build();
940 let mut planner = LogQueryPlanner::new(table_provider, session_state);
941
942 let log_query = LogQuery {
943 table: TableName::new(DEFAULT_CATALOG_NAME, "public", "test_table"),
944 time_filter: TimeFilter {
945 start: Some("2021-01-01T00:00:00Z".to_string()),
946 end: Some("2021-01-02T00:00:00Z".to_string()),
947 span: None,
948 },
949 filters: Default::default(),
950 limit: Limit {
951 skip: None,
952 fetch: Some(100),
953 },
954 context: Context::None,
955 columns: vec![],
956 exprs: vec![LogExpr::ScalarFunc {
957 name: "date_trunc".to_string(),
958 args: vec![
959 LogExpr::Literal("day".to_string()),
960 LogExpr::NamedIdent("timestamp".to_string()),
961 ],
962 alias: Some("time_bucket".to_string()),
963 }],
964 };
965
966 let plan = planner.query_to_plan(log_query).await.unwrap();
967 let expected = "Limit: skip=0, fetch=100 [time_bucket:Timestamp(ms)]\
968 \n Projection: date_trunc(Utf8(\"day\"), greptime.public.test_table.timestamp) AS time_bucket [time_bucket:Timestamp(ms)]\
969 \n Filter: greptime.public.test_table.timestamp >= Utf8(\"2021-01-01T00:00:00Z\") AND greptime.public.test_table.timestamp < Utf8(\"2021-01-02T00:00:00Z\") [message:Utf8, timestamp:Timestamp(ms), host:Utf8;N, is_active:Boolean;N]\
970 \n TableScan: greptime.public.test_table [message:Utf8, timestamp:Timestamp(ms), host:Utf8;N, is_active:Boolean;N]";
971
972 assert_eq!(plan.display_indent_schema().to_string(), expected);
973 }
974
975 #[tokio::test]
976 async fn test_build_content_filter_between() {
977 let table_provider =
978 build_test_table_provider(&[("public".to_string(), "test_table".to_string())]).await;
979 let session_state = SessionStateBuilder::new().with_default_features().build();
980 let planner = LogQueryPlanner::new(table_provider, session_state);
981 let schema = mock_schema();
982
983 let column_filter = ColumnFilters {
984 expr: Box::new(LogExpr::NamedIdent("message".to_string())),
985 filters: vec![ContentFilter::Between {
986 start: "a".to_string(),
987 end: "z".to_string(),
988 start_inclusive: true,
989 end_inclusive: false,
990 }],
991 };
992
993 let expr_option = planner
994 .build_column_filter(&column_filter, schema.arrow_schema())
995 .unwrap();
996 assert!(expr_option.is_some());
997
998 let expr = expr_option.unwrap();
999 let expected_expr = col("message")
1000 .gt_eq(lit(ScalarValue::Utf8(Some("a".to_string()))))
1001 .and(col("message").lt(lit(ScalarValue::Utf8(Some("z".to_string())))));
1002
1003 assert_eq!(format!("{:?}", expr), format!("{:?}", expected_expr));
1004 }
1005
1006 #[tokio::test]
1007 async fn test_query_to_plan_with_date_histogram() {
1008 let table_provider =
1009 build_test_table_provider(&[("public".to_string(), "test_table".to_string())]).await;
1010 let session_state = SessionStateBuilder::new().with_default_features().build();
1011 let mut planner = LogQueryPlanner::new(table_provider, session_state);
1012
1013 let log_query = LogQuery {
1014 table: TableName::new(DEFAULT_CATALOG_NAME, "public", "test_table"),
1015 time_filter: TimeFilter {
1016 start: Some("2021-01-01T00:00:00Z".to_string()),
1017 end: Some("2021-01-02T00:00:00Z".to_string()),
1018 span: None,
1019 },
1020 filters: Default::default(),
1021 limit: Limit {
1022 skip: Some(0),
1023 fetch: None,
1024 },
1025 context: Context::None,
1026 columns: vec![],
1027 exprs: vec![
1028 LogExpr::ScalarFunc {
1029 name: "date_bin".to_string(),
1030 args: vec![
1031 LogExpr::Literal("30 seconds".to_string()),
1032 LogExpr::NamedIdent("timestamp".to_string()),
1033 ],
1034 alias: Some("2__date_histogram__time_bucket".to_string()),
1035 },
1036 LogExpr::AggrFunc {
1037 expr: vec![AggFunc::new(
1038 "count".to_string(),
1039 vec![LogExpr::PositionalIdent(0)],
1040 Some("count_result".to_string()),
1041 )],
1042 by: vec![LogExpr::NamedIdent(
1043 "2__date_histogram__time_bucket".to_string(),
1044 )],
1045 },
1046 ],
1047 };
1048
1049 let plan = planner.query_to_plan(log_query).await.unwrap();
1050 let expected = "Limit: skip=0, fetch=None [2__date_histogram__time_bucket:Timestamp(ns);N, count_result:Int64]\
1051\n Aggregate: groupBy=[[2__date_histogram__time_bucket]], aggr=[[count(2__date_histogram__time_bucket) AS count_result]] [2__date_histogram__time_bucket:Timestamp(ns);N, count_result:Int64]\
1052\n Projection: date_bin(Utf8(\"30 seconds\"), greptime.public.test_table.timestamp) AS 2__date_histogram__time_bucket [2__date_histogram__time_bucket:Timestamp(ns);N]\
1053\n Filter: greptime.public.test_table.timestamp >= Utf8(\"2021-01-01T00:00:00Z\") AND greptime.public.test_table.timestamp < Utf8(\"2021-01-02T00:00:00Z\") [message:Utf8, timestamp:Timestamp(ms), host:Utf8;N, is_active:Boolean;N]\
1054\n TableScan: greptime.public.test_table [message:Utf8, timestamp:Timestamp(ms), host:Utf8;N, is_active:Boolean;N]";
1055
1056 assert_eq!(plan.display_indent_schema().to_string(), expected);
1057 }
1058
1059 #[tokio::test]
1060 async fn test_build_compound_filter() {
1061 let table_provider =
1062 build_test_table_provider(&[("public".to_string(), "test_table".to_string())]).await;
1063 let session_state = SessionStateBuilder::new().with_default_features().build();
1064 let planner = LogQueryPlanner::new(table_provider, session_state);
1065 let schema = mock_schema();
1066
1067 let column_filter = ColumnFilters {
1069 expr: Box::new(LogExpr::NamedIdent("message".to_string())),
1070 filters: vec![
1071 ContentFilter::Contains("error".to_string()),
1072 ContentFilter::Prefix("WARN".to_string()),
1073 ],
1074 };
1075 let expr = planner
1076 .build_column_filter(&column_filter, schema.arrow_schema())
1077 .unwrap()
1078 .unwrap();
1079
1080 let expected_expr = col("message")
1081 .like(lit(ScalarValue::Utf8(Some("%error%".to_string()))))
1082 .and(col("message").like(lit(ScalarValue::Utf8(Some("WARN%".to_string())))));
1083
1084 assert_eq!(format!("{:?}", expr), format!("{:?}", expected_expr));
1085
1086 let column_filter = ColumnFilters {
1088 expr: Box::new(LogExpr::NamedIdent("message".to_string())),
1089 filters: vec![ContentFilter::Compound(
1090 vec![
1091 ContentFilter::Contains("error".to_string()),
1092 ContentFilter::Prefix("WARN".to_string()),
1093 ],
1094 ConjunctionOperator::Or,
1095 )],
1096 };
1097 let expr = planner
1098 .build_column_filter(&column_filter, schema.arrow_schema())
1099 .unwrap()
1100 .unwrap();
1101
1102 let expected_expr = col("message")
1103 .like(lit(ScalarValue::Utf8(Some("%error%".to_string()))))
1104 .or(col("message").like(lit(ScalarValue::Utf8(Some("WARN%".to_string())))));
1105
1106 assert_eq!(format!("{:?}", expr), format!("{:?}", expected_expr));
1107
1108 let column_filter = ColumnFilters {
1110 expr: Box::new(LogExpr::NamedIdent("message".to_string())),
1111 filters: vec![ContentFilter::Compound(
1112 vec![
1113 ContentFilter::Contains("error".to_string()),
1114 ContentFilter::Compound(
1115 vec![
1116 ContentFilter::Prefix("WARN".to_string()),
1117 ContentFilter::Exact("DEBUG".to_string()),
1118 ],
1119 ConjunctionOperator::Or,
1120 ),
1121 ],
1122 ConjunctionOperator::And,
1123 )],
1124 };
1125 let expr = planner
1126 .build_column_filter(&column_filter, schema.arrow_schema())
1127 .unwrap()
1128 .unwrap();
1129
1130 let expected_nested = col("message")
1131 .like(lit(ScalarValue::Utf8(Some("WARN%".to_string()))))
1132 .or(col("message").like(lit(ScalarValue::Utf8(Some("DEBUG".to_string())))));
1133 let expected_expr = col("message")
1134 .like(lit(ScalarValue::Utf8(Some("%error%".to_string()))))
1135 .and(expected_nested);
1136
1137 assert_eq!(format!("{:?}", expr), format!("{:?}", expected_expr));
1138 }
1139
1140 #[tokio::test]
1141 async fn test_build_great_than_filter() {
1142 let table_provider =
1143 build_test_table_provider(&[("public".to_string(), "test_table".to_string())]).await;
1144 let session_state = SessionStateBuilder::new().with_default_features().build();
1145 let planner = LogQueryPlanner::new(table_provider, session_state);
1146 let schema = mock_schema();
1147
1148 let column_filter = ColumnFilters {
1150 expr: Box::new(LogExpr::NamedIdent("message".to_string())),
1151 filters: vec![ContentFilter::GreatThan {
1152 value: "error".to_string(),
1153 inclusive: true,
1154 }],
1155 };
1156
1157 let expr_option = planner
1158 .build_column_filter(&column_filter, schema.arrow_schema())
1159 .unwrap();
1160 assert!(expr_option.is_some());
1161
1162 let expr = expr_option.unwrap();
1163 let expected_expr = col("message").gt_eq(lit(ScalarValue::Utf8(Some("error".to_string()))));
1164
1165 assert_eq!(format!("{:?}", expr), format!("{:?}", expected_expr));
1166
1167 let column_filter = ColumnFilters {
1169 expr: Box::new(LogExpr::NamedIdent("message".to_string())),
1170 filters: vec![ContentFilter::GreatThan {
1171 value: "error".to_string(),
1172 inclusive: false,
1173 }],
1174 };
1175
1176 let expr_option = planner
1177 .build_column_filter(&column_filter, schema.arrow_schema())
1178 .unwrap();
1179 assert!(expr_option.is_some());
1180
1181 let expr = expr_option.unwrap();
1182 let expected_expr = col("message").gt(lit(ScalarValue::Utf8(Some("error".to_string()))));
1183
1184 assert_eq!(format!("{:?}", expr), format!("{:?}", expected_expr));
1185 }
1186
1187 #[tokio::test]
1188 async fn test_build_less_than_filter() {
1189 let table_provider =
1190 build_test_table_provider(&[("public".to_string(), "test_table".to_string())]).await;
1191 let session_state = SessionStateBuilder::new().with_default_features().build();
1192 let planner = LogQueryPlanner::new(table_provider, session_state);
1193 let schema = mock_schema();
1194
1195 let column_filter = ColumnFilters {
1197 expr: Box::new(LogExpr::NamedIdent("message".to_string())),
1198 filters: vec![ContentFilter::LessThan {
1199 value: "error".to_string(),
1200 inclusive: true,
1201 }],
1202 };
1203
1204 let expr_option = planner
1205 .build_column_filter(&column_filter, schema.arrow_schema())
1206 .unwrap();
1207 assert!(expr_option.is_some());
1208
1209 let expr = expr_option.unwrap();
1210 let expected_expr = col("message").lt_eq(lit(ScalarValue::Utf8(Some("error".to_string()))));
1211
1212 assert_eq!(format!("{:?}", expr), format!("{:?}", expected_expr));
1213
1214 let column_filter = ColumnFilters {
1216 expr: Box::new(LogExpr::NamedIdent("message".to_string())),
1217 filters: vec![ContentFilter::LessThan {
1218 value: "error".to_string(),
1219 inclusive: false,
1220 }],
1221 };
1222
1223 let expr_option = planner
1224 .build_column_filter(&column_filter, schema.arrow_schema())
1225 .unwrap();
1226 assert!(expr_option.is_some());
1227
1228 let expr = expr_option.unwrap();
1229 let expected_expr = col("message").lt(lit(ScalarValue::Utf8(Some("error".to_string()))));
1230
1231 assert_eq!(format!("{:?}", expr), format!("{:?}", expected_expr));
1232 }
1233
1234 #[tokio::test]
1235 async fn test_build_in_filter() {
1236 let table_provider =
1237 build_test_table_provider(&[("public".to_string(), "test_table".to_string())]).await;
1238 let session_state = SessionStateBuilder::new().with_default_features().build();
1239 let planner = LogQueryPlanner::new(table_provider, session_state);
1240 let schema = mock_schema();
1241
1242 let column_filter = ColumnFilters {
1244 expr: Box::new(LogExpr::NamedIdent("message".to_string())),
1245 filters: vec![ContentFilter::In(vec![
1246 "error".to_string(),
1247 "warning".to_string(),
1248 "info".to_string(),
1249 ])],
1250 };
1251
1252 let expr_option = planner
1253 .build_column_filter(&column_filter, schema.arrow_schema())
1254 .unwrap();
1255 assert!(expr_option.is_some());
1256
1257 let expr = expr_option.unwrap();
1258 let expected_expr = col("message").in_list(
1259 vec![
1260 lit(ScalarValue::Utf8(Some("error".to_string()))),
1261 lit(ScalarValue::Utf8(Some("warning".to_string()))),
1262 lit(ScalarValue::Utf8(Some("info".to_string()))),
1263 ],
1264 false,
1265 );
1266
1267 assert_eq!(format!("{:?}", expr), format!("{:?}", expected_expr));
1268 }
1269
1270 #[tokio::test]
1271 async fn test_build_is_true_filter() {
1272 let table_provider =
1273 build_test_table_provider(&[("public".to_string(), "test_table".to_string())]).await;
1274 let session_state = SessionStateBuilder::new().with_default_features().build();
1275 let planner = LogQueryPlanner::new(table_provider, session_state);
1276 let schema = mock_schema();
1277
1278 let column_filter = ColumnFilters {
1280 expr: Box::new(LogExpr::NamedIdent("is_active".to_string())),
1281 filters: vec![ContentFilter::IsTrue],
1282 };
1283
1284 let expr_option = planner
1285 .build_column_filter(&column_filter, schema.arrow_schema())
1286 .unwrap();
1287 assert!(expr_option.is_some());
1288
1289 let expr = expr_option.unwrap();
1290 let expected_expr_string =
1291 "IsTrue(Column(Column { relation: None, name: \"is_active\" }))".to_string();
1292
1293 assert_eq!(format!("{:?}", expr), expected_expr_string);
1294 }
1295
1296 #[tokio::test]
1297 async fn test_build_filter_with_scalar_fn() {
1298 let table_provider =
1299 build_test_table_provider(&[("public".to_string(), "test_table".to_string())]).await;
1300 let session_state = SessionStateBuilder::new().with_default_features().build();
1301 let planner = LogQueryPlanner::new(table_provider, session_state);
1302 let schema = mock_schema();
1303
1304 let column_filter = ColumnFilters {
1305 expr: Box::new(LogExpr::BinaryOp {
1306 left: Box::new(LogExpr::ScalarFunc {
1307 name: "character_length".to_string(),
1308 args: vec![LogExpr::NamedIdent("message".to_string())],
1309 alias: None,
1310 }),
1311 op: BinaryOperator::Gt,
1312 right: Box::new(LogExpr::Literal("100".to_string())),
1313 }),
1314 filters: vec![ContentFilter::IsTrue],
1315 };
1316
1317 let expr_option = planner
1318 .build_column_filter(&column_filter, schema.arrow_schema())
1319 .unwrap();
1320 assert!(expr_option.is_some());
1321
1322 let expr = expr_option.unwrap();
1323 let expected_expr_string = "character_length(message) > Int32(100) IS TRUE";
1324
1325 assert_eq!(format!("{}", expr), expected_expr_string);
1326 }
1327
1328 #[tokio::test]
1329 async fn test_type_inference_float_comparison() {
1330 let table_provider = build_test_table_provider_with_typed_columns(&[(
1331 "public".to_string(),
1332 "test_table".to_string(),
1333 )])
1334 .await;
1335 let session_state = SessionStateBuilder::new().with_default_features().build();
1336 let planner = LogQueryPlanner::new(table_provider, session_state);
1337 let schema = mock_schema_with_typed_columns();
1338
1339 let column_filter = ColumnFilters {
1341 expr: Box::new(LogExpr::NamedIdent("score".to_string())),
1342 filters: vec![ContentFilter::Between {
1343 start: "75.5".to_string(),
1344 end: "100.0".to_string(),
1345 start_inclusive: true,
1346 end_inclusive: false,
1347 }],
1348 };
1349
1350 let expr_option = planner
1351 .build_column_filter(&column_filter, schema.arrow_schema())
1352 .unwrap();
1353 assert!(expr_option.is_some());
1354
1355 let expr = expr_option.unwrap();
1356 let expected_expr = col("score")
1358 .gt_eq(lit(ScalarValue::Float64(Some(75.5))))
1359 .and(col("score").lt(lit(ScalarValue::Float64(Some(100.0)))));
1360
1361 assert_eq!(format!("{:?}", expr), format!("{:?}", expected_expr));
1362 }
1363
1364 #[tokio::test]
1365 async fn test_type_inference_boolean_comparison() {
1366 let table_provider = build_test_table_provider_with_typed_columns(&[(
1367 "public".to_string(),
1368 "test_table".to_string(),
1369 )])
1370 .await;
1371 let session_state = SessionStateBuilder::new().with_default_features().build();
1372 let planner = LogQueryPlanner::new(table_provider, session_state);
1373 let schema = mock_schema_with_typed_columns();
1374
1375 let column_filter = ColumnFilters {
1377 expr: Box::new(LogExpr::NamedIdent("is_active".to_string())),
1378 filters: vec![ContentFilter::In(vec![
1379 "true".to_string(),
1380 "1".to_string(),
1381 "false".to_string(),
1382 ])],
1383 };
1384
1385 let expr_option = planner
1386 .build_column_filter(&column_filter, schema.arrow_schema())
1387 .unwrap();
1388 assert!(expr_option.is_some());
1389
1390 let expr = expr_option.unwrap();
1391 let expected_expr = col("is_active").in_list(
1393 vec![
1394 lit(ScalarValue::Boolean(Some(true))),
1395 lit(ScalarValue::Boolean(Some(true))),
1396 lit(ScalarValue::Boolean(Some(false))),
1397 ],
1398 false,
1399 );
1400
1401 assert_eq!(format!("{:?}", expr), format!("{:?}", expected_expr));
1402 }
1403
1404 #[tokio::test]
1405 async fn test_fallback_to_utf8_on_parse_failure() {
1406 let table_provider = build_test_table_provider_with_typed_columns(&[(
1407 "public".to_string(),
1408 "test_table".to_string(),
1409 )])
1410 .await;
1411 let session_state = SessionStateBuilder::new().with_default_features().build();
1412 let planner = LogQueryPlanner::new(table_provider, session_state);
1413 let schema = mock_schema_with_typed_columns();
1414
1415 let column_filter = ColumnFilters {
1417 expr: Box::new(LogExpr::NamedIdent("age".to_string())),
1418 filters: vec![ContentFilter::GreatThan {
1419 value: "not_a_number".to_string(),
1420 inclusive: false,
1421 }],
1422 };
1423
1424 let expr_option = planner
1425 .build_column_filter(&column_filter, schema.arrow_schema())
1426 .unwrap();
1427 assert!(expr_option.is_some());
1428
1429 let expr = expr_option.unwrap();
1430 let expected_expr = col("age").gt(lit(ScalarValue::Utf8(Some("not_a_number".to_string()))));
1432
1433 assert_eq!(format!("{:?}", expr), format!("{:?}", expected_expr));
1434 }
1435
1436 #[tokio::test]
1437 async fn test_string_column_remains_utf8() {
1438 let table_provider = build_test_table_provider_with_typed_columns(&[(
1439 "public".to_string(),
1440 "test_table".to_string(),
1441 )])
1442 .await;
1443 let session_state = SessionStateBuilder::new().with_default_features().build();
1444 let planner = LogQueryPlanner::new(table_provider, session_state);
1445 let schema = mock_schema_with_typed_columns();
1446
1447 let column_filter = ColumnFilters {
1449 expr: Box::new(LogExpr::NamedIdent("message".to_string())),
1450 filters: vec![ContentFilter::GreatThan {
1451 value: "123".to_string(),
1452 inclusive: false,
1453 }],
1454 };
1455
1456 let expr_option = planner
1457 .build_column_filter(&column_filter, schema.arrow_schema())
1458 .unwrap();
1459 assert!(expr_option.is_some());
1460
1461 let expr = expr_option.unwrap();
1462 let expected_expr = col("message").gt(lit(ScalarValue::Utf8(Some("123".to_string()))));
1464
1465 assert_eq!(format!("{:?}", expr), format!("{:?}", expected_expr));
1466 }
1467
1468 #[tokio::test]
1469 async fn test_all_binary_operators() {
1470 let table_provider = build_test_table_provider_with_typed_columns(&[(
1471 "public".to_string(),
1472 "test_table".to_string(),
1473 )])
1474 .await;
1475 let session_state = SessionStateBuilder::new().with_default_features().build();
1476 let planner = LogQueryPlanner::new(table_provider, session_state);
1477 let schema = mock_schema_with_typed_columns();
1478
1479 let df_schema = DFSchema::try_from(schema.arrow_schema().clone()).unwrap();
1480
1481 let test_cases = vec![
1483 (BinaryOperator::Eq, Operator::Eq),
1484 (BinaryOperator::Ne, Operator::NotEq),
1485 (BinaryOperator::Lt, Operator::Lt),
1486 (BinaryOperator::Le, Operator::LtEq),
1487 (BinaryOperator::Gt, Operator::Gt),
1488 (BinaryOperator::Ge, Operator::GtEq),
1489 (BinaryOperator::Plus, Operator::Plus),
1490 (BinaryOperator::Minus, Operator::Minus),
1491 (BinaryOperator::Multiply, Operator::Multiply),
1492 (BinaryOperator::Divide, Operator::Divide),
1493 (BinaryOperator::Modulo, Operator::Modulo),
1494 (BinaryOperator::And, Operator::And),
1495 (BinaryOperator::Or, Operator::Or),
1496 ];
1497
1498 for (binary_op, expected_df_op) in test_cases {
1499 let binary_expr = LogExpr::BinaryOp {
1500 left: Box::new(LogExpr::NamedIdent("age".to_string())),
1501 op: binary_op,
1502 right: Box::new(LogExpr::Literal("25".to_string())),
1503 };
1504
1505 let expr = planner
1506 .log_expr_to_df_expr(&binary_expr, &df_schema)
1507 .unwrap();
1508
1509 let expected_expr = Expr::BinaryExpr(BinaryExpr {
1510 left: Box::new(col("age")),
1511 op: expected_df_op,
1512 right: Box::new(lit(ScalarValue::Int32(Some(25)))),
1513 });
1514
1515 assert_eq!(format!("{:?}", expr), format!("{:?}", expected_expr));
1516 }
1517 }
1518
1519 #[tokio::test]
1520 async fn test_nested_binary_operations() {
1521 let table_provider = build_test_table_provider_with_typed_columns(&[(
1522 "public".to_string(),
1523 "test_table".to_string(),
1524 )])
1525 .await;
1526 let session_state = SessionStateBuilder::new().with_default_features().build();
1527 let planner = LogQueryPlanner::new(table_provider, session_state);
1528 let schema = mock_schema_with_typed_columns();
1529
1530 let df_schema = DFSchema::try_from(schema.arrow_schema().clone()).unwrap();
1531
1532 let nested_binary_expr = LogExpr::BinaryOp {
1534 left: Box::new(LogExpr::BinaryOp {
1535 left: Box::new(LogExpr::NamedIdent("age".to_string())),
1536 op: BinaryOperator::Plus,
1537 right: Box::new(LogExpr::Literal("5".to_string())),
1538 }),
1539 op: BinaryOperator::Gt,
1540 right: Box::new(LogExpr::Literal("30".to_string())),
1541 };
1542
1543 let expr = planner
1544 .log_expr_to_df_expr(&nested_binary_expr, &df_schema)
1545 .unwrap();
1546
1547 let expected_expr_debug = r#"BinaryExpr(BinaryExpr { left: BinaryExpr(BinaryExpr { left: Column(Column { relation: None, name: "age" }), op: Plus, right: Literal(Int32(5), None) }), op: Gt, right: Literal(Int32(30), None) })"#;
1549 assert_eq!(format!("{:?}", expr), expected_expr_debug);
1550 }
1551}