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