From 7cc2d17cd6b6a8cc85b508431e4ca1aa480f3173 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 02:24:21 +0000 Subject: [PATCH 1/7] docs: inventory backend computation against physical compile coverage Co-Authored-By: Claude Opus 5.5 --- docs/develop_docs/README.md | 1 + .../develop_docs/physical-compile-coverage.md | 69 +++++++++++++++++++ 2 files changed, 70 insertions(+) create mode 100644 docs/develop_docs/physical-compile-coverage.md diff --git a/docs/develop_docs/README.md b/docs/develop_docs/README.md index 2593f3f5..817722ad 100644 --- a/docs/develop_docs/README.md +++ b/docs/develop_docs/README.md @@ -15,5 +15,6 @@ formats, evidence, and verification workflows. - [Metrics-observability corpora](metrics-observability-corpora.md) - [Physical handoff cost references](physical-handoff-costs.md), [storage operations](storage-operation-costs.md) - [Replacement explanations](replacement-explanations.md) +- [Physical compile coverage for deployment computation](physical-compile-coverage.md) - [Planner vocabulary migration (#427)](planner-vocabulary-migration.md) diff --git a/docs/develop_docs/physical-compile-coverage.md b/docs/develop_docs/physical-compile-coverage.md new file mode 100644 index 00000000..4dc309ae --- /dev/null +++ b/docs/develop_docs/physical-compile-coverage.md @@ -0,0 +1,69 @@ +# Physical compile coverage for deployment computation + +Audience: developers moving computation from ASAPQuery-backend into +`asap_physical_operators::physical_planner`. + +## Contract + +Logical selection decides what to compute. The maintenance lifecycle sets node +timing. `physical_planner::compile` turns a timed `PostAsapDag` into physical +operator DAGs. The backend owns ingestion, panes, storage, stored-state +readout, external exact engines, pricing/selection, and execution scheduling. + +A backend lowering is *covered* when `compile` accepts the corresponding +`PostAsapDag` node and produces operators with the same result. The backend +should then pass the timed DAG and its input contracts to `compile`. It should +not rebuild operator choices from PromQL text or construct operators itself. + +## Inventory + +Surveyed backend: `ASAPQuery-backend` branch `perf/788-startup-search`. +Planner base: `split/462-f-physical-planner` (#475). + +Status values: + +- **Supported**: `compile` or `compile_node` already covers this computation. +- **Partial**: some shapes are covered. The Notes column lists the gap. +- **Missing**: `compile` rejects this computation. +- **Backend**: not computation, or owned by the backend. + +| # | Backend site | Computation | Planner node | Status at #475 | Notes | +|---|---|---|---|---|---| +| 1 | `query_time.rs` `Lower::lower`, `compile_logical` | PromQL AST → `QueryTimeOperator` graph for a native query | `Fallback { QueryExpr }` subtrees plus value payloads | Missing | `compile` lowers `Fallback` only as a raw `Scan` source. | +| 2 | `QueryTimeOperator::Aggregate` (sum/min/max/avg/count) | Grouped value aggregation | `Value::Exact(Aggregate)`; `SummaryAgg{ExactAggregate, Reduce}` over finalized values | Supported | Also `promql_values::compile_aggregate`. | +| 3 | `QueryTimeOperator::Sort`, `Limit` (topk, sort, sort_desc) | Ordering and per-group limits | `Value::Sort`, `Value::Limit` | Supported | | +| 4 | `QueryTimeOperator::Binary`, `QueryPlanNode::Binary` (vector ⊗ scalar) | Arithmetic with a scalar operand | `Binary` whose operand is `Fallback{PromqlScalarBridge(Literal)}` | Missing | Query-time `Binary` accepts only label-map vector schemas. The literal node has no native binding. | +| 5 | `QueryTimeOperator::Binary` (vector ⊗ vector) | One-to-one label matching and arithmetic | `Binary` over grouped value rows | Missing | Only the ingestion-time `aligned_binary` and label-map `vector_binary` exist. | +| 6 | `binary_operator` CheckedDiv / FiniteDiv | Guarded division | `BinaryOperator` checked flags | Partial | Flags are evaluated, but only where rows 4/5 are covered. | +| 7 | `QueryTimeOperator::Binary` comparisons, `bool` | Filter or 0/1 comparison | `Binary{Compare}` | Missing | `Payload::Binary` does not carry `return_bool`. | +| 8 | `QueryTimeOperator::UnaryNegate` | Negation | `Binary{Mul}` by literal `-1` | Missing | The frontend emits `* -1`; same gap as row 4. | +| 9 | `QueryTimeOperator::VectorToScalar` | `scalar()` | `Fallback{PromqlScalarFromVector}` | Missing | Only `promql_values::compile_vector_to_scalar`. | +| 10 | `QueryTimeOperator::HistogramQuantile` | Bucket interpolation | `Fallback` / `AggIntent::HistogramQuantile` | Missing | Only `promql_values::compile_histogram_quantile`. | +| 11 | `QueryTimeOperator::Temporal` (rate, increase, `*_over_time`) | Per-series window functions | `SummaryAgg{PerEntity}` over `TimeRange(Scan)` | Partial | Supported with closed series identity. Not supported over `Fallback` matrices (`compile_temporal` only). | +| 12 | `logical_dag.rs` `Subquery`, `subquery_grid`, `expanded_inputs` | Re-evaluate the child on a step grid and assemble a matrix | `Fallback{PromqlSubquery}` | Missing | No Planner operator. | +| 13 | `QueryPlanNode::Scalar`, `DagCompiler::lower` scalar literal | Scalar constant | `Fallback{PromqlScalarBridge(Literal)}` | Missing | Only `promql_values::compile_scalar`. | +| 14 | `DagCompiler::lower` `ReduceSum`; `physical_values.rs` PerEntity projection | Sum over finalized values; per-entity identity | `SummaryAgg{ExactAggregate(Sum)}` | Supported | The backend builds an identity `Operator::project` itself for PerEntity. | +| 15 | `DagCompiler::lower` `ExactReadout`; `post_asap_readout.rs` ExactReadout | Finalize exact state (sum/count/min/max/rate/increase) | `Value::FinalizeExactAccumulator` | Partial | Count yields Int64 against a declared Float64 PromQL value. `compile` rejects it. | +| 16 | `post_asap_readout.rs` SummaryEstimate (`readout_bound`, `expand_item_rows`) | Sketch estimate per group; TopK item expansion | `SummaryEstimate` | Partial | The backend's label-map state layout and MetricsQL `__name__` rules have no Planner equivalent. `compile_exact_readout` has no sketch counterpart. | +| 17 | `post_asap_readout.rs` SummaryMerge (`merge_bound_states`) | Merge states by group | `SummaryMerge` | Supported | Union plus `summary_merge`. | +| 18 | `post_asap_readout.rs` counter range parameters | Counter lookback for rate/increase | `TimeRange` ancestor of finalization | Supported | Applied through `with_counter_lookback`. | +| 19 | `post_asap_readout.rs` `execute_value_fragment` | Per-timestamp binding of a value fragment | n/a | Backend | Evaluation scheduling. | +| 20 | `DagCompiler::lower` SummaryJoin / Subtract / Delete | Summary algebra | `SummaryJoin`, `SummarySubtract`, `SummaryDelete` | Missing | The backend also rejects these (`ExactFallback`). | +| 21 | `current_series.rs` Snapshot + TopK | Current-series ranking | `ReadPopulation{TopK}` | Supported | | +| 22 | `current_series.rs` Sum / Count / Average | Current-series aggregates | `ReadPopulation{Sum,Count,Average}` | Missing | `compile` accepts only TopK. | +| 23 | `current_series.rs` Quantile | Current-series quantile | `ReadPopulation{Quantile}` | Missing | No exact quantile reduction. | +| 24 | `raw_dag.rs` weight `Column` | Summary update from a sample/projected value | `SummaryAgg` | Supported | | +| 25 | `raw_dag.rs` weight `Constant` | Unit/constant-weight update | `SummaryAgg` | Missing | `compile_node` requires a column weight. | +| 26 | `raw_dag.rs` item `Column` / `Tuple` | Keyed update item | `SummaryAgg{item}` | Supported | `keyed_summary_build`. | +| 27 | `raw_dag.rs` item `EntityIdentity` | Series-identity item | `SummaryAgg{item}` | Missing | Needs the series-identity column. | +| 28 | `physical_values.rs` `compile`, `combine` | Translate `QueryTimeOperator` to `promql_values::*`; compose fragments | n/a | Supported | Exists only because of row 1. `CompiledPhysicalDag::compose` is Planner API. | +| 29 | `query_plan.rs` `compile_native_fragment` (Semi join, Exact aggregate, Sort, Limit, Filter) | Relational value ops | `RelationalJoin`, `Value::*` | Supported | Already calls `compile`. | +| 30 | `query_time.rs` `selected_query_time_nodes`, `selected_native_expression`, `selected_aggregate_operator` | Recover operator identity from original PromQL text | Payload variants (`ExactKind::Min`/`Max`, `AggIntent`) | Supported | Payloads already carry the identity. These witnesses are needed only while row 1 remains. | +| 31 | Scan, ExactSubquery, CandidateExactSubquery, CurrentSeries ingest, ReadMaterialization, ExternalExact | Storage reads and external engines | Input contracts | Backend | | + +Totals at #475: 11 Supported, 4 Partial, 14 Missing, 2 Backend. + +## Remaining + +Every Missing or Partial row above, in order of backend usage. Rows 1 and 28 +are removed from the backend only after rows 4–13 are covered. From a8c5d3f6c4584878491767d6d8821736b2c8cc41 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 02:44:41 +0000 Subject: [PATCH 2/7] feat(physical): compile maintained-population aggregate readouts ReadPopulation Sum/Count/Average/Quantile now compile to a grouped aggregate over the population snapshot, so deployments no longer evaluate these readouts in their current-series store. Quantile uses a new exact Reduction::Quantile with PromQL rank interpolation. Co-Authored-By: Claude Opus 5.5 --- .../src/operators/aggregate/mod.rs | 55 +++++ .../src/physical_planner/mod.rs | 13 +- .../src/physical_planner/row_values.rs | 31 +++ .../tests/deployment_computation.rs | 200 ++++++++++++++++++ 4 files changed, 294 insertions(+), 5 deletions(-) create mode 100644 crates/asap-physical-operators/src/physical_planner/row_values.rs create mode 100644 crates/asap-physical-operators/tests/deployment_computation.rs diff --git a/crates/asap-physical-operators/src/operators/aggregate/mod.rs b/crates/asap-physical-operators/src/operators/aggregate/mod.rs index 71e7134c..0df3c157 100644 --- a/crates/asap-physical-operators/src/operators/aggregate/mod.rs +++ b/crates/asap-physical-operators/src/operators/aggregate/mod.rs @@ -27,6 +27,12 @@ impl Operator { false, ) } + Reduction::Quantile { column, q } => { + if plain(&input, *column)?.0 != &DataType::Float64 || q.is_nan() { + return Err(invalid("quantile requires Float64 input and a numeric q")); + } + (DataType::Float64, false) + } Reduction::Min(i) | Reduction::Max(i) => { let (t, nullable) = plain(&input, *i)?; if !ordered(t) { @@ -120,6 +126,11 @@ pub enum Reduction { Avg(usize), Min(usize), Max(usize), + /// PromQL `quantile`: linear interpolation between closest ranks. + Quantile { + column: usize, + q: f64, + }, } pub(super) fn execute<'a>( operator: &'a Operator, @@ -198,6 +209,25 @@ async fn reduce( Ok(output) } +// Matches Prometheus `quantile`: NaN for no values, ±Inf outside [0, 1]. +fn quantile(q: f64, mut values: Vec) -> f64 { + if values.is_empty() { + return f64::NAN; + } + if q < 0. { + return f64::NEG_INFINITY; + } + if q > 1. { + return f64::INFINITY; + } + values.sort_by(f64::total_cmp); + let rank = q * (values.len() - 1) as f64; + let low = rank.floor() as usize; + let high = (low + 1).min(values.len() - 1); + let weight = rank - low as f64; + values[low] * (1. - weight) + values[high] * weight +} + async fn reduce_one( rows: &[Vec], measure: &Reduction, @@ -211,6 +241,18 @@ async fn reduce_one( )) } Reduction::Sum(i) | Reduction::Avg(i) | Reduction::Min(i) | Reduction::Max(i) => *i, + Reduction::Quantile { column, q } => { + let mut values = Vec::with_capacity(rows.len()); + for row in rows { + work.checkpoint().await?; + match &row[*column] { + Value::Float64(value) => values.push(*value), + Value::Null => {} + _ => return Err(invalid("floating quantile value required")), + } + } + return Ok(Value::Float64(quantile(*q, values))); + } }; let values = rows .iter() @@ -291,3 +333,16 @@ async fn reduce_one( sum })) } + +#[cfg(test)] +mod tests { + // Quantile follows Prometheus: interpolate ranks, NaN when empty, ±Inf outside [0, 1]. + #[test] + fn quantile_matches_prometheus_edge_cases() { + assert!(super::quantile(0.5, vec![]).is_nan()); + assert_eq!(super::quantile(0.5, vec![3.]), 3.); + assert_eq!(super::quantile(0.75, vec![4., 1., 2., 3.]), 3.25); + assert_eq!(super::quantile(-0.1, vec![1.]), f64::NEG_INFINITY); + assert_eq!(super::quantile(1.1, vec![1.]), f64::INFINITY); + } +} diff --git a/crates/asap-physical-operators/src/physical_planner/mod.rs b/crates/asap-physical-operators/src/physical_planner/mod.rs index fede3777..3896e2a2 100644 --- a/crates/asap-physical-operators/src/physical_planner/mod.rs +++ b/crates/asap-physical-operators/src/physical_planner/mod.rs @@ -43,6 +43,8 @@ pub use candidates::{ mod compiled; pub use compiled::{CompiledPhysicalDag, InputContract}; +mod row_values; + /// Compile computation without opening or retaining deployment readers. /// Input contracts identify explicit boundaries selected by maintenance planning. pub fn compile( @@ -247,11 +249,6 @@ fn compile_internal( use planner_types::post_asap::maintained_population::{ PopulationInput, PopulationReadout, }; - let PopulationReadout::TopK { k } = readout else { - return Err(invalid( - "native population readout does not support this operation", - )); - }; let [producer] = inputs.as_slice() else { return Err(invalid("population readout requires one input")); }; @@ -272,6 +269,12 @@ fn compile_internal( )); } let input = schemas[0].clone(); + let PopulationReadout::TopK { k } = readout else { + let aggregate = + row_values::population_aggregate(&input, &spec.grouping, readout)?; + graph.add(id, inputs, aggregate.with_output_schema(output)?)?; + continue; + }; let groups = spec .grouping .iter() diff --git a/crates/asap-physical-operators/src/physical_planner/row_values.rs b/crates/asap-physical-operators/src/physical_planner/row_values.rs new file mode 100644 index 00000000..954da893 --- /dev/null +++ b/crates/asap-physical-operators/src/physical_planner/row_values.rs @@ -0,0 +1,31 @@ +//! Query-time PromQL value computation over logical row schemas. +use super::*; +use planner_types::post_asap::maintained_population::PopulationReadout; + +/// Aggregate readouts of a maintained current-series population. +pub(super) fn population_aggregate( + input: &Schema, + grouping: &[String], + readout: &PopulationReadout, +) -> Result { + let groups = grouping + .iter() + .map(|name| named_column(input, &ColumnRef::Named(name.clone()))) + .collect::, _>>()?; + let value = named_column(input, &ColumnRef::SampleValue)?; + let reduction = match readout { + PopulationReadout::Sum => Reduction::Sum(value), + PopulationReadout::Count => Reduction::Count, + PopulationReadout::Average => Reduction::Avg(value), + PopulationReadout::Quantile { q } => Reduction::Quantile { + column: value, + q: *q, + }, + PopulationReadout::TopK { .. } => { + return Err(invalid( + "TopK population readout ranks; it does not aggregate", + )) + } + }; + Operator::aggregate(input.clone(), groups, vec![("value".into(), reduction)]) +} diff --git a/crates/asap-physical-operators/tests/deployment_computation.rs b/crates/asap-physical-operators/tests/deployment_computation.rs new file mode 100644 index 00000000..70880fd1 --- /dev/null +++ b/crates/asap-physical-operators/tests/deployment_computation.rs @@ -0,0 +1,200 @@ +//! Planner-selected PromQL computation compiles from the timed DAG alone; +//! the deployment supplies only raw rows at the ingestion frontier. +use asap_physical_operators::{ + operators::Operator, + physical_planner::{compile, promql_rows, CompiledPhysicalDag, InputContract, Source}, + runtime::{Limits, RunContext, Scope}, + values::{Batch, Value}, +}; +use futures::{executor::block_on, StreamExt}; +use planner_types::{post_asap::*, pre_asap::QueryExpr, types::AccuracyTarget, workload::*}; +use std::{collections::BTreeMap, rc::Rc, sync::Arc}; + +fn lower(query: &str) -> QueryExpr { + let workload = PlanningWorkload { + query_workload: QueryWorkload { + language: QueryLanguage::PromQL, + query_batch: Some(vec![BatchEntry { + query: Query(query.into()), + requirements: QueryRequirements { + accuracy: AccuracyRequirement::Explicit(AccuracyTarget::Exact), + ..Default::default() + }, + predictability: Predictability::Unknown, + invocations: 1, + execute_at: None, + time_selection: TimeSelection::default(), + }]), + repeating_queries: None, + }, + data_workload: Some(DataWorkload { + data_ingestion_interval: Evidence { + value: Some(DurationMs(60_000)), + ..Default::default() + }, + ..Default::default() + }), + }; + asap_frontend_promql::lower_promql_workload(&workload, 0) + .unwrap() + .remove(0) +} + +fn population_dag(query: &str) -> PostAsapDag { + let root = Rc::new(promql_rows::with_series_identity(&lower(query)).unwrap()); + let selected = asap_aware_mapping::maintained_population::MaintainedPopulationStrategy::new( + std::slice::from_ref(&root), + ) + .candidate(&root) + .unwrap(); + compile_post_asap_dag(&selected).unwrap() +} + +/// Raw scan nodes are the frontier; everything above them is compiled. +fn raw_inputs(dag: &PostAsapDag) -> Vec<(u64, Arc, String)> { + dag.nodes + .iter() + .filter_map(|node| match &node.payload { + PostAsapOperatorPayload::Fallback { + expression: QueryExpr::TimeRange { child, .. }, + } => match child.as_ref() { + QueryExpr::Scan { + source: planner_types::pre_asap::Source::TimeSeries { metric }, + .. + } => Some(( + u64::from(node.id.0), + Arc::new(node.output_schema.clone()), + metric.clone(), + )), + _ => None, + }, + _ => None, + }) + .collect() +} + +type Sample = (&'static str, &'static str, &'static str, i64, f64); + +/// Compile, round-trip, bind raw `(metric, job, instance, ts, value)` samples, +/// and return `(job, value)` rows of the root. +fn run(dag: &PostAsapDag, samples: &[Sample], end: i64) -> Result, String> { + let inputs = raw_inputs(dag); + let program = compile( + dag, + inputs + .iter() + .map(|(id, schema, _)| (*id, InputContract::bounded(schema.clone()))) + .collect(), + &[u64::from(dag.root.0)], + ) + .map_err(|e| e.to_string())?; + let program: CompiledPhysicalDag = + serde_json::from_slice(&serde_json::to_vec(&program).unwrap()).unwrap(); + let sources = inputs + .iter() + .map(|(id, schema, metric)| { + let rows = samples + .iter() + .filter(|sample| sample.0 == metric) + .map(|(name, job, instance, at, value)| { + let labels = BTreeMap::from([ + ("__name__".to_string(), name.to_string()), + ("job".into(), job.to_string()), + ("instance".into(), instance.to_string()), + ]); + if schema + .fields + .iter() + .any(|f| f.name == promql_rows::SERIES_IDENTITY_COLUMN) + { + promql_rows::series_row(schema, &labels, *at, *value).unwrap() + } else { + schema + .fields + .iter() + .enumerate() + .map(|(i, f)| match f.name.as_str() { + _ if Some(i) == schema.time_index => Value::Timestamp(*at), + "value" => Value::Float64(*value), + label => Value::Utf8(labels[label].clone().into()), + }) + .collect() + } + }) + .collect(); + let batch = Batch::try_new(schema.clone(), rows).unwrap(); + ( + *id, + Box::new(Operator::source(schema.clone(), vec![batch]).unwrap()) as Source<'_>, + ) + }) + .collect(); + let graph = program.instantiate(sources).map_err(|e| e.to_string())?; + let context = RunContext::new( + Scope::Query { + evaluation_time_ms: end, + revision: 0, + }, + Limits::default(), + ) + .unwrap(); + block_on(async { + let mut stream = graph + .execute(program.roots(), context) + .map_err(|e| e.to_string())? + .remove(0); + let mut rows = BTreeMap::new(); + while let Some(batch) = stream.next().await { + let batch = batch.map_err(|e| e.to_string())?; + let job = batch.schema().fields.iter().position(|f| f.name == "job"); + for row in batch.rows() { + let key = match job.map(|i| &row[i]) { + Some(Value::Utf8(job)) => job.to_string(), + _ => String::new(), + }; + let value = match row.last() { + Some(Value::Float64(v)) => *v, + Some(Value::Int64(v)) => *v as f64, + other => return Err(format!("unexpected value {other:?}")), + }; + assert!(rows.insert(key, value).is_none(), "duplicate output group"); + } + } + Ok(rows) + }) +} + +const SAMPLES: &[Sample] = &[ + ("m", "api", "a", 10_000, 4.), + ("m", "api", "a", 50_000, 1.), + ("m", "api", "b", 40_000, 7.), + ("m", "api", "c", 30_000, 2.), + ("m", "db", "d", 20_000, 5.), +]; + +fn reference(pairs: &[(&str, f64)]) -> BTreeMap { + pairs.iter().map(|(k, v)| (k.to_string(), *v)).collect() +} + +// Current-series aggregates read the latest member values, matching the +// backend CurrentSeriesStore formulas (PromQL quantile interpolation). +#[test] +fn population_aggregates_match_current_series_reference() { + // Latest values: api = {a: 1, b: 7, c: 2}; db = {d: 5}. + for (query, expected) in [ + ("sum by (job) (m)", reference(&[("api", 10.), ("db", 5.)])), + ("count by (job) (m)", reference(&[("api", 3.), ("db", 1.)])), + ( + "avg by (job) (m)", + reference(&[("api", 10. / 3.), ("db", 5.)]), + ), + // Sorted api = [1, 2, 7]; rank 0.25 * 2 = 0.5 → 1.5. + ( + "quantile by (job) (0.25, m)", + reference(&[("api", 1.5), ("db", 5.)]), + ), + ] { + let dag = population_dag(query); + assert_eq!(run(&dag, SAMPLES, 60_000).unwrap(), expected, "{query}"); + } +} From 362203a43d8dc4040a8794708a6263055350f4ba Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 02:45:59 +0000 Subject: [PATCH 3/7] feat(physical): compile query-time arithmetic over grouped value rows A query-time Binary with a PromQL scalar-literal operand folds the literal into a projection, which also covers unary negation. Two grouped row inputs match one-to-one on equal label columns through an inner equi-join before the operator is applied. Comparisons and per-series rows still fail at compile time. Co-Authored-By: Claude Opus 5.5 --- .../src/physical_planner/mod.rs | 56 +++++- .../src/physical_planner/row_values.rs | 170 +++++++++++++++++- .../tests/deployment_computation.rs | 84 +++++++++ 3 files changed, 308 insertions(+), 2 deletions(-) diff --git a/crates/asap-physical-operators/src/physical_planner/mod.rs b/crates/asap-physical-operators/src/physical_planner/mod.rs index 3896e2a2..102f32b0 100644 --- a/crates/asap-physical-operators/src/physical_planner/mod.rs +++ b/crates/asap-physical-operators/src/physical_planner/mod.rs @@ -14,7 +14,8 @@ use planner_types::{ SketchQuery, SummaryFamilyType, SummaryInputExpr, ValueOperation, }, pre_asap::{ - AggIntent, ColumnRef, CompareOpKind, GroupKeys, QueryExpr, Reduction as PlannerReduction, + AggIntent, ColumnRef, CompareOpKind, DataType, GroupKeys, QueryExpr, + Reduction as PlannerReduction, }, }; use std::{ @@ -145,7 +146,28 @@ fn compile_internal( }, ) }); + // Scalar literal operands of query-time arithmetic are folded into the consumer. + let mut literals = BTreeMap::::new(); for edge in edges { + let consumer = u64::from(edge.consumer.0); + if let ( + Payload::Fallback { expression }, + Some(PostAsapDagNode { + payload: Payload::Binary { .. }, + .. + }), + ) = ( + &nodes[&u64::from(edge.producer.0)].payload, + nodes.get(&consumer), + ) { + if let Some(value) = row_values::scalar_literal(expression) { + let left = edge.role == planner_types::post_asap::EdgeRole::Left; + if literals.insert(consumer, (value, left)).is_some() { + return Err(invalid("binary with two scalar literals is not folded")); + } + continue; + } + } dependencies .entry(u64::from(edge.consumer.0)) .or_default() @@ -360,6 +382,38 @@ fn compile_internal( )?; continue; } + if let Payload::Binary { operator } = &node.payload { + let query_time = node.output_state.timing + == planner_types::post_asap::ExecutionTiming::QueryTime; + if let Some(&(value, left)) = literals.get(&id) { + let [input] = schemas.as_slice() else { + return Err(invalid("scalar binary requires one row input")); + }; + if !query_time { + return Err(invalid("scalar literal binary must run at query time")); + } + let project = row_values::scalar_binary(input, operator, value, left) + .map_err(|error| invalid(format!("node {id}: {error}")))?; + graph.add(id, inputs, project.with_output_schema(output)?)?; + continue; + } + let label_map = |schema: &Schema| { + schema + .fields + .iter() + .any(|f| matches!(f.dtype, SummaryFamilyType::Plain(DataType::Map { .. }))) + }; + if let (true, [left, right]) = (query_time, schemas.as_slice()) { + if !label_map(left) && !label_map(right) { + let (join, project) = row_values::grouped_binary(left, right, operator) + .map_err(|error| invalid(format!("node {id}: {error}")))?; + graph.add(auxiliary, inputs, join)?; + graph.add(id, vec![auxiliary], project.with_output_schema(output)?)?; + auxiliary -= 1; + continue; + } + } + } let mut operator = compile_node(node, &schemas) .map_err(|error| invalid(format!("node {id}: {error}")))?; if operator.is_counter_readout() { diff --git a/crates/asap-physical-operators/src/physical_planner/row_values.rs b/crates/asap-physical-operators/src/physical_planner/row_values.rs index 954da893..10c93c63 100644 --- a/crates/asap-physical-operators/src/physical_planner/row_values.rs +++ b/crates/asap-physical-operators/src/physical_planner/row_values.rs @@ -1,6 +1,174 @@ //! Query-time PromQL value computation over logical row schemas. use super::*; -use planner_types::post_asap::maintained_population::PopulationReadout; +use planner_types::post_asap::{maintained_population::PopulationReadout, BinaryOperator}; +use planner_types::pre_asap::{BinaryOpKind, DataType, Predicate, ScalarValue}; +use std::rc::Rc; + +/// A PromQL number literal has no row schema; its consumer folds it in. +pub(super) fn scalar_literal(expression: &QueryExpr) -> Option { + match expression { + QueryExpr::PromqlScalarBridge(child) => scalar_literal(child), + QueryExpr::Literal(ScalarValue::Float64(value)) => Some(*value), + _ => None, + } +} + +/// Rows without a time column, label map, or series identity carry only +/// their group labels, so those labels are the complete PromQL identity. +fn grouped_value(input: &Schema) -> Result<(usize, Vec), Error> { + if input.time_index.is_some() + || input + .fields + .iter() + .any(|field| field.name == promql_rows::SERIES_IDENTITY_COLUMN) + { + return Err(invalid( + "row binary requires grouped rows; per-series matching needs a name-free identity", + )); + } + let mut value = None; + let mut labels = Vec::new(); + for (i, field) in input.fields.iter().enumerate() { + match &field.dtype { + SummaryFamilyType::Plain(DataType::Float64) if value.is_none() => value = Some(i), + SummaryFamilyType::Plain(DataType::Utf8) => labels.push(i), + _ => { + return Err(invalid( + "row binary requires Utf8 labels and one Float64 value", + )) + } + } + } + Ok(( + value.ok_or_else(|| invalid("row binary requires a Float64 value"))?, + labels, + )) +} + +fn arithmetic(operator: &BinaryOperator) -> Result<(), Error> { + if !matches!(operator.kind, BinaryOpKind::Arithmetic(_)) { + return Err(invalid( + "row comparison requires filter or bool semantics, which Binary does not carry", + )); + } + Ok(()) +} + +/// Apply `vector op scalar` (or `scalar op vector`) to each row's value. +pub(super) fn scalar_binary( + input: &Schema, + operator: &BinaryOperator, + literal: f64, + literal_left: bool, +) -> Result { + arithmetic(operator)?; + let (value, _) = grouped_value(input)?; + let literal = Expression::Literal { + value: crate::values::Value::Float64(literal), + dtype: DataType::Float64, + }; + let columns = input + .fields + .iter() + .enumerate() + .map(|(i, field)| { + let expression = if i != value { + Expression::Column(i) + } else if literal_left { + binary(operator, literal.clone(), Expression::Column(i)) + } else { + binary(operator, Expression::Column(i), literal.clone()) + }; + (field.name.clone(), expression) + }) + .collect(); + Operator::project(input.clone(), columns) +} + +/// One-to-one PromQL matching of grouped rows on equal label sets. Returns +/// the inner equi-join and the projection that applies the operator. +pub(super) fn grouped_binary( + left: &Schema, + right: &Schema, + operator: &BinaryOperator, +) -> Result<(Operator, Operator), Error> { + arithmetic(operator)?; + let (left_value, left_labels) = grouped_value(left)?; + let (right_value, right_labels) = grouped_value(right)?; + if left_labels.len() != right_labels.len() { + return Err(invalid("row binary inputs have different label sets")); + } + let width = left.fields.len(); + let keys = left_labels + .iter() + .map(|&l| { + let name = &left.fields[l].name; + let r = right_labels + .iter() + .copied() + .find(|&r| &right.fields[r].name == name) + .ok_or_else(|| invalid("row binary inputs have different label sets"))?; + let (a, b) = ( + Rc::new(QueryExpr::Column(l)), + Rc::new(QueryExpr::Column(width + r)), + ); + let equal = QueryExpr::Compare { + left: a.clone(), + op: CompareOpKind::Eq, + right: b.clone(), + }; + // A nullable label compares like PromQL's empty label: absent on both sides matches. + Ok(if left.fields[l].nullable || right.fields[r].nullable { + QueryExpr::BoolOr(vec![ + equal, + QueryExpr::BoolAnd(vec![QueryExpr::IsNull(a), QueryExpr::IsNull(b)]), + ]) + } else { + equal + }) + }) + .collect::, Error>>()?; + let predicate = Predicate(Rc::new(QueryExpr::BoolAnd(keys))); + let mut joined = left.fields.clone(); + joined.extend(right.fields.iter().cloned()); + let join = Operator::relational_join( + left.clone(), + right.clone(), + planner_types::pre_asap::JoinKind::Inner, + &predicate, + Arc::new(planner_types::post_asap::SummarySchema { + fields: joined, + time_index: None, + }), + )?; + let columns = left + .fields + .iter() + .enumerate() + .map(|(i, field)| { + let expression = if i == left_value { + binary( + operator, + Expression::Column(i), + Expression::Column(width + right_value), + ) + } else { + Expression::Column(i) + }; + (field.name.clone(), expression) + }) + .collect(); + let project = Operator::project(join.schema(), columns)?; + Ok((join, project)) +} + +fn binary(operator: &BinaryOperator, left: Expression, right: Expression) -> Expression { + Expression::Binary { + operator: operator.clone(), + left: Box::new(left), + right: Box::new(right), + } +} /// Aggregate readouts of a maintained current-series population. pub(super) fn population_aggregate( diff --git a/crates/asap-physical-operators/tests/deployment_computation.rs b/crates/asap-physical-operators/tests/deployment_computation.rs index 70880fd1..25899ec6 100644 --- a/crates/asap-physical-operators/tests/deployment_computation.rs +++ b/crates/asap-physical-operators/tests/deployment_computation.rs @@ -40,6 +40,27 @@ fn lower(query: &str) -> QueryExpr { .remove(0) } +/// The first exact summary candidate, as Planner selection would hand it over. +fn exact_dag(query: &str) -> PostAsapDag { + use asap_aware_mapping::{Replacement, ReplacementStrategy, TargetSubDAG}; + let expression = lower(query); + let root = Rc::new(promql_rows::with_series_identity(&expression).unwrap_or(expression)); + asap_aware_mapping::SketchAlgorithmStrategy::new(&asap_aware_mapping::DefaultCostModel) + .replacements(&TargetSubDAG::new(&root)) + .into_iter() + .find_map(|candidate| match candidate.replacement { + Replacement::Summary(node) => { + let dag = compile_post_asap_dag(&node).ok()?; + dag.nodes + .iter() + .all(|n| !matches!(&n.payload, PostAsapOperatorPayload::SummaryAgg { family, .. } if !matches!(family, SummaryFamilyType::ExactAggregate(..)))) + .then_some(dag) + } + _ => None, + }) + .unwrap() +} + fn population_dag(query: &str) -> PostAsapDag { let root = Rc::new(promql_rows::with_series_identity(&lower(query)).unwrap()); let selected = asap_aware_mapping::maintained_population::MaintainedPopulationStrategy::new( @@ -198,3 +219,66 @@ fn population_aggregates_match_current_series_reference() { assert_eq!(run(&dag, SAMPLES, 60_000).unwrap(), expected, "{query}"); } } + +// Scalar operands on either side apply to every grouped value, including negation. +#[test] +fn scalar_literal_arithmetic_applies_to_grouped_values() { + // sum_over_time over 5m per job: api = 4 + 1 + 7 + 2 = 14, db = 5. + for (query, expected) in [ + ( + "sum by (job) (sum_over_time(m[5m])) * 2", + reference(&[("api", 28.), ("db", 10.)]), + ), + ( + "100 - sum by (job) (sum_over_time(m[5m]))", + reference(&[("api", 86.), ("db", 95.)]), + ), + ( + "-sum by (job) (sum_over_time(m[5m]))", + reference(&[("api", -14.), ("db", -5.)]), + ), + ] { + assert_eq!( + run(&exact_dag(query), SAMPLES, 60_000).unwrap(), + expected, + "{query}" + ); + } +} + +// Grouped vectors match one-to-one on labels; unmatched groups are dropped and +// unchecked division by zero yields +Inf as in PromQL. +#[test] +fn grouped_vector_arithmetic_matches_labels() { + let samples: &[Sample] = &[ + ("a", "api", "x", 10_000, 6.), + ("a", "api", "y", 20_000, 3.), + ("a", "db", "x", 10_000, 1.), + ("a", "web", "x", 10_000, 1.), + ("b", "api", "x", 10_000, 3.), + ("b", "db", "x", 10_000, 0.), + ("b", "cache", "x", 10_000, 1.), + ]; + let dag = + exact_dag("sum by (job) (sum_over_time(a[5m])) / sum by (job) (sum_over_time(b[5m]))"); + assert_eq!( + run(&dag, samples, 60_000).unwrap(), + reference(&[("api", 3.), ("db", f64::INFINITY)]) + ); +} + +// Comparisons need filter/bool semantics that `Binary` does not carry, so +// they fail at compile time instead of emitting 0/1 values. +#[test] +fn row_comparison_fails_closed() { + let mut dag = exact_dag("sum by (job) (sum_over_time(m[5m])) * 2"); + for node in &mut dag.nodes { + if let PostAsapOperatorPayload::Binary { operator } = &mut node.payload { + operator.kind = planner_types::pre_asap::BinaryOpKind::Compare( + planner_types::pre_asap::CompareOpKind::Gt, + ); + } + } + let error = run(&dag, SAMPLES, 60_000).unwrap_err(); + assert!(error.contains("comparison"), "{error}"); +} From 49cf0446dfdb062bff27ed27c872df5b2e8f1367 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 02:45:59 +0000 Subject: [PATCH 4/7] fix(physical): finalize exact counts to declared Float64 values Exact count readout yields Int64, but PromQL declares a Float64 sample, so compile rejected count finalization. Convert exactly. Co-Authored-By: Claude Opus 5.5 --- .../src/physical_planner/mod.rs | 35 +++++++++++++++ .../tests/deployment_computation.rs | 44 +++++++++++++++++++ 2 files changed, 79 insertions(+) diff --git a/crates/asap-physical-operators/src/physical_planner/mod.rs b/crates/asap-physical-operators/src/physical_planner/mod.rs index 102f32b0..afc1ff1f 100644 --- a/crates/asap-physical-operators/src/physical_planner/mod.rs +++ b/crates/asap-physical-operators/src/physical_planner/mod.rs @@ -414,6 +414,41 @@ fn compile_internal( } } } + if let Payload::Value { + operation: ValueOperation::FinalizeExactAccumulator, + } = &node.payload + { + // Exact counts read out as Int64; PromQL declares a Float64 sample. + let readout = bind_operation(node, &schemas) + .map_err(|error| invalid(format!("node {id}: {error}")))?; + let actual = readout.schema(); + let converted = actual.fields.iter().zip(&output.fields).position(|(a, d)| { + a.dtype == SummaryFamilyType::Plain(DataType::Int64) + && d.dtype == SummaryFamilyType::Plain(DataType::Float64) + }); + if let Some(column) = converted { + let columns = actual + .fields + .iter() + .enumerate() + .map(|(i, field)| { + ( + field.name.clone(), + if i == column { + Expression::ExactFloat64(i) + } else { + Expression::Column(i) + }, + ) + }) + .collect(); + let project = Operator::project(actual, columns)?.with_output_schema(output)?; + graph.add(auxiliary, inputs, readout)?; + graph.add(id, vec![auxiliary], project)?; + auxiliary -= 1; + continue; + } + } let mut operator = compile_node(node, &schemas) .map_err(|error| invalid(format!("node {id}: {error}")))?; if operator.is_counter_readout() { diff --git a/crates/asap-physical-operators/tests/deployment_computation.rs b/crates/asap-physical-operators/tests/deployment_computation.rs index 25899ec6..f08cfcae 100644 --- a/crates/asap-physical-operators/tests/deployment_computation.rs +++ b/crates/asap-physical-operators/tests/deployment_computation.rs @@ -267,6 +267,50 @@ fn grouped_vector_arithmetic_matches_labels() { ); } +// Exact observation counts finalize to the Float64 value PromQL declares, +// then roll up per job: api has 2 + 1 + 1 samples in 5m, db has 1. +#[test] +fn exact_count_finalizes_to_declared_float_value() { + let mut dag = exact_dag("sum by (job) (count_over_time(m[5m]))"); + let finalize = dag + .nodes + .iter() + .find(|node| { + matches!( + node.payload, + PostAsapOperatorPayload::Value { + operation: ValueOperation::FinalizeExactAccumulator + } + ) + }) + .unwrap() + .clone(); + let root = dag.nodes.iter().find(|n| n.id == dag.root).unwrap().clone(); + let mut edge = dag + .edges + .iter() + .find(|e| e.producer == finalize.id) + .unwrap() + .clone(); + // Read the rolled-up exact state the same way the query path does. + let mut read = finalize.clone(); + read.id = PostAsapNodeId(root.id.0 + 1); + read.output_schema = root.output_schema.clone(); + read.output_schema.fields.last_mut().unwrap().dtype = + SummaryFamilyType::Plain(planner_types::pre_asap::DataType::Float64); + edge.producer = root.id; + edge.consumer = read.id; + edge.intermediate_schema = root.output_schema.clone(); + edge.data_state = root.output_state.clone(); + dag.root = read.id; + dag.nodes.push(read); + dag.edges.push(edge); + assert_eq!( + run(&dag, SAMPLES, 60_000).unwrap(), + reference(&[("api", 4.), ("db", 1.)]) + ); +} + // Comparisons need filter/bool semantics that `Binary` does not carry, so // they fail at compile time instead of emitting 0/1 values. #[test] From 3e2609a14cfba141cc0f9aa7d299ec18e0a3129b Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 02:45:59 +0000 Subject: [PATCH 5/7] docs: record covered and remaining physical compile gaps Co-Authored-By: Claude Opus 5.5 --- .../develop_docs/physical-compile-coverage.md | 27 +++++++++++++++++-- 1 file changed, 25 insertions(+), 2 deletions(-) diff --git a/docs/develop_docs/physical-compile-coverage.md b/docs/develop_docs/physical-compile-coverage.md index 4dc309ae..5db0edce 100644 --- a/docs/develop_docs/physical-compile-coverage.md +++ b/docs/develop_docs/physical-compile-coverage.md @@ -63,7 +63,30 @@ Status values: Totals at #475: 11 Supported, 4 Partial, 14 Missing, 2 Backend. +## Covered after this change + +| Row | Change | +|---|---| +| 4, 8, 13 | Query-time `Binary` folds a scalar-literal operand into a projection over grouped value rows. | +| 5 | Query-time `Binary` over grouped value rows performs an inner equi-join on equal label columns, then applies the operator. Per-series rows remain Partial. | +| 15 | Count finalization converts exactly to the declared Float64 value. | +| 22, 23 | `ReadPopulation` Sum/Count/Average/Quantile compile to grouped aggregation. `Reduction::Quantile` implements PromQL interpolation. | + +Totals after this change: 17 Supported, 4 Partial, 8 Missing, 2 Backend. + ## Remaining -Every Missing or Partial row above, in order of backend usage. Rows 1 and 28 -are removed from the backend only after rows 4–13 are covered. +In order of backend usage: + +1. Row 1 and rows 9–12: lower PromQL-shaped `Fallback{QueryExpr}` subtrees + (range functions over matrices, `scalar()`, `histogram_quantile`, `sort`, + subquery grids). After that, rows 28 and 30 can be deleted from the backend. +2. Row 7: comparison filters and `bool` comparisons. This needs `return_bool` + in the `Binary` payload. `compile` currently rejects comparisons. +3. Row 5 for per-series rows: matching needs a metric-name-free series + identity, not the full `$promql_series_identity`. +4. Rows 25 and 27: constant weights and `EntityIdentity` items for precompute + `SummaryAgg`. +5. Row 16: a label-map sketch-state readout, the counterpart of + `compile_exact_readout`, and MetricsQL `__name__` retention rules. +6. Row 20: summary join, subtract, and delete. From 644cd95802f6ca1c552f3084f00cd83587dde610 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 02:50:21 +0000 Subject: [PATCH 6/7] fix(physical): empty global population readouts and NaN quantile order A global population aggregate with no members emitted one row; PromQL returns an empty vector. Quantile now orders NaN samples first, as Prometheus does. Found in independent review. Co-Authored-By: Claude Opus 5.5 --- .../src/operators/aggregate/mod.rs | 11 ++++-- .../src/physical_planner/mod.rs | 11 ++++-- .../src/physical_planner/row_values.rs | 34 +++++++++++++++++-- .../tests/deployment_computation.rs | 14 +++++++- 4 files changed, 62 insertions(+), 8 deletions(-) diff --git a/crates/asap-physical-operators/src/operators/aggregate/mod.rs b/crates/asap-physical-operators/src/operators/aggregate/mod.rs index 0df3c157..37b4e96d 100644 --- a/crates/asap-physical-operators/src/operators/aggregate/mod.rs +++ b/crates/asap-physical-operators/src/operators/aggregate/mod.rs @@ -209,7 +209,8 @@ async fn reduce( Ok(output) } -// Matches Prometheus `quantile`: NaN for no values, ±Inf outside [0, 1]. +// Matches Prometheus `quantile`: NaN for no values, ±Inf outside [0, 1], +// and NaN samples ordered first. fn quantile(q: f64, mut values: Vec) -> f64 { if values.is_empty() { return f64::NAN; @@ -220,7 +221,12 @@ fn quantile(q: f64, mut values: Vec) -> f64 { if q > 1. { return f64::INFINITY; } - values.sort_by(f64::total_cmp); + values.sort_by(|a, b| match (a.is_nan(), b.is_nan()) { + (true, true) => std::cmp::Ordering::Equal, + (true, false) => std::cmp::Ordering::Less, + (false, true) => std::cmp::Ordering::Greater, + _ => a.total_cmp(b), + }); let rank = q * (values.len() - 1) as f64; let low = rank.floor() as usize; let high = (low + 1).min(values.len() - 1); @@ -344,5 +350,6 @@ mod tests { assert_eq!(super::quantile(0.75, vec![4., 1., 2., 3.]), 3.25); assert_eq!(super::quantile(-0.1, vec![1.]), f64::NEG_INFINITY); assert_eq!(super::quantile(1.1, vec![1.]), f64::INFINITY); + assert_eq!(super::quantile(1., vec![2., f64::NAN, 1.]), 2.); } } diff --git a/crates/asap-physical-operators/src/physical_planner/mod.rs b/crates/asap-physical-operators/src/physical_planner/mod.rs index afc1ff1f..10fae689 100644 --- a/crates/asap-physical-operators/src/physical_planner/mod.rs +++ b/crates/asap-physical-operators/src/physical_planner/mod.rs @@ -292,9 +292,16 @@ fn compile_internal( } let input = schemas[0].clone(); let PopulationReadout::TopK { k } = readout else { - let aggregate = + let mut chain = row_values::population_aggregate(&input, &spec.grouping, readout)?; - graph.add(id, inputs, aggregate.with_output_schema(output)?)?; + let last = chain.pop().expect("nonempty chain"); + let mut inputs = inputs; + for operator in chain { + graph.add(auxiliary, inputs, operator)?; + inputs = vec![auxiliary]; + auxiliary -= 1; + } + graph.add(id, inputs, last.with_output_schema(output)?)?; continue; }; let groups = spec diff --git a/crates/asap-physical-operators/src/physical_planner/row_values.rs b/crates/asap-physical-operators/src/physical_planner/row_values.rs index 10c93c63..737c8065 100644 --- a/crates/asap-physical-operators/src/physical_planner/row_values.rs +++ b/crates/asap-physical-operators/src/physical_planner/row_values.rs @@ -170,12 +170,12 @@ fn binary(operator: &BinaryOperator, left: Expression, right: Expression) -> Exp } } -/// Aggregate readouts of a maintained current-series population. +/// Aggregate readouts of a maintained current-series population, as a chain. pub(super) fn population_aggregate( input: &Schema, grouping: &[String], readout: &PopulationReadout, -) -> Result { +) -> Result, Error> { let groups = grouping .iter() .map(|name| named_column(input, &ColumnRef::Named(name.clone()))) @@ -195,5 +195,33 @@ pub(super) fn population_aggregate( )) } }; - Operator::aggregate(input.clone(), groups, vec![("value".into(), reduction)]) + if !groups.is_empty() { + return Ok(vec![Operator::aggregate( + input.clone(), + groups, + vec![("value".into(), reduction)], + )?]); + } + // A global aggregate over no members is an empty PromQL vector, not one row. + let aggregate = Operator::aggregate( + input.clone(), + vec![], + vec![ + ("value".into(), reduction), + ("members".into(), Reduction::Count), + ], + )?; + let zero = Expression::Literal { + value: crate::values::Value::Int64(0), + dtype: DataType::Int64, + }; + let filter = Operator::filter( + aggregate.schema(), + Expression::Less(Box::new(zero), Box::new(Expression::Column(1))), + )?; + let project = Operator::project( + filter.schema(), + vec![("value".into(), Expression::Column(0))], + )?; + Ok(vec![aggregate, filter, project]) } diff --git a/crates/asap-physical-operators/tests/deployment_computation.rs b/crates/asap-physical-operators/tests/deployment_computation.rs index f08cfcae..b026f47f 100644 --- a/crates/asap-physical-operators/tests/deployment_computation.rs +++ b/crates/asap-physical-operators/tests/deployment_computation.rs @@ -220,6 +220,18 @@ fn population_aggregates_match_current_series_reference() { } } +// A global readout of an empty population is an empty vector, as in PromQL. +#[test] +fn global_population_aggregate_of_no_members_is_empty() { + // Latest values are [1, 2, 5, 7] at 60s; every member has expired by 1000s. + for (query, expected) in [("sum(m)", 15.), ("count(m)", 4.), ("quantile(0.5, m)", 3.5)] { + let dag = population_dag(query); + let live = run(&dag, SAMPLES, 60_000).unwrap(); + assert_eq!(live, reference(&[("", expected)]), "{query}"); + assert!(run(&dag, SAMPLES, 1_000_000).unwrap().is_empty(), "{query}"); + } +} + // Scalar operands on either side apply to every grouped value, including negation. #[test] fn scalar_literal_arithmetic_applies_to_grouped_values() { @@ -301,7 +313,7 @@ fn exact_count_finalizes_to_declared_float_value() { edge.producer = root.id; edge.consumer = read.id; edge.intermediate_schema = root.output_schema.clone(); - edge.data_state = root.output_state.clone(); + edge.data_state = root.output_state; dag.root = read.id; dag.nodes.push(read); dag.edges.push(edge); From f9d8a3a80bf55d716c77cdaee67a7d344ca2d8d1 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Wed, 30 Sep 2026 18:09:28 +0000 Subject: [PATCH 7/7] integrate: let coverage lowering chains use per-node helper indices #479 (d41f201) allows several helper operators per Planner node, so the coverage lowerings added here advance `auxiliary`; make it mutable. Taken from integration commit da77b78. Co-Authored-By: Claude Opus 5.5 --- crates/asap-physical-operators/src/physical_planner/mod.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/crates/asap-physical-operators/src/physical_planner/mod.rs b/crates/asap-physical-operators/src/physical_planner/mod.rs index 10fae689..c7967288 100644 --- a/crates/asap-physical-operators/src/physical_planner/mod.rs +++ b/crates/asap-physical-operators/src/physical_planner/mod.rs @@ -200,7 +200,7 @@ fn compile_internal( let mut graph = CompiledPhysicalDag::new(roots.to_vec()); for id in ordered { let node = nodes[&id]; - let auxiliary = helper_id(id, 0); + let mut auxiliary = helper_id(id, 0); let output = Arc::new(node.output_schema.clone()); crate::values::validate_schema(&output)?; if let Some(source) = sources.remove(&id) {