Skip to content
Merged
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
6 changes: 3 additions & 3 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

6 changes: 3 additions & 3 deletions control_plane/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -93,16 +93,16 @@ asap_types.workspace = true
# scaffolding, unaware that `data_plane`'s `summary_executor.rs` in *this*
# repo is a real one. Vendored locally instead of chased upstream -- see
# `data_plane/src/query_engines/asap_query_engine/summary_exec.rs`.
planner-types = { package = "asap-types", git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "cb70086b4c4a7ba89baf2516be81d0b192137a3a" }
asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "cb70086b4c4a7ba89baf2516be81d0b192137a3a" }
planner-types = { package = "asap-types", git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "3afcba68f4e8397fb81e2be988f47120f63f7a39" }
asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "3afcba68f4e8397fb81e2be988f47120f63f7a39" }

# L1 adoption (design-target-architecture.md Part B): the PromQL front
# end itself, replacing control_plane's own query_parser/promql.rs.
# Pinned via `rev`, not a floating branch reference. Same rev as
# `planner-types`/`asap-aware-mapping` above -- these three MUST move
# together (two revs of the same upstream repo's types in one workspace
# resolve to distinct Rust types that won't unify).
asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "cb70086b4c4a7ba89baf2516be81d0b192137a3a" }
asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "3afcba68f4e8397fb81e2be988f47120f63f7a39" }

[dev-dependencies]
tokio = { version = "1", features = ["full", "test-util"] }
Expand Down
6 changes: 3 additions & 3 deletions control_plane/src/asap_tier_analysis.rs
Original file line number Diff line number Diff line change
Expand Up @@ -356,11 +356,11 @@ pub(crate) fn collect_agg_intents(expr: &QueryExpr, out: &mut Vec<AggIntent>) {
// at construction time (`intent_algebra::lower`).
QueryExpr::Filter { child, .. }
| QueryExpr::Project { child, .. }
| QueryExpr::Distinct { child, .. }
| QueryExpr::Dedup { child, .. }
| QueryExpr::Sort { child, .. }
| QueryExpr::Limit { child, .. }
| QueryExpr::Subquery { child, .. } => collect_agg_intents(child, out),
QueryExpr::Merge { children } => {
| QueryExpr::PromqlSubquery { child, .. } => collect_agg_intents(child, out),
QueryExpr::Concat { children } => {
for c in children {
collect_agg_intents(c, out);
}
Expand Down
15 changes: 8 additions & 7 deletions control_plane/src/asap_tier_implement.rs
Original file line number Diff line number Diff line change
Expand Up @@ -73,7 +73,7 @@

use std::rc::Rc;

use asap_aware_mapping::{implement_tree_with, DefaultCostModel, ImplementError};
use asap_aware_mapping::DefaultCostModel;
use planner_types::post_asap::SummaryNode;

use crate::intent_algebra::query_expr::QueryExpr;
Expand Down Expand Up @@ -130,11 +130,11 @@ fn collect_aggregate_roots<'a>(expr: &'a QueryExpr, out: &mut Vec<&'a QueryExpr>
// at construction time (`intent_algebra::lower`).
QueryExpr::Filter { child, .. }
| QueryExpr::Project { child, .. }
| QueryExpr::Distinct { child, .. }
| QueryExpr::Dedup { child, .. }
| QueryExpr::Sort { child, .. }
| QueryExpr::Limit { child, .. }
| QueryExpr::Subquery { child, .. } => collect_aggregate_roots(child, out),
QueryExpr::Merge { children } => {
| QueryExpr::PromqlSubquery { child, .. } => collect_aggregate_roots(child, out),
QueryExpr::Concat { children } => {
for c in children {
collect_aggregate_roots(c, out);
}
Expand Down Expand Up @@ -165,7 +165,7 @@ pub enum ImplementPromqlError {
UnparseableMetricsql(String),
/// L3→L4 implementation failed for a found `Aggregate` root (schema
/// derivation error — see `asap_aware_mapping::bind::ImplementError`).
Implement(ImplementError),
Implement(crate::planner_selection::SelectionError),
}

/// Parse `metricsql`, find every independently-realizable `Aggregate`
Expand Down Expand Up @@ -194,7 +194,8 @@ pub fn implement_promql_for_asap_tier(
roots
.into_iter()
.map(|root| {
implement_tree_with(root, &DefaultCostModel).map_err(ImplementPromqlError::Implement)
crate::planner_selection::select_summary(root, &DefaultCostModel)
.map_err(ImplementPromqlError::Implement)
})
.collect()
}
Expand Down Expand Up @@ -268,7 +269,7 @@ mod tests {
.expect("parses and implements");
assert_eq!(roots.len(), 1);
assert!(
matches!(roots[0].expr, SummaryExpr::Logical(_)),
matches!(roots[0].expr, SummaryExpr::KeepPreAsap(_)),
"Avg has no ASAP-tier realization yet on either path: {:?}",
roots[0].expr,
);
Expand Down
4 changes: 2 additions & 2 deletions control_plane/src/emit/backend_push.rs
Original file line number Diff line number Diff line change
Expand Up @@ -724,12 +724,12 @@ mod tests {
AggregationInput, BackendAggregation, BackendReadout,
};
use planner_types::post_asap::SketchQuery;
use planner_types::post_asap::{SketchKind, SketchParams};
use planner_types::post_asap::{SketchAlgorithm, SketchParams};
BackendStageConfig {
aggregations: vec![BackendAggregation {
aggregation_id: agg_id.to_string(),
metric_name: metric.to_string(),
sketch_kind: SketchKind::DDSketch.into(),
sketch_kind: SketchAlgorithm::DDSketch.into(),
sketch_params: SketchParams::DDSketch { alpha: 0.01 }.into(),
grouping: vec![],
item_label: None,
Expand Down
57 changes: 27 additions & 30 deletions control_plane/src/emit/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,7 @@ use crate::sketch_algebra::physical_expr::L4Plan;
use crate::sketch_algebra::PhysicalExpr;
use crate::store::WorkloadStore;
use anyhow::Result;
use planner_types::post_asap::{SketchKind, SummaryExpr, SummaryNode};
use planner_types::post_asap::{SketchAlgorithm, SummaryExpr, SummaryNode};
use std::rc::Rc;

/// Phase ε.1.5 — which edge runtime an agent identifies as.
Expand Down Expand Up @@ -287,15 +287,15 @@ fn apply_cold_format_from_env(edge_cfg: &mut EdgeStageConfig) {
/// (`Logical`-only, unresolved `Ref`, raw Mode-3 archive). These map
/// onto the raw-passthrough default pipeline in the routing emitter,
/// which is correct.
pub fn extract_root_sketch_kind(expr: &PhysicalExpr) -> Option<SketchKind> {
pub fn extract_root_sketch_kind(expr: &PhysicalExpr) -> Option<SketchAlgorithm> {
match expr {
PhysicalExpr::Committed(plan) => extract_from_plan(plan),
PhysicalExpr::RawAtEdgeSketchAtBackend { family, .. } => Some(family.clone()),
PhysicalExpr::RawAtEdgePrometheusArchive { .. } => None,
}
}

fn extract_from_plan(plan: &L4Plan) -> Option<SketchKind> {
fn extract_from_plan(plan: &L4Plan) -> Option<SketchAlgorithm> {
match plan {
L4Plan::Summary(node) => extract_from_node(node),
L4Plan::LetBinding { expr, child, .. } => {
Expand All @@ -305,7 +305,7 @@ fn extract_from_plan(plan: &L4Plan) -> Option<SketchKind> {
}
}

fn extract_from_node(node: &Rc<SummaryNode>) -> Option<SketchKind> {
fn extract_from_node(node: &Rc<SummaryNode>) -> Option<SketchAlgorithm> {
match &node.expr {
// `SummaryAgg`'s `kind`/`params` collapsed into one `family:
// SummaryFamilyType` field (ASAPPlanner#218 -- see
Expand All @@ -315,7 +315,7 @@ fn extract_from_node(node: &Rc<SummaryNode>) -> Option<SketchKind> {
SummaryExpr::SummaryAgg {
family: planner_types::post_asap::SummaryFamilyType::Sketch(kind, _),
..
} => Some(kind.clone()),
} => Some(kind.algorithm().clone()),
// An exact accumulator has no sketch family beneath it (its own
// child is always a plain `Logical` leaf) — same as the old
// `ExactAgg` case.
Expand All @@ -327,7 +327,7 @@ fn extract_from_node(node: &Rc<SummaryNode>) -> Option<SketchKind> {
SummaryExpr::SummaryJoin { .. }
| SummaryExpr::SummarySubtract { .. }
| SummaryExpr::SummaryDelete { .. }
| SummaryExpr::Logical(_) => None,
| SummaryExpr::KeepPreAsap(_) => None,
}
}

Expand All @@ -342,7 +342,7 @@ fn extract_from_node(node: &Rc<SummaryNode>) -> Option<SketchKind> {
/// (`quantile_over_time` → DDSketch, `count`-distinct → HLL, `topk` →
/// CountSketch, …). We therefore collect the UNION of every workload
/// entry's committed sketch family per metric into a
/// `BTreeSet<SketchKind>` (deterministic order). The emitter routes the
/// `BTreeSet<SketchAlgorithm>` (deterministic order). The emitter routes the
/// metric to EACH family in its set and prunes pipelines/processors to
/// the union of all sets — eliminating the prior all-5 fan-out that
/// shipped sketch state through every family regardless of need.
Expand All @@ -362,8 +362,8 @@ fn extract_from_node(node: &Rc<SummaryNode>) -> Option<SketchKind> {
pub fn collect_metric_to_family(
registry: &WorkloadRegistry,
workload_store: &WorkloadStore,
) -> std::collections::HashMap<String, std::collections::BTreeSet<SketchKind>> {
let mut out: std::collections::HashMap<String, std::collections::BTreeSet<SketchKind>> =
) -> std::collections::HashMap<String, std::collections::BTreeSet<SketchAlgorithm>> {
let mut out: std::collections::HashMap<String, std::collections::BTreeSet<SketchAlgorithm>> =
std::collections::HashMap::new();
for entry in registry.entries() {
// B2 (metric, role) restructure: walk EVERY role registered for
Expand Down Expand Up @@ -773,7 +773,7 @@ mod runtime_tests {
window_secs: Some(60),
sketch_processors: vec![EdgeSketchProcessor {
processor_name: "ddsketch".to_string(),
sketch_kind: SketchKind::DDSketch,
sketch_kind: SketchAlgorithm::DDSketch,
sketch_params: SketchParams::DDSketch { alpha: 0.01 },
aggregation_id: "agg0".to_string(),
}],
Expand Down Expand Up @@ -950,7 +950,7 @@ mod runtime_tests {

#[test]
fn collect_metric_to_family_binds_all_six_contract_metrics_from_live_yaml() {
use planner_types::post_asap::SketchKind;
use planner_types::post_asap::SketchAlgorithm;

// The 6 contract metrics reproduced inline (mirrors
// deploy/configs/mvp-workload.yaml entries 1, 5, 6, 7, 8 plus the
Expand Down Expand Up @@ -1004,31 +1004,28 @@ mod runtime_tests {
// metric needs. For THIS workload every sketched metric is
// queried by exactly one capability, so each set has size 1.
use std::collections::BTreeSet;
let expected: Vec<(&str, Option<BTreeSet<SketchKind>>)> = vec![
let expected: Vec<(&str, Option<BTreeSet<SketchAlgorithm>>)> = vec![
(
"http_latency_ms",
Some(BTreeSet::from([SketchKind::DDSketch])),
Some(BTreeSet::from([SketchAlgorithm::DDSketch])),
),
("http_requests_total", None), // raw passthrough
(
"request_size_bytes",
Some(BTreeSet::from([SketchKind::Kll])),
Some(BTreeSet::from([SketchAlgorithm::Kll])),
),
(
"unique_users_per_min",
Some(BTreeSet::from([SketchKind::Hll])),
),
(
"top_endpoint_qps",
Some(BTreeSet::from([SketchKind::CountSketchWithHeap])),
Some(BTreeSet::from([SketchAlgorithm::Hll])),
),
("top_endpoint_qps", None),
// `CountMinSketch` override re-derives statistic to
// `Frequency`, `AggIntent::Extension`-shaped — now binds via
// `ControlPlaneCostModel::realize_extension` (ASAPController#150,
// see `optimizer::rules::tests::typed_binding_endpoint_request_freq_binds_cms`).
(
"endpoint_request_freq",
Some(BTreeSet::from([SketchKind::Cms])),
Some(BTreeSet::from([SketchAlgorithm::Cms])),
),
];
for (metric, want) in &expected {
Expand All @@ -1039,12 +1036,12 @@ mod runtime_tests {
full map: {map:?}",
);
}
// Routing table covers the 5 sketched metrics (only
// http_requests_total declines, as raw passthrough).
// TopK also declines until a fresh membership-margin certificate is
// supplied to the physical compiler.
assert_eq!(
map.len(),
5,
"routing table should have 5 entries (5 sketches; only raw passthrough declines), got: {map:?}"
4,
"routing table should have 4 evidence-valid sketch entries; raw passthrough and uncertified TopK decline, got: {map:?}"
);
}

Expand Down Expand Up @@ -1160,7 +1157,7 @@ mod runtime_tests {
fn collect_metric_to_family_unions_multiple_capabilities_per_metric() {
use crate::types::{AggType, QueryWorkload, SketchType, WorkloadCharacteristics};
use crate::workload::AggRole;
use planner_types::post_asap::SketchKind;
use planner_types::post_asap::SketchAlgorithm;
use std::collections::BTreeSet;
use std::time::Duration;

Expand Down Expand Up @@ -1236,9 +1233,9 @@ mod runtime_tests {
assert_eq!(
got,
BTreeSet::from([
SketchKind::DDSketch,
SketchKind::Hll,
SketchKind::Cms
SketchAlgorithm::DDSketch,
SketchAlgorithm::Hll,
SketchAlgorithm::Cms
]),
"a metric queried by 3 capabilities must accumulate 3 families (UNION, not first-wins)\nmap: {map:?}"
);
Expand Down Expand Up @@ -1306,7 +1303,7 @@ mod runtime_tests {
let _env = crate::test_support::env_lock();
use crate::physical::colored_dag::emitter::{EdgeStageConfig, ExportTarget};
use crate::physical::colored_dag::stage_id::StageId;
use planner_types::post_asap::SketchKind;
use planner_types::post_asap::SketchAlgorithm;

let yaml = r#"
- metric_name: http_requests_total_latency_ms
Expand All @@ -1332,7 +1329,7 @@ mod runtime_tests {
warm_passthrough_metrics: Vec::new(),
metric_to_family: std::collections::HashMap::from([(
"http_requests_total_latency_ms".to_string(),
std::collections::BTreeSet::from([SketchKind::DDSketch]),
std::collections::BTreeSet::from([SketchAlgorithm::DDSketch]),
)]),
metric_to_grouping_labels: std::collections::HashMap::new(),
cumulative_counter_metrics: Vec::new(),
Expand Down
26 changes: 13 additions & 13 deletions control_plane/src/emit/otap.rs
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,7 @@ use std::collections::BTreeMap;

use crate::physical::colored_dag::emitter::{EdgeSketchProcessor, EdgeStageConfig, ExportTarget};
use crate::physical::colored_dag::stage_id::StageId;
use planner_types::post_asap::{SketchKind, SketchParams};
use planner_types::post_asap::{SketchAlgorithm, SketchParams};

/// Default URL for Prometheus's native OTLP HTTP receiver.
/// Matches `super::stage_config::emit_edge_yaml`'s placeholder so the
Expand Down Expand Up @@ -310,7 +310,7 @@ fn build_asap_sketches_config(sp: &EdgeSketchProcessor, window_secs: Option<u64>
}
// Heap-bearing width/depth extraction is identical to the bare
// kind — this path never distinguished `with_heap` even before
// `SketchKind` split it into its own variant (heap_size wasn't
// `SketchAlgorithm` split it into its own variant (heap_size wasn't
// emitted here either way).
SketchParams::Cms { width, depth } | SketchParams::CmsWithHeap { width, depth, .. } => {
m.insert("rows".into(), Value::Number((*depth as u64).into()));
Expand All @@ -328,24 +328,24 @@ fn build_asap_sketches_config(sp: &EdgeSketchProcessor, window_secs: Option<u64>
SketchParams::Kmv { .. } | SketchParams::Theta { .. } => {
unreachable!(
"edge sketch processor config requested for a non-sketch or unsupported \
SketchKind; no Bind* rule in this repo produces one"
SketchAlgorithm; no Bind* rule in this repo produces one"
)
}
}
Value::Mapping(m)
}

fn sketch_kind_tag(kind: &SketchKind) -> &'static str {
fn sketch_kind_tag(kind: &SketchAlgorithm) -> &'static str {
match kind {
SketchKind::Kll => "kll",
SketchKind::DDSketch => "ddsketch",
SketchKind::Hll => "hll",
SketchKind::Cms | SketchKind::CmsWithHeap => "cms",
SketchKind::CountSketch | SketchKind::CountSketchWithHeap => "count_sketch",
SketchKind::Kmv | SketchKind::Theta => {
SketchAlgorithm::Kll => "kll",
SketchAlgorithm::DDSketch => "ddsketch",
SketchAlgorithm::Hll => "hll",
SketchAlgorithm::Cms | SketchAlgorithm::CmsWithHeap => "cms",
SketchAlgorithm::CountSketch | SketchAlgorithm::CountSketchWithHeap => "count_sketch",
SketchAlgorithm::Kmv | SketchAlgorithm::Theta => {
unreachable!(
"edge sketch processor config requested for a non-sketch or unsupported \
SketchKind; no Bind* rule in this repo produces one"
SketchAlgorithm; no Bind* rule in this repo produces one"
)
}
}
Expand All @@ -357,7 +357,7 @@ fn sketch_kind_tag(kind: &SketchKind) -> &'static str {
mod tests {
use super::*;
use crate::physical::colored_dag::emitter::{EdgeSketchProcessor, PrometheusArchiveMetric};
use planner_types::post_asap::{SketchKind, SketchParams};
use planner_types::post_asap::{SketchAlgorithm, SketchParams};

/// Minimal struct-stub used to validate the emitted DAG parses as the
/// otap-dataflow schema. We don't pull in the otap-df-config crate
Expand Down Expand Up @@ -404,7 +404,7 @@ mod tests {
window_secs: Some(60),
sketch_processors: vec![EdgeSketchProcessor {
processor_name: "ddsketch".to_string(),
sketch_kind: SketchKind::DDSketch,
sketch_kind: SketchAlgorithm::DDSketch,
sketch_params: SketchParams::DDSketch { alpha: 0.01 },
aggregation_id: "agg0".to_string(),
}],
Expand Down
Loading