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
165 changes: 161 additions & 4 deletions control_plane/src/physical/compiler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -372,6 +372,8 @@ pub struct ProducerContract {

#[derive(Debug, Error, PartialEq, Eq)]
pub enum PrecomputePlanError {
#[error("raw counter state {0} cannot preserve independent series resets/timestamps")]
UnsupportedRawCounter(u64),
#[error("PrecomputePlan envelope does not match BackendPlan identity/lifecycle")]
PlanIdentityMismatch,
#[error("unsupported precompute ingest protocol/endpoint/identity contract")]
Expand Down Expand Up @@ -520,6 +522,20 @@ impl PrecomputePlan {
}
let mut materializations = BTreeSet::new();
for materialization in &self.materializations {
if self.ingest.protocol == IngestProtocol::PrometheusRemoteWriteV1
&& (matches!(
materialization.aggregation_type,
asap_types::AggregationType::Increase
) || (materialization.aggregation_type
== asap_types::AggregationType::SingleSubpopulation
&& materialization
.aggregation_sub_type
.eq_ignore_ascii_case("increase")))
{
return Err(PrecomputePlanError::UnsupportedRawCounter(
materialization.policy_fp_u64(),
));
}
if !materializations.insert(materialization.policy_fingerprint()) {
return Err(PrecomputePlanError::DuplicateMaterialization(
materialization.policy_fp_u64(),
Expand Down Expand Up @@ -1666,6 +1682,7 @@ impl BackendLocalPlanningSnapshot {
});
}
select_workload_roots(&mut queries, canonical_roots, &topk_evidence_by_id)?;
preserve_native_raw_counters(&mut queries)?;
Ok((
PlanningRequest {
queries,
Expand All @@ -1677,6 +1694,45 @@ impl BackendLocalPlanningSnapshot {
}
}

/// Preserve native execution for counters until raw producers retain series identity.
fn preserve_native_raw_counters(queries: &mut [PlanningQuery]) -> Result<(), CompileError> {
for query in queries {
let selected = collect_selected_materializations(&query.post_asap).map_err(|reason| {
CompileError::Query {
query_id: query.query_id.clone(),
reason,
}
})?;
if selected.iter().any(|state| {
matches!(
state.family,
SummaryFamilyType::ExactAggregate(
planner_types::post_asap::ExactKind::Increase
| planner_types::post_asap::ExactKind::Rate,
_
)
)
}) {
let parsed = crate::query_parser::parse_query_expr_canonical(
&query.query_string,
query.accuracy.clone(),
)
.map_err(|error| CompileError::Query {
query_id: query.query_id.clone(),
reason: error.to_string(),
})?;
query.post_asap =
crate::planner_selection::keep_pre_asap(&parsed).map_err(|error| {
CompileError::Query {
query_id: query.query_id.clone(),
reason: error.to_string(),
}
})?;
}
}
Ok(())
}

impl PhysicalCompiler {
pub fn compile(
&self,
Expand All @@ -1698,6 +1754,10 @@ impl PhysicalCompiler {
});
}

if environment.target == PhysicalDeploymentTarget::BackendLocalRemoteWrite {
preserve_native_raw_counters(&mut request.queries)?;
}

let roots = request
.queries
.iter()
Expand Down Expand Up @@ -2738,6 +2798,92 @@ fn stable_workload_plan_id(
mod tests {
use super::*;

// A pooled raw counter cannot distinguish same-timestamp series or independent resets.
#[test]
fn backend_local_rejects_pooled_counter_materialization() {
let request = request("counter", "sum(rate(m[1m]))");
let mut environment = environment(10_000);
environment.target = PhysicalDeploymentTarget::BackendLocalRemoteWrite;
environment.collector_ids.clear();
let plan = PhysicalCompiler.compile(request, environment).unwrap();
assert!(
plan.precompute_plan.materializations.is_empty(),
"unsafe pooled raw counter was installed"
);
assert!(plan.query_plan.entries.values().all(|entry| matches!(
entry.nodes[&entry.root],
crate::query_plan::QueryPlanNode::ExactFallback { .. }
)));
}

// Old precompiled raw counter artifacts are rejected before installation too.
#[test]
fn raw_counter_artifact_is_rejected_while_envelope_counter_remains_valid() {
let plan = PhysicalCompiler
.compile(request("counter", "rate(m[1m])"), environment(10_000))
.unwrap();
plan.precompute_plan.validate().unwrap();
let mut raw = plan.precompute_plan;
raw.ingest.protocol = IngestProtocol::PrometheusRemoteWriteV1;
raw.ingest.endpoint_path = "/api/v1/write".into();
raw.ingest.timestamp_unit = TimestampUnit::UnixMilliseconds;
raw.ingest.require_plan_identity = false;
raw.ingest.require_materialization_identity = false;
raw.ingest.require_registered_producer = false;
raw.producers.clear();
assert!(matches!(
raw.validate(),
Err(PrecomputePlanError::UnsupportedRawCounter(_))
));
}

// Counter fallback does not disable an independent safe summary in the same workload.
#[test]
fn counter_fallback_preserves_other_workload_summaries() {
let mut workload = request("counter", "sum(rate(m[1m]))");
workload
.queries
.extend(request("gauge", "sum_over_time(g[1m])").queries);
let mut environment = environment(10_000);
environment.target = PhysicalDeploymentTarget::BackendLocalRemoteWrite;
environment.collector_ids.clear();
let plan = PhysicalCompiler.compile(workload, environment).unwrap();
assert_eq!(plan.precompute_plan.materializations.len(), 1);
assert!(plan
.query_plan
.lookup("sum(rate(m[1m]))")
.unwrap()
.materialization_bindings()
.is_empty());
assert_eq!(
plan.query_plan
.lookup("sum_over_time(g[1m])")
.unwrap()
.materialization_bindings()
.len(),
1
);
}

// Capability normalization precedes candidate enumeration, avoiding duplicate exact quotes.
#[test]
fn counter_only_snapshot_has_one_exact_cost_alternative() {
let mut snapshot: BackendLocalPlanningSnapshot = serde_json::from_str(include_str!(
"../../../docs/examples/asapquery-planning-snapshot.json"
))
.unwrap();
let entry = &mut snapshot.query_workload.repeating_queries.as_mut().unwrap()[0];
entry.query = Query("rate(m[1m])".into());
entry.requirements.accuracy = AccuracyRequirement::Explicit(AccuracyTarget::Exact);
let (request, _) = snapshot.planning_request().unwrap();
assert_eq!(
super::super::workload_cost::with_exact_alternative(request)
.unwrap()
.len(),
1
);
}

// Count and value rankings must configure different state update contracts.
#[test]
fn temporal_topk_binds_planner_update_weight() {
Expand Down Expand Up @@ -3407,9 +3553,20 @@ mod tests {
let mut deployment = environment(10_000);
deployment.target = PhysicalDeploymentTarget::BackendLocalRemoteWrite;
deployment.collector_ids.clear();
let plan = PhysicalCompiler
.compile(request(query_id, promql), deployment)
.unwrap_or_else(|error| panic!("{promql} must compile: {error}"));
let compiled = PhysicalCompiler.compile(request(query_id, promql), deployment);
if matches!(
expected_readout,
crate::query_plan::ExactReadout::Rate | crate::query_plan::ExactReadout::Increase
) {
let plan = compiled.unwrap();
assert!(plan.precompute_plan.materializations.is_empty());
assert!(plan.query_plan.entries.values().all(|entry| matches!(
entry.nodes[&entry.root],
crate::query_plan::QueryPlanNode::ExactFallback { .. }
)));
continue;
}
let plan = compiled.unwrap_or_else(|error| panic!("{promql} must compile: {error}"));
assert_eq!(plan.backend_plan.materializations.len(), 1, "{promql}");
assert_eq!(plan.query_plan.entries.len(), 1, "{promql}");
assert!(plan.collector_plans.is_empty(), "{promql}");
Expand Down Expand Up @@ -3462,7 +3619,7 @@ mod tests {
assert!(plan.collector_plans.is_empty());
assert!(plan.transmission_plan.rules.is_empty());
assert_eq!(plan.query_plan.entries.len(), 6);
assert_eq!(plan.precompute_plan.materializations.len(), 5);
assert_eq!(plan.precompute_plan.materializations.len(), 4);
for query in [
"rate(asap_demo_counter_total[5s])",
"increase(asap_demo_counter_total[5s])",
Expand Down
48 changes: 48 additions & 0 deletions data_plane/src/precompute_engine/worker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3259,6 +3259,54 @@ aggregations:
// sketch (`sketch_panes`) paths.
// -----------------------------------------------------------------------

// This single-series updater cannot implement a grouped sum of counter increases.
// The physical compiler rejects raw counter producers until series state is preserved.
#[test]
fn pooled_counter_samples_lose_independent_same_timestamp_reset() {
use crate::precompute_engine::operators::IncreaseAccumulator;
let config = make_agg_config(
1,
"requests_total",
AggregationType::SingleSubpopulation,
"Increase",
10,
0,
vec![],
);
let sink = Arc::new(CapturingOutputSink::new());
let mut worker = make_worker(
HashMap::from([(1, config)]),
sink.clone(),
false,
0,
LateDataPolicy::Drop,
);
worker
.process_group_samples(
1,
PolicyFingerprint(1),
"",
vec![
("requests_total{instance=\"a\"}".into(), 1000, 100.0),
("requests_total{instance=\"b\"}".into(), 1000, 50.0),
("requests_total{instance=\"a\"}".into(), 2000, 110.0),
("requests_total{instance=\"b\"}".into(), 2000, 5.0),
],
)
.unwrap();
worker.force_close_all().unwrap();
let captured = sink.drain();
let accumulator = captured[0]
.1
.as_any()
.downcast_ref::<IncreaseAccumulator>()
.unwrap();
assert_eq!(accumulator.total_increase, 10.0);
let independent_increases = (110.0 - 100.0) + 5.0;
assert_eq!(independent_increases, 15.0);
assert_ne!(accumulator.total_increase, independent_increases);
}

#[test]
fn shutdown_force_close_emits_trailing_sample_window() {
// 10s tumbling window; make_worker uses grace=0, isolating the
Expand Down
Loading
Loading