Skip to main content

PromPlanner

Struct PromPlanner 

Source
pub struct PromPlanner {
    table_provider: DfTableSourceProvider,
    ctx: PromPlannerContext,
    promql_annotations: Option<PromqlAnnotationCollector>,
}

Fields§

§table_provider: DfTableSourceProvider§ctx: PromPlannerContext§promql_annotations: Option<PromqlAnnotationCollector>

Optional collector passed to native histogram UDFs.

Implementations§

Source§

impl PromPlanner

Source

fn at_ref_time( &self, at: &Option<AtModifier>, offset: &Option<Offset>, ) -> Result<Option<Millisecond>>

Resolve the @ modifier of a vector or matrix selector into the timestamp its sample window is anchored at, in milliseconds since the Unix epoch.

Prometheus semantics:

  • @ <unix_ts> anchors at the given timestamp,
  • @ start() / @ end() anchor at the evaluation range of the whole statement,
  • offset shifts the anchor backwards: the window ends at anchor - offset.

Returns None when the selector has no @ modifier.

Source

pub(crate) fn anchor_sub( lhs: Millisecond, rhs: Millisecond, ) -> Result<Millisecond>

Subtracts rhs from lhs on the millisecond timeline of an @ anchor.

A negative result is valid: @ and offset accept timestamps before the Unix epoch. A result outside the representable millisecond range is rejected like an unrepresentable anchor (Self::system_time_to_millis) instead of clamping it, so the same class of input always gets the same answer.

Source

pub(crate) fn at_modifier_offset( &self, at: &Option<AtModifier>, offset: &Option<Offset>, ) -> Result<Option<Millisecond>>

The offset a selector with an @ modifier is evaluated with.

Prometheus anchors such a selector by rewriting its offset to eval_time - anchor (setOffsetForAtModifier), so that the selector always selects its samples around anchor regardless of the step being evaluated. eval_time is the start of the evaluation the selector belongs to, which is ctx.start.

Returns None when the selector has no @ modifier.

Source

fn promotes_anchored_range_call(expr: &PromExpr) -> bool

Whether expr is a call that has to be evaluated once for the whole grid, because it folds a range selector anchored by @ — the range argument of the call’s parser signature, which is a [MatrixSelector] here; only a call with such an argument can take this path, so no function-name registry is involved.

This is the shape that needs Prometheus’ StepInvariantExpr the most: a range function such as rate derives its result from the evaluation instant it is called at, so folding the window once per step would let the outer evaluation grid change the result of a window that @ fixed. Evaluated once, at the start of the grid, the rewritten offset of Self::at_modifier_offset places the anchor at that instant, and Self::replay_over_grid reports the result at every step.

Unlike Prometheus, which wraps the whole step-invariant subtree (preprocessExprHelper), only the call itself is promoted here. The operators above it are not: they are still planned at every step over the replayed result, which is safe for the row-wise ones and keeps the promotion root narrow. The promotion root has to stay a call over one range selector, because the replay needs one series per batch (Self::series_divide_plan) and only such a call guarantees that the rows it emits still describe the series it was divided by. The operators left out — an aggregation, a join, or a label rewriting call such as label_join — mix or re-label the rows of different series, so they are unsafe as promotion roots even though evaluating them after the replay is fine. A call whose input is an anchored instant selector (abs(some_metric @ 300)) needs no promotion either: the selector anchors and replays its sample per series on its own, and the call above it is row-wise.

predict_linear is the exception among the range functions: it predicts from the evaluation instant of each step (Self::create_range_eval_ts_expr), so it has to stay outside the promoted subtree and follow the grid. The remaining arguments of the call have to be literals, since the replay of the promoted result has no second vector input to divide. Parentheses around the range argument are transparent (rate((m[5m] @ 300))), so they are looked through and the call is promoted as if they were absent. Nothing else of the subtree is unwrapped, so an outer parenthesis promotes no operator above the call.

Source

pub(crate) async fn promote_anchored_range_call( &mut self, prom_expr: &PromExpr, timestamp_fn: bool, query_engine_state: &QueryEngineState, ) -> Result<Option<LogicalPlan>>

Plans the anchored range call prom_expr (Self::promotes_anchored_range_call) as a step-invariant subtree: the call is evaluated on a single evaluation instant (grid_start, the start of the outer evaluation) and its result is then reported at every step of [grid_start, ctx.end] by Self::replay_over_grid.

This is the planner’s counterpart of Prometheus’ StepInvariantExpr for the one shape it promotes. Evaluating the call once matters for the functions that derive their result from the step being evaluated: rate(m[5m] @ 300) folds its window around the anchor once, and the extrapolation boundaries of rate must be derived from that same window at every step instead of following the outer evaluation timestamp.

Only the call itself is promoted; the operators above it are planned as usual over the replayed result. The result of the promoted call is split into one series per batch before it is replayed (Self::series_divide_plan), because the row-wise projection of the call does not preserve the batch layout of the selector.

Returns None when prom_expr is not such a call, so that the caller plans it as usual. The selector inside the promoted call keeps its own @ anchoring (see Self::at_modifier_offset), and planning it with ctx.end == ctx.start folds its window once for that single instant instead of expanding it over the grid, which the replay of the call result above already does.

Source

fn system_time_to_millis(time: &SystemTime) -> Result<Millisecond>

Convert the timestamp of an @ modifier into milliseconds since the Unix epoch.

Source

pub(crate) fn replay_over_grid( &self, anchored: LogicalPlan, grid_start: Millisecond, grid_end: Millisecond, time_index_column: String, ) -> LogicalPlan

Report the samples of anchored at every step of the evaluation grid [grid_start, grid_end].

A selector with an @ modifier is anchored: the sample window is selected once, around the anchor timestamp, and every evaluation step reports that same window. Prometheus does this by rewriting the selector’s offset to eval_time - anchor and only fetching the samples on the first step (setOffsetForAtModifier plus the refetch shortcut in rangeEval).

The expansion reuses [InstantManipulate] with a lookback that spans the whole grid: every step then picks the same sample (or the same already computed value, when anchored ends with a function call such as rate) and stamps it with the step’s timestamp.

Every input batch of anchored must hold exactly one series, because [InstantManipulate] takes a batch as one timeline. A leaf-level replay (m @ 300) consumes the [SeriesDivide] of its selector directly. A promoted call is guaranteed that layout by Self::series_divide_plan, which is why Self::promote_anchored_range_call splits its result before calling this method.

Source

fn series_divide_plan( &self, input: LogicalPlan, time_index_column: &str, ) -> Result<LogicalPlan>

Sorts input by its series key and time index and splits it into one batch per series.

[InstantManipulate] reads every input batch as one series (it takes the timeline of the batch and reports the row selected at every step), so a batch holding several series would lose all but one of them. A selector establishes that layout with its own [SeriesDivide], but the per-series distribution requirement does not reach a promoted subtree above it: the row-wise projection of a call sits in between, so the batch boundaries of the selector are not preserved — in a distributed plan the promoted result can be delivered as one batch holding every series. Sorting and dividing here restores the layout, exactly like Self::prom_matrix_selector_to_plan does for the input of a range function.

Series keys that are not present in input are dropped, since ctx.tag_columns may have drifted from the actual output schema. A plan without any series key column is returned unchanged: there is nothing to divide by.

Source§

impl PromPlanner

Source

pub(super) async fn create_histogram_plan( &mut self, function_name: &str, args: &PromFunctionArgs, query_engine_state: &QueryEngineState, ) -> Result<LogicalPlan>

Create a classic, native, or mixed histogram helper plan.

Source

fn create_native_histogram_expr( &self, function: HistogramFoldOperation, field_column: &str, ) -> DfExpr

Source

fn create_native_histogram_plan( &mut self, function: HistogramFoldOperation, input_plan: LogicalPlan, ) -> Result<LogicalPlan>

Source

fn create_mixed_histogram_plan( &mut self, function: HistogramFoldOperation, input_plan: LogicalPlan, float_field: String, histogram_field: String, ) -> Result<LogicalPlan>

Source

pub(super) async fn create_vector_plan( &mut self, args: &PromFunctionArgs, ) -> Result<LogicalPlan>

Source

pub(super) async fn create_scalar_plan( &mut self, args: &PromFunctionArgs, query_engine_state: &QueryEngineState, ) -> Result<LogicalPlan>

Create a SCALAR_FUNCTION plan

Source

pub(super) async fn create_absent_plan( &mut self, args: &PromFunctionArgs, query_engine_state: &QueryEngineState, ) -> Result<LogicalPlan>

Source§

impl PromPlanner

Source

pub(super) async fn try_plan_binary_island( &mut self, binary_expr: &PromBinaryExpr, ) -> Result<Option<LogicalPlan>>

Source

fn binary_island_join_contexts_supported(leaves: &[PlannedIslandLeaf]) -> bool

Source

fn join_binary_island_leaf( &self, left: LogicalPlan, first_leaf: &PlannedIslandLeaf, right_leaf: &PlannedIslandLeaf, ) -> Result<LogicalPlan>

Source

fn build_binary_island_field_exprs( expr: &IslandExpr, leaves: &[PlannedIslandLeaf], schema: &DFSchemaRef, ) -> Result<IslandFieldExprs>

Source

fn project_binary_island( &mut self, input: LogicalPlan, base_alias: &TableReference, base_ctx: &PromPlannerContext, field_exprs: IslandFieldExprs, ) -> Result<LogicalPlan>

Source§

impl PromPlanner

Source

pub(super) fn or_operator( &mut self, left: LogicalPlan, right: LogicalPlan, left_tag_cols_set: HashSet<String>, right_tag_cols_set: HashSet<String>, left_context: PromPlannerContext, right_context: PromPlannerContext, modifier: &Option<BinModifier>, ) -> Result<LogicalPlan>

Source

pub(super) fn set_op_on_non_field_columns( &mut self, left: LogicalPlan, right: LogicalPlan, left_context: PromPlannerContext, right_context: PromPlannerContext, op: TokenType, modifier: &Option<BinModifier>, ) -> Result<LogicalPlan>

Build a set operator (AND/OR/UNLESS)

Source§

impl PromPlanner

Source

pub async fn stmt_to_plan( table_provider: DfTableSourceProvider, stmt: &EvalStmt, query_engine_state: &QueryEngineState, ) -> Result<LogicalPlan>

Source

pub async fn stmt_to_plan_with_annotations( table_provider: DfTableSourceProvider, stmt: &EvalStmt, query_engine_state: &QueryEngineState, promql_annotations: Option<PromqlAnnotationCollector>, ) -> Result<LogicalPlan>

Plans a PromQL statement and passes the optional collector to histogram UDFs.

Source

pub async fn prom_expr_to_plan( &mut self, prom_expr: &PromExpr, query_engine_state: &QueryEngineState, ) -> Result<LogicalPlan>

Source

fn prom_expr_to_plan_inner<'life0, 'life1, 'life_self, 'async_recursion>( &'life_self mut self, prom_expr: &'life0 PromExpr, timestamp_fn: bool, query_engine_state: &'life1 QueryEngineState, ) -> Pin<Box<dyn Future<Output = Result<LogicalPlan>> + Send + 'async_recursion>>
where 'life0: 'async_recursion, 'life1: 'async_recursion, 'life_self: 'async_recursion,

Converts a PromQL expression to a logical plan.

NOTE: The timestamp_fn indicates whether the PromQL timestamp() function is being evaluated in the current context. If true, the planner generates a logical plan that projects the timestamp (time index) column as the value column for each input row, implementing the PromQL timestamp() function semantics. If false, the planner generates the standard logical plan for the given PromQL expression.

Source

async fn prom_subquery_expr_to_plan( &mut self, query_engine_state: &QueryEngineState, subquery_expr: &SubqueryExpr, ) -> Result<LogicalPlan>

Source

async fn prom_aggr_expr_to_plan( &mut self, query_engine_state: &QueryEngineState, aggr_expr: &AggregateExpr, ) -> Result<LogicalPlan>

Source

async fn prom_topk_bottomk_to_plan( &mut self, aggr_expr: &AggregateExpr, input: LogicalPlan, ) -> Result<LogicalPlan>

Create logical plan for PromQL topk and bottomk expr.

Source

async fn prom_unary_expr_to_plan( &mut self, query_engine_state: &QueryEngineState, unary_expr: &UnaryExpr, ) -> Result<LogicalPlan>

Source

fn negate_field_columns(&mut self, input: LogicalPlan) -> Result<LogicalPlan>

Source

async fn prom_binary_expr_to_plan( &mut self, query_engine_state: &QueryEngineState, binary_expr: &PromBinaryExpr, ) -> Result<LogicalPlan>

Source

fn filter_binary_projection( &mut self, input: LogicalPlan, has_native_histogram: bool, preserve_any_value: bool, retain_field_columns: Vec<bool>, ) -> Result<LogicalPlan>

Source

fn project_binary_join_side( &mut self, input: LogicalPlan, table_ref: &TableReference, context: &PromPlannerContext, result_labels: Option<&BinaryResultLabels>, ) -> Result<LogicalPlan>

Source

fn prom_number_lit_to_plan( &mut self, number_literal: &NumberLiteral, ) -> Result<LogicalPlan>

Source

fn prom_string_lit_to_plan( &mut self, string_literal: &StringLiteral, ) -> Result<LogicalPlan>

Source

fn offset_millis(offset: &Option<Offset>) -> Millisecond

The offset of a selector in milliseconds. A positive offset selects samples from an earlier time and moves them forward into the evaluation timeline.

Source

fn series_key_columns(&self) -> Vec<String>

The columns that identify one series, which is the series key expected by the PromQL plan nodes that hold exactly one series per input batch.

Source

fn series_key_columns_for_schema(&self, schema: &DFSchemaRef) -> Vec<String>

Keep replay and series division on the same effective keys when a call rewrites labels.

Source

async fn prom_vector_selector_to_plan( &mut self, vector_selector: &VectorSelector, timestamp_fn: bool, ) -> Result<LogicalPlan>

Source

fn timestamp_seconds_expr(column: &str, schema: &DFSchema) -> Result<DfExpr>

Converts the timestamp column column into PromQL seconds, truncated to milliseconds.

Source

fn create_timestamp_func_plan( &mut self, input: LogicalPlan, timestamp_value: DfExpr, ) -> Result<LogicalPlan>

Builds a projection plan for the PromQL timestamp() function, which reports timestamp_value as the value of each row, along with the original tag and time index columns.

Updates the planner context’s field columns to the timestamp column name.

Source

async fn prom_matrix_selector_to_plan( &mut self, matrix_selector: &MatrixSelector, ) -> Result<LogicalPlan>

Source

async fn prom_call_expr_to_plan( &mut self, query_engine_state: &QueryEngineState, call_expr: &Call, ) -> Result<LogicalPlan>

Source

async fn prom_ext_expr_to_plan( &mut self, query_engine_state: &QueryEngineState, ext_expr: &Extension, ) -> Result<LogicalPlan>

Source

fn preprocess_label_matchers( &mut self, label_matchers: &Matchers, name: &Option<String>, ) -> Result<Matchers>

Extract metric name from __name__ matcher and set it into PromPlannerContext. Returns a new [Matchers] that doesn’t contain metric name matcher.

Each call to this function means new selector is started. Thus, the context will be reset at first.

Name rule:

  • if name is some, then the matchers MUST NOT contain __name__ matcher.
  • if name is none, then the matchers MAY contain NONE OR MULTIPLE __name__ matchers.
Source

async fn selector_to_series_normalize_plan( &mut self, offset_duration: Millisecond, label_matchers: Matchers, is_range_selector: bool, ) -> Result<LogicalPlan>

Source

fn agg_modifier_to_col( &mut self, input_schema: &DFSchemaRef, modifier: &Option<LabelModifier>, update_ctx: bool, ) -> Result<Vec<DfExpr>>

Convert [LabelModifier] to [Column] exprs for aggregation. Timestamp column and tag columns will be included.

§Side effect

This method will also change the tag columns in ctx if update_ctx is true.

Source

pub fn matchers_to_expr( label_matchers: Matchers, table_schema: &DFSchemaRef, ) -> Result<Vec<DfExpr>>

Source

fn find_case_sensitive_column( schema: &DFSchemaRef, column: &str, ) -> Option<String>

Source

fn table_from_source(&self, source: &Arc<dyn TableSource>) -> Result<TableRef>

Source

fn table_ref(&self) -> Result<TableReference>

Source

fn build_time_index_filter( &self, offset_duration: i64, schema: &DFSchemaRef, ) -> Result<Option<DfExpr>>

Source

async fn create_table_scan_plan( &mut self, table_ref: TableReference, ) -> Result<LogicalPlan>

Create a table scan plan and a filter plan with given filter.

§Panic

If the filter is empty

Source

fn collect_row_key_tag_columns_from_plan( &self, plan: &LogicalPlan, ) -> Result<BTreeSet<String>>

Source

async fn setup_context(&mut self) -> Result<Option<LogicalPlan>>

Setup PromPlannerContext’s state fields.

Returns a logical plan for an empty metric.

Source

fn setup_context_for_empty_metric(&mut self) -> Result<LogicalPlan>

Setup PromPlannerContext’s state fields for a non existent table without any rows.

Source

fn create_function_args(&self, args: &[Box<PromExpr>]) -> Result<FunctionArgs>

Source

fn create_mixed_range_function_exprs( &mut self, func: &Function, other_input_exprs: VecDeque<DfExpr>, float_field: &str, histogram_field: &str, input_schema: &DFSchemaRef, range_fold_offset: Option<Millisecond>, ) -> Result<Option<Vec<DfExpr>>>

Source

fn create_function_expr( &mut self, func: &Function, other_input_exprs: Vec<DfExpr>, input_schema: &DFSchemaRef, query_engine_state: &QueryEngineState, range_fold_offset: Option<Millisecond>, ) -> Result<(Vec<DfExpr>, Vec<String>)>

Creates function expressions for projection and returns the expressions and new tags.

§Side Effects

This method will update PromPlannerContext’s fields and tags if needed.

Source

fn select_delta_range_math( &self, function: &str, input_schema: &DFSchemaRef, range_length: Millisecond, delta_sum: DfExpr, cumulative: DfExpr, ) -> Result<DfExpr>

Source

fn validate_label_name(label_name: &str) -> Result<()>

Validate label name according to Prometheus specification. Label names must match the regex: [a-zA-Z_][a-zA-Z0-9_]* Additionally, label names starting with double underscores are reserved for internal use.

Source

fn build_regexp_replace_label_expr( &self, other_input_exprs: &mut VecDeque<DfExpr>, input_schema: &DFSchemaRef, ) -> Result<Option<(DfExpr, String)>>

Build expr for label_replace function

Source

fn build_concat_labels_expr( other_input_exprs: &mut VecDeque<DfExpr>, ctx: &PromPlannerContext, input_schema: &DFSchemaRef, query_engine_state: &QueryEngineState, ) -> Result<(DfExpr, String)>

Build expr for label_join function

Source

fn label_value_expr(label: &str, input_schema: &DFSchemaRef) -> Result<DfExpr>

The value of label as a string, where NULL (the series has no such label) reads as the empty string, as in PromQL.

Source

fn empty_label_to_null(value: DfExpr) -> DfExpr

An empty label value means the label is absent in PromQL. Label functions represent it as NULL, like a series that never had the label, so that both compare equal when label sets are matched and neither is reported as a label.

Source

fn create_time_index_column_expr(&self) -> Result<DfExpr>

Source

fn create_range_eval_ts_expr( &self, fold_offset: Millisecond, input_schema: &DFSchemaRef, ) -> Result<DfExpr>

Builds the evaluation instant the window of the last planned range selector is folded for, as a Timestamp(Millisecond) expression.

The timestamp payload of a folded window is shifted onto the evaluation timeline by the offset the window is folded with (fold_offset), while the time index column of a folded row keeps the evaluation timestamp of its step. Adding the offset back yields the evaluation instant on the payload timeline, which is where the regression of predict_linear is centered: neither a plain window (which may end before the step, and is additionally shifted by offset on the payload timeline) nor an @-anchored one (whose end is the anchor, while the payload is shifted by at_offset) ends at the step it is evaluated at.

The sum is computed on the millisecond representation and cast back, so that the result keeps the Timestamp(Millisecond) type the range functions declare for it.

Source

fn name_without_last_arg(expr: &DfExpr) -> String

The name of expr without its last argument.

A predict_linear call appends the evaluation instant of its window as a private last argument (Self::create_range_eval_ts_expr). Naming the output column after the call the user wrote, without that argument, keeps the injected expression out of the user-visible schema; the expression itself keeps every argument it needs.

Source

fn create_tag_column_exprs(&self) -> Result<Vec<DfExpr>>

Source

fn create_field_column_exprs(&self) -> Result<Vec<DfExpr>>

Source

fn create_tag_and_time_index_column_sort_exprs(&self) -> Result<Vec<SortExpr>>

Source

fn create_field_columns_sort_exprs(&self, asc: bool) -> Vec<SortExpr>

Source

fn create_sort_exprs_by_tags( func: &str, tags: Vec<DfExpr>, asc: bool, ) -> Result<Vec<SortExpr>>

Source

fn create_empty_values_filter_expr( &self, preserve_any_value: bool, ) -> Result<DfExpr>

Source

fn create_aggregate_exprs( &mut self, op: TokenType, param: &Option<Box<PromExpr>>, input_plan: &LogicalPlan, ) -> Result<(Vec<DfExpr>, Vec<DfExpr>)>

Creates a set of DataFusion DfExpr::AggregateFunction expressions for each value column using the specified aggregate function.

§Side Effects

This method modifies the value columns in the context by replacing them with the new columns created by the aggregate function application.

§Returns

Returns a tuple of (aggregate_expressions, previous_field_expressions) where:

  • aggregate_expressions: Expressions that apply the aggregate function to the original fields
  • previous_field_expressions: Field expressions naming the pre-aggregation values. This is non-empty only when the operation is count_values, which groups by the sample value and projects it as the generated label, so these expressions are passed through the same formatting as that label (prom_float_to_string).
Source

fn create_numeric_aggregate_expr( op: TokenType, param: &Option<Box<PromExpr>>, input: DfExpr, ) -> Result<DfExpr>

Source

fn create_mixed_aggregate_exprs( &mut self, op: TokenType, param: &Option<Box<PromExpr>>, float_column: &str, histogram_column: &str, ) -> Result<(Vec<DfExpr>, Vec<DfExpr>)>

Source

fn mixed_sample_count_column(column: &str) -> DfExpr

Source

fn mixed_sample_count_name(column: &str) -> String

Source

fn mixed_aggregate_filter_expr( &self, op: TokenType, float_column: &str, histogram_column: &str, ) -> Result<DfExpr>

Source

fn mixed_ignored_histogram_filter_expr( &self, op: TokenType, histogram_column: &str, ) -> Result<DfExpr>

Source

fn create_native_histogram_aggregate_expr( &self, op: TokenType, column: &str, ) -> Result<DfExpr>

Source

fn create_native_histogram_aggregate_exprs( &mut self, op: TokenType, input_plan: &LogicalPlan, ) -> Result<(Vec<DfExpr>, Vec<DfExpr>)>

Source

fn get_param_value_as_str( op: TokenType, param: &Option<Box<PromExpr>>, ) -> Result<&str>

Source

fn get_param_as_literal_expr( param: Option<&PromExpr>, op: Option<TokenType>, expected_type: Option<ArrowDataType>, ) -> Result<DfExpr>

Source

fn create_window_exprs( &mut self, op: TokenType, group_exprs: Vec<DfExpr>, input_plan: &LogicalPlan, ) -> Result<Vec<DfExpr>>

Create [DfExpr::WindowFunction] expr for each value column with given window function.

Source

fn try_build_literal_expr(expr: &PromExpr) -> Option<DfExpr>

Try to build a DataFusion Literal Expression from PromQL Expr, return None if the input is not a literal expression.

Source

fn try_build_special_time_expr_with_context( &self, expr: &PromExpr, ) -> Option<DfExpr>

Source

fn native_histogram_binary_expr( token: TokenType, lhs: DfExpr, lhs_is_histogram: bool, rhs: DfExpr, rhs_is_histogram: bool, filter_context: bool, promql_annotations: Option<PromqlAnnotationCollector>, ) -> Result<Option<DfExpr>>

Source

fn prom_token_to_binary_expr_builder( token: TokenType, ) -> Result<Box<dyn Fn(DfExpr, DfExpr) -> Result<DfExpr>>>

Return a lambda to build binary expression from token. Because some binary operator are function in DataFusion like atan2 or ^.

Source

fn is_token_a_comparison_op(token: TokenType) -> bool

Check if the given op is a comparison operator.

Source

fn is_token_a_set_op(token: TokenType) -> bool

Check if the given op is a set operator (UNION, INTERSECT and EXCEPT in SQL).

Source

fn align_binary_field_columns<'a>( left_schema: &DFSchemaRef, right_schema: &DFSchemaRef, left_field_columns: &'a [String], right_field_columns: &'a [String], op: TokenType, left_is_scalar: bool, right_is_scalar: bool, ) -> (Vec<(String, Vec<(&'a String, &'a String)>)>, Vec<(&'a String, &'a String)>)

Source

fn binary_result_is_histogram( token: TokenType, lhs_is_histogram: bool, rhs_is_histogram: bool, ) -> Option<bool>

Source

fn plan_has_tsid_column(plan: &LogicalPlan) -> bool

Source

fn is_empty_metric(plan: &LogicalPlan) -> bool

Source

fn native_histogram_arrow_type() -> ArrowDataType

Source

fn field_column_type<'a>( schema: &'a DFSchemaRef, field_column: &str, ) -> Option<&'a ArrowDataType>

Source

fn field_column_is_native_histogram( schema: &DFSchemaRef, field_column: &str, ) -> bool

Source

fn field_columns_contain_native_histogram( schema: &DFSchemaRef, field_columns: &[String], ) -> bool

Source

fn field_column_is_float_range(schema: &DFSchemaRef, field_column: &str) -> bool

Source

fn field_columns_are_alternative_samples( schema: &DFSchemaRef, field_columns: &[String], ) -> bool

Source

fn alternative_sample_columns<'a>( schema: &DFSchemaRef, field_columns: &'a [String], ) -> Option<(&'a str, &'a str)>

Source

fn alternative_sample_range_columns<'a>( schema: &DFSchemaRef, field_columns: &'a [String], ) -> Option<(&'a str, &'a str)>

Source

fn field_column_is_native_histogram_range( schema: &DFSchemaRef, field_column: &str, ) -> bool

Source

fn all_field_columns_are_native_histograms(&self, schema: &DFSchemaRef) -> bool

Source

fn all_field_columns_are_native_histogram_ranges( &self, schema: &DFSchemaRef, ) -> bool

Source

fn optional_tsid_projection( schema: &DFSchemaRef, table_ref: Option<&TableReference>, keep_tsid: bool, ) -> Option<DfExpr>

Source

fn binary_join_key_columns( &self, left_schema: &DFSchemaRef, right_schema: &DFSchemaRef, left_context: &PromPlannerContext, right_context: &PromPlannerContext, only_join_time_index: bool, modifier: &Option<BinModifier>, ) -> Result<(BTreeSet<String>, BTreeSet<String>, bool)>

Source

fn binary_result_labels( left_context: &PromPlannerContext, right_context: &PromPlannerContext, modifier: &Option<BinModifier>, ) -> Option<Vec<(bool, String)>>

Result labels of a vector-vector binary operation, following Prometheus resultMetric: on(...) keeps only the matching labels, ignoring(...) drops them, and a group modifier keeps the “many” side’s labels plus the group_x(...) labels taken from the “one” side.

The flag of each entry tells which operand the label is projected from. None means the operation keeps a whole operand tag set, which the default projection already does.

Source

fn binary_result_label_projection( schema: &DFSchemaRef, left_table_ref: &TableReference, right_table_ref: &TableReference, left_context: &PromPlannerContext, right_context: &PromPlannerContext, labels: Vec<(bool, String)>, ) -> Result<BinaryResultLabels>

Resolve Self::binary_result_labels against the join output.

Source

fn binary_result_labels_may_repeat( left_context: &PromPlannerContext, right_context: &PromPlannerContext, modifier: &Option<BinModifier>, ) -> bool

Whether the result of a binary operation can hold two series with the same labels, which only Self::binary_result_labels can introduce: one-to-one matching on a subset of the tags, or a group modifier that overwrites a label of the “many” side.

Source

fn assert_unique_match_group( plan: LogicalPlan, group_exprs: Vec<DfExpr>, group_labels: Vec<String>, time_index_expr: DfExpr, violation: MatchGroupViolation, ) -> Result<LogicalPlan>

Wrap plan in a check that fails the query when a match group holds more than one row at a timestamp. group_exprs are the label columns of the group, resolved against plan.

Source

fn binary_modifier_preserves_tsid_join_key( &self, left_context: &PromPlannerContext, right_context: &PromPlannerContext, modifier: &Option<BinModifier>, ) -> bool

Source

fn join_on_non_field_columns( &self, left: LogicalPlan, right: LogicalPlan, left_table_ref: TableReference, right_table_ref: TableReference, left_time_index_column: Option<String>, right_time_index_column: Option<String>, only_join_time_index: bool, modifier: &Option<BinModifier>, left_context: &PromPlannerContext, right_context: &PromPlannerContext, ) -> Result<LogicalPlan>

Build a inner join on time index column and tag columns to concat two logical plans. When only_join_time_index == true we only join on the time index, because these two plan may not have the same tag columns

Source

fn assert_unique_one_side( plan: LogicalPlan, join_keys: &BTreeSet<String>, tag_columns: &[String], time_index_column: Option<&str>, one_side_is_left: bool, ) -> Result<LogicalPlan>

Guard the side of a vector matching that must hold one series per match group.

Source

fn selected_binary_match_labels( left_context: &PromPlannerContext, right_context: &PromPlannerContext, modifier: &Option<BinModifier>, ) -> BTreeSet<String>

Source

fn only_temporality_match_label_mismatches( left_context: &PromPlannerContext, right_context: &PromPlannerContext, modifier: &Option<BinModifier>, ) -> bool

Source

fn align_temporality_match_column( left: LogicalPlan, right: LogicalPlan, left_context: &mut PromPlannerContext, right_context: &mut PromPlannerContext, ) -> Result<(LogicalPlan, LogicalPlan, bool)>

Source

fn normalized_match_key_expr( label: &str, field: Option<(Option<TableReference>, ArrowDataType)>, value_type: &ArrowDataType, internal_name: &str, ) -> DfExpr

Source

fn is_zero_row_empty_relation(plan: &LogicalPlan) -> bool

Source

fn string_value_data_type(data_type: &ArrowDataType) -> Option<&ArrowDataType>

Source

fn string_scalar_value( data_type: &ArrowDataType, value: Option<String>, ) -> Option<ScalarValue>

Source

fn common_label_data_type( left: Option<&ArrowDataType>, right: Option<&ArrowDataType>, ) -> Option<ArrowDataType>

Source

fn projection_for_each_field_column<F>( &mut self, input: LogicalPlan, name_to_expr: F, ) -> Result<LogicalPlan>
where F: FnMut(&String) -> Result<DfExpr>,

Build a projection that project and perform operation expr for every value columns. Non-value columns (tag and timestamp) will be preserved in the projection.

§Side effect

This function will update the value columns in the context. Those new column names don’t contains qualifier.

Source

fn projection_for_each_field_column_with_labels<F>( &mut self, input: LogicalPlan, result_labels: Option<&BinaryResultLabels>, name_to_expr: F, ) -> Result<LogicalPlan>
where F: FnMut(&String) -> Result<DfExpr>,

Like Self::projection_for_each_field_column, but projects result_labels instead of the context tag columns when a binary operation derived its own result label set.

Source

fn filter_on_field_column<F>( &self, input: LogicalPlan, name_to_expr: F, ) -> Result<LogicalPlan>
where F: FnMut(&String) -> Result<DfExpr>,

Build a filter plan on one value column or a float/histogram alternative pair.

Source

fn date_part_on_time_index(&self, date_part: &str) -> Result<DfExpr>

Generate an expr like date_part("hour", <TIME_INDEX>). Caller should ensure the time index column in context is set

Source

fn strip_tsid_column(&self, plan: LogicalPlan) -> Result<LogicalPlan>

Source

fn apply_alias( &mut self, plan: LogicalPlan, alias_name: String, ) -> Result<LogicalPlan>

Apply an alias to the query result by adding a projection with the alias name

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
§

impl<T> Conv for T

§

fn conv<T>(self) -> T
where Self: Into<T>,

Converts self into T using Into<T>. Read more
§

impl<T> FmtForward for T

§

fn fmt_binary(self) -> FmtBinary<Self>
where Self: Binary,

Causes self to use its Binary implementation when Debug-formatted.
§

fn fmt_display(self) -> FmtDisplay<Self>
where Self: Display,

Causes self to use its Display implementation when Debug-formatted.
§

fn fmt_lower_exp(self) -> FmtLowerExp<Self>
where Self: LowerExp,

Causes self to use its LowerExp implementation when Debug-formatted.
§

fn fmt_lower_hex(self) -> FmtLowerHex<Self>
where Self: LowerHex,

Causes self to use its LowerHex implementation when Debug-formatted.
§

fn fmt_octal(self) -> FmtOctal<Self>
where Self: Octal,

Causes self to use its Octal implementation when Debug-formatted.
§

fn fmt_pointer(self) -> FmtPointer<Self>
where Self: Pointer,

Causes self to use its Pointer implementation when Debug-formatted.
§

fn fmt_upper_exp(self) -> FmtUpperExp<Self>
where Self: UpperExp,

Causes self to use its UpperExp implementation when Debug-formatted.
§

fn fmt_upper_hex(self) -> FmtUpperHex<Self>
where Self: UpperHex,

Causes self to use its UpperHex implementation when Debug-formatted.
§

fn fmt_list(self) -> FmtList<Self>
where &'a Self: for<'a> IntoIterator,

Formats each item in a sequence. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

§

impl<T> FutureExt for T

§

fn with_context(self, otel_cx: Context) -> WithContext<Self>

Attaches the provided Context to this type, returning a WithContext wrapper. Read more
§

fn with_current_context(self) -> WithContext<Self>

Attaches the current Context to this type, returning a WithContext wrapper. Read more
§

impl<T> Instrument for T

§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided [Span], returning an Instrumented wrapper. Read more
§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self>

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
§

impl<T> IntoRequest<T> for T

§

fn into_request(self) -> Request<T>

Wrap the input message T in a tonic::Request
§

impl<L> LayerExt<L> for L

§

fn named_layer<S>(&self, service: S) -> Layered<<L as Layer<S>>::Service, S>
where L: Layer<S>,

Applies the layer to a service and wraps it in [Layered].
§

impl<T> Pipe for T
where T: ?Sized,

§

fn pipe<R>(self, func: impl FnOnce(Self) -> R) -> R
where Self: Sized,

Pipes by value. This is generally the method you want to use. Read more
§

fn pipe_ref<'a, R>(&'a self, func: impl FnOnce(&'a Self) -> R) -> R
where R: 'a,

Borrows self and passes that borrow into the pipe function. Read more
§

fn pipe_ref_mut<'a, R>(&'a mut self, func: impl FnOnce(&'a mut Self) -> R) -> R
where R: 'a,

Mutably borrows self and passes that borrow into the pipe function. Read more
§

fn pipe_borrow<'a, B, R>(&'a self, func: impl FnOnce(&'a B) -> R) -> R
where Self: Borrow<B>, B: 'a + ?Sized, R: 'a,

Borrows self, then passes self.borrow() into the pipe function. Read more
§

fn pipe_borrow_mut<'a, B, R>( &'a mut self, func: impl FnOnce(&'a mut B) -> R, ) -> R
where Self: BorrowMut<B>, B: 'a + ?Sized, R: 'a,

Mutably borrows self, then passes self.borrow_mut() into the pipe function. Read more
§

fn pipe_as_ref<'a, U, R>(&'a self, func: impl FnOnce(&'a U) -> R) -> R
where Self: AsRef<U>, U: 'a + ?Sized, R: 'a,

Borrows self, then passes self.as_ref() into the pipe function.
§

fn pipe_as_mut<'a, U, R>(&'a mut self, func: impl FnOnce(&'a mut U) -> R) -> R
where Self: AsMut<U>, U: 'a + ?Sized, R: 'a,

Mutably borrows self, then passes self.as_mut() into the pipe function.
§

fn pipe_deref<'a, T, R>(&'a self, func: impl FnOnce(&'a T) -> R) -> R
where Self: Deref<Target = T>, T: 'a + ?Sized, R: 'a,

Borrows self, then passes self.deref() into the pipe function.
§

fn pipe_deref_mut<'a, T, R>( &'a mut self, func: impl FnOnce(&'a mut T) -> R, ) -> R
where Self: DerefMut<Target = T> + Deref, T: 'a + ?Sized, R: 'a,

Mutably borrows self, then passes self.deref_mut() into the pipe function.
§

impl<T> Pointable for T

§

const ALIGN: usize

The alignment of pointer.
§

type Init = T

The type for initializers.
§

unsafe fn init(init: <T as Pointable>::Init) -> usize

Initializes a with the given initializer. Read more
§

unsafe fn deref<'a>(ptr: usize) -> &'a T

Dereferences the given pointer. Read more
§

unsafe fn deref_mut<'a>(ptr: usize) -> &'a mut T

Mutably dereferences the given pointer. Read more
§

unsafe fn drop(ptr: usize)

Drops the object pointed to by the given pointer. Read more
§

impl<T> PolicyExt for T
where T: ?Sized,

§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns [Action::Follow] only if self and other return Action::Follow. Read more
§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns [Action::Follow] if either self or other returns Action::Follow. Read more
Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
§

impl<T> ServiceExt for T

§

fn propagate_header(self, header: HeaderName) -> PropagateHeader<Self>
where Self: Sized,

Propagate a header from the request to the response. Read more
§

fn add_extension<T>(self, value: T) -> AddExtension<Self, T>
where Self: Sized,

Add some shareable value to request extensions. Read more
§

fn map_request_body<F>(self, f: F) -> MapRequestBody<Self, F>
where Self: Sized,

Apply a transformation to the request body. Read more
§

fn map_response_body<F>(self, f: F) -> MapResponseBody<Self, F>
where Self: Sized,

Apply a transformation to the response body. Read more
§

fn compression(self) -> Compression<Self>
where Self: Sized,

Compresses response bodies. Read more
§

fn decompression(self) -> Decompression<Self>
where Self: Sized,

Decompress response bodies. Read more
§

fn trace_for_http(self) -> Trace<Self, SharedClassifier<ServerErrorsAsFailures>>
where Self: Sized,

High level tracing that classifies responses using HTTP status codes. Read more
§

fn trace_for_grpc(self) -> Trace<Self, SharedClassifier<GrpcErrorsAsFailures>>
where Self: Sized,

High level tracing that classifies responses using gRPC headers. Read more
§

fn follow_redirects(self) -> FollowRedirect<Self>
where Self: Sized,

Follow redirect resposes using the Standard policy. Read more
§

fn sensitive_headers( self, headers: impl IntoIterator<Item = HeaderName>, ) -> SetSensitiveRequestHeaders<SetSensitiveResponseHeaders<Self>>
where Self: Sized,

Mark headers as sensitive on both requests and responses. Read more
§

fn sensitive_request_headers( self, headers: impl IntoIterator<Item = HeaderName>, ) -> SetSensitiveRequestHeaders<Self>
where Self: Sized,

Mark headers as sensitive on requests. Read more
§

fn sensitive_response_headers( self, headers: impl IntoIterator<Item = HeaderName>, ) -> SetSensitiveResponseHeaders<Self>
where Self: Sized,

Mark headers as sensitive on responses. Read more
§

fn override_request_header<M>( self, header_name: HeaderName, make: M, ) -> SetRequestHeader<Self, M>
where Self: Sized,

Insert a header into the request. Read more
§

fn append_request_header<M>( self, header_name: HeaderName, make: M, ) -> SetRequestHeader<Self, M>
where Self: Sized,

Append a header into the request. Read more
§

fn insert_request_header_if_not_present<M>( self, header_name: HeaderName, make: M, ) -> SetRequestHeader<Self, M>
where Self: Sized,

Insert a header into the request, if the header is not already present. Read more
§

fn override_response_header<M>( self, header_name: HeaderName, make: M, ) -> SetResponseHeader<Self, M>
where Self: Sized,

Insert a header into the response. Read more
§

fn append_response_header<M>( self, header_name: HeaderName, make: M, ) -> SetResponseHeader<Self, M>
where Self: Sized,

Append a header into the response. Read more
§

fn insert_response_header_if_not_present<M>( self, header_name: HeaderName, make: M, ) -> SetResponseHeader<Self, M>
where Self: Sized,

Insert a header into the response, if the header is not already present. Read more
§

fn set_request_id<M>( self, header_name: HeaderName, make_request_id: M, ) -> SetRequestId<Self, M>
where Self: Sized, M: MakeRequestId,

Add request id header and extension. Read more
§

fn set_x_request_id<M>(self, make_request_id: M) -> SetRequestId<Self, M>
where Self: Sized, M: MakeRequestId,

Add request id header and extension, using x-request-id as the header name. Read more
§

fn propagate_request_id( self, header_name: HeaderName, ) -> PropagateRequestId<Self>
where Self: Sized,

Propgate request ids from requests to responses. Read more
§

fn propagate_x_request_id(self) -> PropagateRequestId<Self>
where Self: Sized,

Propgate request ids from requests to responses, using x-request-id as the header name. Read more
§

fn catch_panic(self) -> CatchPanic<Self, DefaultResponseForPanic>
where Self: Sized,

Catch panics and convert them into 500 Internal Server responses. Read more
§

fn request_body_limit(self, limit: usize) -> RequestBodyLimit<Self>
where Self: Sized,

Intercept requests with over-sized payloads and convert them into 413 Payload Too Large responses. Read more
§

fn trim_trailing_slash(self) -> NormalizePath<Self>
where Self: Sized,

Remove trailing slashes from paths. Read more
§

fn append_trailing_slash(self) -> NormalizePath<Self>
where Self: Sized,

Append trailing slash to paths. Read more
§

impl<SS, SP> SupersetOf<SS> for SP
where SS: SubsetOf<SP>,

§

fn to_subset(&self) -> Option<SS>

The inverse inclusion map: attempts to construct self from the equivalent element of its superset. Read more
§

fn is_in_subset(&self) -> bool

Checks if self is actually part of its subset T (and can be converted to it).
§

fn to_subset_unchecked(&self) -> SS

Use with care! Same as self.to_subset but without any property checks. Always succeeds.
§

fn from_subset(element: &SS) -> SP

The inclusion map: converts self to the equivalent element of its superset.
§

impl<T> Tap for T

§

fn tap(self, func: impl FnOnce(&Self)) -> Self

Immutable access to a value. Read more
§

fn tap_mut(self, func: impl FnOnce(&mut Self)) -> Self

Mutable access to a value. Read more
§

fn tap_borrow<B>(self, func: impl FnOnce(&B)) -> Self
where Self: Borrow<B>, B: ?Sized,

Immutable access to the Borrow<B> of a value. Read more
§

fn tap_borrow_mut<B>(self, func: impl FnOnce(&mut B)) -> Self
where Self: BorrowMut<B>, B: ?Sized,

Mutable access to the BorrowMut<B> of a value. Read more
§

fn tap_ref<R>(self, func: impl FnOnce(&R)) -> Self
where Self: AsRef<R>, R: ?Sized,

Immutable access to the AsRef<R> view of a value. Read more
§

fn tap_ref_mut<R>(self, func: impl FnOnce(&mut R)) -> Self
where Self: AsMut<R>, R: ?Sized,

Mutable access to the AsMut<R> view of a value. Read more
§

fn tap_deref<T>(self, func: impl FnOnce(&T)) -> Self
where Self: Deref<Target = T>, T: ?Sized,

Immutable access to the Deref::Target of a value. Read more
§

fn tap_deref_mut<T>(self, func: impl FnOnce(&mut T)) -> Self
where Self: DerefMut<Target = T> + Deref, T: ?Sized,

Mutable access to the Deref::Target of a value. Read more
§

fn tap_dbg(self, func: impl FnOnce(&Self)) -> Self

Calls .tap() only in debug builds, and is erased in release builds.
§

fn tap_mut_dbg(self, func: impl FnOnce(&mut Self)) -> Self

Calls .tap_mut() only in debug builds, and is erased in release builds.
§

fn tap_borrow_dbg<B>(self, func: impl FnOnce(&B)) -> Self
where Self: Borrow<B>, B: ?Sized,

Calls .tap_borrow() only in debug builds, and is erased in release builds.
§

fn tap_borrow_mut_dbg<B>(self, func: impl FnOnce(&mut B)) -> Self
where Self: BorrowMut<B>, B: ?Sized,

Calls .tap_borrow_mut() only in debug builds, and is erased in release builds.
§

fn tap_ref_dbg<R>(self, func: impl FnOnce(&R)) -> Self
where Self: AsRef<R>, R: ?Sized,

Calls .tap_ref() only in debug builds, and is erased in release builds.
§

fn tap_ref_mut_dbg<R>(self, func: impl FnOnce(&mut R)) -> Self
where Self: AsMut<R>, R: ?Sized,

Calls .tap_ref_mut() only in debug builds, and is erased in release builds.
§

fn tap_deref_dbg<T>(self, func: impl FnOnce(&T)) -> Self
where Self: Deref<Target = T>, T: ?Sized,

Calls .tap_deref() only in debug builds, and is erased in release builds.
§

fn tap_deref_mut_dbg<T>(self, func: impl FnOnce(&mut T)) -> Self
where Self: DerefMut<Target = T> + Deref, T: ?Sized,

Calls .tap_deref_mut() only in debug builds, and is erased in release builds.
§

impl<T> TryConv for T

§

fn try_conv<T>(self) -> Result<T, Self::Error>
where Self: TryInto<T>,

Attempts to convert self into T using TryInto<T>. Read more
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

§

fn vzip(self) -> V

§

impl<T> WithSubscriber for T

§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a [WithDispatch] wrapper. Read more
§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a [WithDispatch] wrapper. Read more
§

impl<G1, G2> Within<G2> for G1
where G2: Contains<G1>,

§

fn is_within(&self, b: &G2) -> bool

§

impl<T> Any for T
where T: Any,

§

impl<T> ErasedDestructor for T
where T: 'static,

§

impl<T> MaybeSend for T
where T: Send,

§

impl<T> MaybeSend for T
where T: Send,