Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions .github/workflows/rust.yml
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ on:
- 'Cargo.toml'
- 'Cargo.lock'
- '.github/workflows/rust.yml'
# The types tests validate the viewer's exported-kind contract.
# The devtools viewer_contract test validates the viewer's operator-kind contract.
- 'tools/dag-viewer/**'
pull_request:
types: [opened, synchronize, reopened, ready_for_review]
Expand All @@ -17,7 +17,7 @@ on:
- 'Cargo.toml'
- 'Cargo.lock'
- '.github/workflows/rust.yml'
# The types tests validate the viewer's exported-kind contract.
# The devtools viewer_contract test validates the viewer's operator-kind contract.
- 'tools/dag-viewer/**'
workflow_dispatch:

Expand Down
10 changes: 4 additions & 6 deletions crates/devtools/tests/viewer_contract.rs
Original file line number Diff line number Diff line change
@@ -1,9 +1,7 @@
//! `tools/dag-viewer` ↔ `asap_types::dag_export` contract: the viewer's
//! `KIND_CATEGORY_JSON` must categorize exactly the `kind` strings
//! [`asap_types::dag_export::export`] can emit — `Operator::kind_name()` of
//! every `NonASAPOp` and `ASAPOp` variant — no more (a stale kind the IR no
//! longer has) and no less (an exported kind the viewer would render
//! uncategorized).
//! `tools/dag-viewer` ↔ IR contract: the viewer's `KIND_CATEGORY_JSON` must
//! categorize exactly the `Operator::kind_name()` strings of every
//! `NonASAPOp` and `ASAPOp` variant — no more (a stale kind the IR no longer
//! has) and no less (a kind the viewer would render uncategorized).

use std::collections::{BTreeMap, BTreeSet};

Expand Down
49 changes: 5 additions & 44 deletions crates/integration-tests/tests/exact_composition.rs
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
//! Issue #171 — composing exact operators with summary plans across
//! explicit update/evaluation boundaries, end to end through
//! `search_workload_with` → `candidate_selection::global_selection` →
//! `GlobalSelection::assemble_selected_dag` → `dag_export`.
//! `GlobalSelection::assemble_selected_dag`.
//!
//! Covers the issue's integration matrix: both nesting directions, grouped
//! fine-to-coarse and identity folds, one inner summary shared by several
Expand All @@ -13,7 +13,7 @@
use std::rc::Rc;

use asap_integration_tests::fixtures::lower_promql;
use asap_integration_tests::post_asap::{maintained, post_asap_dag, timed};
use asap_integration_tests::post_asap::{maintained, timed};
use asap_logical_optimizer::pass1::exact_composition::ExactOperation;
use asap_logical_optimizer::pass1::replacement::{
default_strategies, search_workload_with, ASAPStrategies, Replacement, ReplacementProvenance,
Expand All @@ -25,11 +25,9 @@ use asap_plan_selection::cost::cost_model::{
CostProvenance, CostUnit, ExactCompositionCostInputs, ExactCompositionCostRequest,
};
use asap_plan_selection::{CostModel, DefaultCostModel, EvaluationRate};
use asap_types::dag_export;
use asap_types::ir::data_state;
use asap_types::ir::operator::{default_quantile, AggIntent};
use asap_types::ir::operator::{Reduction, Source};
use asap_types::ir::physical_export::PhysicalASAPOperatorPayload;
use asap_types::ir::operator::agg_intent::{default_quantile, AggIntent};
use asap_types::ir::operator::operator_properties::{Reduction, Source};
use asap_types::ir::properties::timing::data_state;
use asap_types::ir::properties::{ExecutionDataState, ExecutionTiming};
use asap_types::ir::schema::{DataType, Field, Schema};
use asap_types::ir::schema::{ExactKind, FieldDataType, SketchAlgorithm, SummaryUpdate};
Expand Down Expand Up @@ -703,43 +701,6 @@ fn missing_cost_statistics_preserve_the_conservative_retain_exact() {
.any(|e| e.kind == ExplanationKind::ExactComposition));
}

// ── DAG export: explicit stage, schema, provenance ───────────────────────

#[test]
fn dag_export_carries_explicit_stage_and_plain_schema_for_a_composed_plan() {
let root = agg(vec![0], AggIntent::Max { col: None }, fine_quantile());
let space = plan(vec![("q", root)]);
let root = &space.roots[0].1;
let composed = global_selection(&space, &StatsModel)
.assemble_selected_dag(root)
.unwrap()
.unwrap();
let dag = dag_export::export_summary(&composed);
let node = &dag.nodes[dag.root as usize];
assert_eq!(node.kind, "aggregate");
assert!(node.detail["measures"].is_array());
// Timing is explicit in the wire-6 DAG: the root is a relational
// aggregate placed at query time.
let wire = post_asap_dag(&composed);
let wire_root = wire.nodes.iter().find(|n| n.id == wire.roots[0]).unwrap();
assert!(matches!(
wire_root.payload,
PhysicalASAPOperatorPayload::NonASAP(NonASAPOp::Aggregate { .. })
));
assert_eq!(wire_root.output_state.timing, ExecutionTiming::QueryTime);

// Pre-ASAP export of the same target still describes the same columns.
let pre = dag_export::export(root);
let pre_root = &pre.nodes[pre.root as usize];
let pre_cols: Vec<String> = pre_root.schema.as_ref().unwrap()["fields"]
.as_array()
.unwrap()
.iter()
.map(|c| c["name"].as_str().unwrap().to_string())
.collect();
assert_eq!(pre_cols, names(&composed));
}

/// The PromQL front end produces the exact issue shape and it composes.
#[test]
fn promql_max_by_zone_over_quantile_over_time_composes() {
Expand Down
41 changes: 7 additions & 34 deletions crates/logical-optimizer/src/pass1/explanation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -232,12 +232,10 @@ pub enum ExplanationKind {
/// [`crate::pass1::replacement::ReplacementSubDAG::rationale`]).
///
/// `node_hash` is [`structural_hash`](asap_types::ir::cse::structural_hash)
/// of the `TargetSubDAG`'s own `target` sub-DAG — the same function, on the
/// same `Rc<OperatorNode>` shape, that [`asap_types::dag_export::DAGNode::hash`]
/// is computed with. A downstream consumer that independently exported the
/// same node (e.g. via `asap_types::dag_export::export`) can match
/// this explanation to a `DAGNode` by first comparing hashes and then
/// confirming structural equality with [`ReplacementExplanation::target`].
/// of the `TargetSubDAG`'s own `target` sub-DAG. A downstream consumer that
/// hashed the same node can match this explanation by first comparing hashes
/// and then confirming structural equality with
/// [`ReplacementExplanation::target`].
#[derive(Debug, Clone, PartialEq)]
pub struct ReplacementExplanation {
pub kind: ExplanationKind,
Expand Down Expand Up @@ -300,10 +298,9 @@ fn findings_from_candidate_logical_asap_dags(
space: &CandidateLogicalASAPDAGs<String>,
) -> Vec<ReplacementExplanation> {
let locations = collect_locations(&space.roots);
// One cache for the whole pass, mirroring `dag_export::export`'s own
// `HashCache` reuse — this is a bottom-up pass over every discovered
// group, so amortizing the cache across groups (rather than resetting it
// per group) is real, not just a micro-optimization.
// One cache for the whole pass — this is a bottom-up pass over every
// discovered group, so amortizing the cache across groups (rather than
// resetting it per group) is real, not just a micro-optimization.
let mut hash_cache = HashCache::new();
let mut findings = Vec::new();
for group in space.target_subdag_candidates() {
Expand Down Expand Up @@ -581,30 +578,6 @@ mod tests {
assert!(sketch[0].reason.to_lowercase().contains("kll"));
}

/// `node_hash` must be the literal `structural_hash` a downstream
/// consumer would compute over the *same* `OperatorNode` sub-DAG via
/// `asap_types::dag_export::export` — the whole point of carrying it is
/// that two independent exports of the same DAG agree, with no
/// string-matching against `location` required.
#[test]
fn node_hash_matches_dag_export_hash_for_the_same_sub_dag() {
let q = agg(vec![2], default_quantile(0.99), metric_scan(&["job"]));
let dag = asap_types::dag_export::export(&q);
let expected_hash = dag.nodes[dag.root as usize].hash;

let findings = explain_replacements(vec![("dashboard_p99", q)]);
let sketch = findings
.iter()
.find(|f| f.kind == ExplanationKind::SketchApproximation)
.expect("expected a sketch finding");
assert_eq!(
Some(sketch.node_hash),
expected_hash,
"ReplacementExplanation::node_hash must match dag_export's DAGNode::hash \
for the same OperatorNode sub_dag"
);
}

#[test]
fn exact_quantile_is_not_a_sketch_applicability_finding() {
let q = agg(
Expand Down
3 changes: 1 addition & 2 deletions crates/logical-optimizer/src/pass1/replacement.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8192,8 +8192,7 @@ mod tests {
let Replacement::SubDAG(node) = &candidate.replacement else {
unreachable!()
};
let exported = asap_types::dag_export::export_summary(node);
assert!(exported.nodes[exported.root as usize]
assert!(node
.guarantee
.as_ref()
.is_some_and(ResultGuarantee::has_unknown));
Expand Down
6 changes: 2 additions & 4 deletions crates/types/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -10,10 +10,8 @@ edition = "2021"
[dependencies]
# "rc" — OperatorNode child fields are Rc<OperatorNode> (issue #212, #222:
# shared sub-expressions), and Rc<T>'s Serialize/Deserialize impls live behind
# this feature flag. dag_export.rs / DAGNode already flatten the DAG to a
# node+edge list for JSON export, so this does not change wire format — a
# shared Rc still (de)serializes as an ordinary inline value, once per
# reference, exactly like the old Box.
# this feature flag. A shared Rc still (de)serializes as an ordinary inline
# value, once per reference, exactly like the old Box.
serde = { version = "1", features = ["derive", "rc"] }
serde_json = "1"
thiserror = "2"
3 changes: 1 addition & 2 deletions crates/types/src/cost.rs
Original file line number Diff line number Diff line change
Expand Up @@ -420,8 +420,7 @@ where
/// Whole-selected-workload cost/benefit — one query's (or one workload
/// batch's) aggregate baseline, selected, and benefit, built from
/// [`sum_workload_costs`] over that scope's own per-decision node
/// annotations. See [`crate::dag_export::NamedDAG::workload_cost`] /
/// [`crate::dag_export::WorkloadDAG::workload_cost`].
/// annotations.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct WorkloadCostSummary {
pub baseline_cost: CostAnnotation,
Expand Down
Loading