diff --git a/control_plane/src/main.rs b/control_plane/src/main.rs index 8e0ac885..263bfed3 100644 --- a/control_plane/src/main.rs +++ b/control_plane/src/main.rs @@ -618,6 +618,7 @@ struct CompileAndPublishPhysicalPlanResponse { status: &'static str, generated_at_unix_ms: u64, collector_ids: Vec, + lifecycle_estimates: Vec, } /// Compile one Planner IR decision into matching Collector, Precompute, and Backend views @@ -712,6 +713,7 @@ async fn handle_compile_and_publish_physical_plan( status: "active", generated_at_unix_ms: bundle.envelope.generated_at_unix_ms, collector_ids, + lifecycle_estimates: bundle.lifecycle_estimates, }) .into_response() } @@ -759,6 +761,7 @@ fn compile_physical_plan_request( .unwrap_or_default() .as_millis() as u64; let mut queries = Vec::with_capacity(request.queries.len()); + let mut canonical_roots = Vec::with_capacity(request.queries.len()); for query in request.queries { if query.query_id.trim().is_empty() || query.metric.trim().is_empty() @@ -773,15 +776,11 @@ fn compile_physical_plan_request( Ok(expr) => expr, Err(error) => return Err((StatusCode::UNPROCESSABLE_ENTITY, error.to_string())), }; - let post_asap = match physical::compiler::select_post_asap( - &expr, - query.accuracy.clone(), - &query.lifecycle, - request.evidence.get(&query.query_id), - ) { + let post_asap = match control_plane::planner_selection::keep_pre_asap(&expr) { Ok(plan) => plan, Err(error) => return Err((StatusCode::UNPROCESSABLE_ENTITY, error.to_string())), }; + canonical_roots.push(std::rc::Rc::new(expr)); queries.push(physical::compiler::PlanningQuery { query_id: query.query_id, query_string: query.query_string, @@ -798,6 +797,12 @@ fn compile_physical_plan_request( }); } + if let Err(error) = + physical::compiler::select_workload_roots(&mut queries, canonical_roots, &request.evidence) + { + return Err((StatusCode::UNPROCESSABLE_ENTITY, error.to_string())); + } + let bundle = match physical::compiler::PhysicalCompiler.compile( physical::compiler::PlanningRequest { queries, diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index a764fc4e..2b280480 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -618,6 +618,19 @@ pub struct PhysicalPlan { pub transmission_plan: TransmissionPlan, pub backend_plan: BackendPlan, pub query_plan: QueryPlan, + /// Lifecycle component only, not a complete physical-plan comparison. + pub lifecycle_estimates: Vec, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +pub struct MaterializationLifecycleEstimate { + pub materialization: asap_types::PolicyFingerprint, + pub consumer_query_ids: Vec, + pub window_implementation_id: String, + pub horizon_seconds: f64, + pub expected_reads: f64, + pub expected_updates: f64, + pub lifecycle_cost: f64, } #[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, PartialOrd, Ord)] @@ -1619,25 +1632,7 @@ impl BackendLocalPlanningSnapshot { runtime_policy: RuntimeRulePolicy::default(), }); } - // A cohort shares the same end-to-end requirement, not an inferred - // weakest common accuracy. Different targets are searched separately. - let mut cohorts: Vec<(AccuracyTarget, Vec<(usize, Rc)>)> = Vec::new(); - for (index, root) in canonical_roots.into_iter().enumerate() { - let accuracy = &queries[index].accuracy; - if let Some((_, roots)) = cohorts.iter_mut().find(|(target, _)| target == accuracy) { - roots.push((index, root)); - } else { - cohorts.push((accuracy.clone(), vec![(index, root)])); - } - } - for (accuracy, roots) in cohorts { - let model = ControlPlaneCostModel::new(accuracy.clone()); - let selected = crate::planner_selection::select_workload(roots, accuracy, &model) - .map_err(|error| CompileError::Snapshot(error.to_string()))?; - for (index, node) in selected { - queries[index].post_asap = node; - } - } + select_workload_roots(&mut queries, canonical_roots, &HashMap::new())?; PhysicalCompiler.compile( PlanningRequest { queries, @@ -1690,6 +1685,8 @@ impl PhysicalCompiler { // the planner DAG node identity across the workload; serving never scans // BackendPlan candidates to rediscover this decision. let mut node_bindings = HashMap::::new(); + let consumers = materialization_consumers(&request.queries, environment.target)?; + let mut lifecycle_estimates = BTreeMap::new(); for query in &request.queries { let evidence = request.evidence.get(&query.query_id); @@ -1743,16 +1740,6 @@ impl PhysicalCompiler { if selected.is_empty() { continue; } - let planner_selection = select_lifecycle(query, &node, &model, &environment)?; - let window_implementation = query - .window_implementations - .iter() - .filter(|candidate| candidate.framework == planner_selection.window_framework) - .min_by(|left, right| left.cost.weighted_cost.total_cmp(&right.cost.weighted_cost)) - .ok_or_else(|| CompileError::Lifecycle { - query_id: query.query_id.clone(), - reason: "Planner selected a window framework without a retained concrete implementation".into(), - })?; match &query.source { Source::TimeSeries { .. } => {} Source::Table { .. } => { @@ -1773,26 +1760,49 @@ impl PhysicalCompiler { SummaryFamilyType::ExactAggregate(kind, _) => { format!("{kind:?}").to_ascii_lowercase() } - _ => selected.algorithm, - }; - let aggregation = BackendAggregation { - aggregation_id: aggregation_id.clone(), - metric_name: metric.clone(), - family: physical_family, - window_secs: query.window_secs, - spatial_filter: String::new(), - grouping: query.group_by.clone(), - item_label: None, - aggregation_input: match environment.target { - PhysicalDeploymentTarget::DistributedCollectors => { - AggregationInput::SketchEnvelope - } - PhysicalDeploymentTarget::BackendLocalRemoteWrite => AggregationInput::Raw, - }, + _ => selected.algorithm.clone(), }; + let aggregation = physical_aggregation( + query, + &selected, + aggregation_id.clone(), + environment.target, + ); let precompute_materialization = backend_plan::aggregation_config_for_materialization(&aggregation)?; let materialization = precompute_materialization.policy_fingerprint(); + let state_consumers = consumers[&materialization] + .iter() + .map(|index| &request.queries[*index]) + .collect::>(); + let planner_selection = select_lifecycle( + query, + &selected.node, + &model, + &environment, + &state_consumers, + )?; + let window_implementation = query.window_implementations.iter() + .filter(|candidate| candidate.framework == planner_selection.window_framework) + .min_by(|left, right| left.cost.weighted_cost.total_cmp(&right.cost.weighted_cost)) + .ok_or_else(|| CompileError::Lifecycle { + query_id: query.query_id.clone(), + reason: "Planner selected a window framework without a retained concrete implementation".into(), + })?; + lifecycle_estimates + .entry(materialization) + .or_insert_with(|| MaterializationLifecycleEstimate { + materialization, + consumer_query_ids: state_consumers + .iter() + .map(|query| query.query_id.clone()) + .collect(), + window_implementation_id: window_implementation.implementation_id.clone(), + horizon_seconds: query.lifecycle.horizon_seconds, + expected_reads: planner_selection.expected_reads, + expected_updates: planner_selection.expected_updates, + lifecycle_cost: planner_selection.lifecycle_cost, + }); let binding_key = selected.node_identity; if let Some(existing) = node_bindings.insert(binding_key, materialization) { if existing != materialization { @@ -2029,6 +2039,7 @@ impl PhysicalCompiler { transmission_plan, backend_plan, query_plan, + lifecycle_estimates: lifecycle_estimates.into_values().collect(), }) } } @@ -2072,6 +2083,52 @@ fn summary_agg_metric(node: &SummaryNode) -> Option { .flatten() } +/// Shared selection boundary for canonical startup and compile-and-publish. +/// Certificate-bearing roots stay isolated: equal certificate values do not +/// establish that the certificate's source scope covers another query. +pub fn select_workload_roots( + queries: &mut [PlanningQuery], + roots: Vec>, + evidence: &HashMap, +) -> Result<(), CompileError> { + if roots.len() != queries.len() { + return Err(CompileError::Snapshot( + "canonical root/query mapping is incomplete".into(), + )); + } + let mut cohorts: Vec<(AccuracyTarget, Option, Vec<(usize, Rc)>)> = + Vec::new(); + for (index, root) in roots.into_iter().enumerate() { + let accuracy = &queries[index].accuracy; + let certificate_scope = evidence + .contains_key(&queries[index].query_id) + .then(|| queries[index].query_id.clone()); + if let Some((_, _, roots)) = cohorts + .iter_mut() + .find(|(target, scope, _)| target == accuracy && scope == &certificate_scope) + { + roots.push((index, root)); + } else { + cohorts.push((accuracy.clone(), certificate_scope, vec![(index, root)])); + } + } + for (accuracy, scope, roots) in cohorts { + let model = ControlPlaneCostModel::new(accuracy.clone()); + let certificate = scope.as_ref().and_then(|id| evidence.get(id)); + let selected = crate::planner_selection::select_workload_with_evidence( + roots, + accuracy, + &model, + &QueryEvidence(certificate), + ) + .map_err(|error| CompileError::Snapshot(error.to_string()))?; + for (index, node) in selected { + queries[index].post_asap = node; + } + } + Ok(()) +} + /// Planner-adapter selection step used before physical compilation. Keeping /// this separate makes the ownership boundary explicit: callers supply the /// selected post-ASAP DAG to [`PhysicalCompiler::compile`]. @@ -2246,6 +2303,9 @@ fn validate_window_implementations( struct PlannerPhysicalSelection { lifecycle: CollectorLifecycle, window_framework: SummaryWindowFramework, + expected_reads: f64, + expected_updates: f64, + lifecycle_cost: f64, } fn select_lifecycle( @@ -2253,28 +2313,49 @@ fn select_lifecycle( node: &SummaryNode, model: &ControlPlaneCostModel, environment: &DeploymentEnvironment, + consumers: &[&PlanningQuery], ) -> Result { + // Current lifecycle evidence is per producer with one unit read cost. + // Conflicting source/rate/horizon/cost snapshots cannot be averaged into + // invented evidence. Only recurrence may differ between consumers. + let mut common = query.lifecycle.clone(); + common.evaluation_interval_ms = 0; + for consumer in consumers { + let mut input = consumer.lifecycle.clone(); + input.evaluation_interval_ms = 0; + if input != common { + return Err(CompileError::Lifecycle { + query_id: consumer.query_id.clone(), + reason: "shared producer consumers have conflicting lifecycle evidence".into(), + }); + } + } let workload = QueryWorkload { language: QueryLanguage::PromQL, query_batch: None, - repeating_queries: Some(vec![RepeatingEntry { - query: Query(query.query_id.clone()), - demand: RepeatedDemand::FixedInterval(RepetitionInterval( - query.lifecycle.evaluation_interval_ms, - )), - requirements: QueryRequirements { - accuracy: AccuracyRequirement::Explicit(query.accuracy.clone()), - ..QueryRequirements::default() - }, - predictability: Predictability::Predictable { known_at: None }, - time_selection: TimeSelection { - // Collector windows are retired as whole states. They do not - // claim deletion support for moving-window retractions. - scope: QueryTimeScope::Unknown, - lookback: Some(DurationMs(query.window_secs.saturating_mul(1_000))), - as_of: None, - }, - }]), + repeating_queries: Some( + consumers + .iter() + .map(|query| RepeatingEntry { + query: Query(query.query_id.clone()), + demand: RepeatedDemand::FixedInterval(RepetitionInterval( + query.lifecycle.evaluation_interval_ms, + )), + requirements: QueryRequirements { + accuracy: AccuracyRequirement::Explicit(query.accuracy.clone()), + ..QueryRequirements::default() + }, + predictability: Predictability::Predictable { known_at: None }, + time_selection: TimeSelection { + // Collector windows are retired as whole states. They do not + // claim deletion support for moving-window retractions. + scope: QueryTimeScope::Unknown, + lookback: Some(DurationMs(query.window_secs.saturating_mul(1_000))), + as_of: None, + }, + }) + .collect(), + ), data_workload: Some(DataWorkload { arrival: DataArrival::ContinuouslyIngesting, ingestion_rate: Evidence { @@ -2288,7 +2369,7 @@ fn select_lifecycle( }; let plan = plan_summary_maintenance_lifecycles( Rc::new(node.clone()), - WorkloadDemand::new(&workload, &[0]), + WorkloadDemand::new(&workload, &(0..consumers.len()).collect::>()), environment.observed_at_unix_ms, Some(Horizon(query.lifecycle.horizon_seconds)), SummaryMaintenanceLifecycleCapabilities { @@ -2320,6 +2401,32 @@ fn select_lifecycle( reason: "latest ASAPPlanner selected no window framework from the supplied physical evidence".into(), })?; Ok(PlannerPhysicalSelection { + expected_reads: plan.expected_reads.ok_or_else(|| CompileError::Lifecycle { + query_id: query.query_id.clone(), + reason: "missing joint read demand".into(), + })?, + expected_updates: plan + .update_rate + .ok_or_else(|| CompileError::Lifecycle { + query_id: query.query_id.clone(), + reason: "missing source update demand".into(), + })? + .0 + * query.lifecycle.horizon_seconds, + lifecycle_cost: plan.deployments[0] + .alternatives + .iter() + .find(|alternative| { + alternative.rejection.is_none() + && alternative.summary_maintenance_lifecycle + == guarantee.summary_maintenance_lifecycle + }) + .and_then(|alternative| alternative.total_cost) + .ok_or_else(|| CompileError::Lifecycle { + query_id: query.query_id.clone(), + reason: "missing selected lifecycle cost".into(), + })? + .0, lifecycle: CollectorLifecycle { kind: match guarantee.summary_maintenance_lifecycle { SummaryMaintenanceLifecycle::Ephemeral => "ephemeral", @@ -2351,6 +2458,7 @@ fn select_lifecycle( } struct SelectedMaterialization { + node: Rc, node_identity: usize, metric: String, family: SummaryFamilyType, @@ -2359,6 +2467,52 @@ struct SelectedMaterialization { parameters: Value, } +fn physical_aggregation( + query: &PlanningQuery, + selected: &SelectedMaterialization, + aggregation_id: String, + target: PhysicalDeploymentTarget, +) -> BackendAggregation { + BackendAggregation { + aggregation_id, + metric_name: selected.metric.clone(), + family: physical_materialization_family(&selected.family), + window_secs: query.window_secs, + spatial_filter: String::new(), + grouping: query.group_by.clone(), + item_label: None, + aggregation_input: match target { + PhysicalDeploymentTarget::DistributedCollectors => AggregationInput::SketchEnvelope, + PhysicalDeploymentTarget::BackendLocalRemoteWrite => AggregationInput::Raw, + }, + } +} + +fn materialization_consumers( + queries: &[PlanningQuery], + target: PhysicalDeploymentTarget, +) -> Result>, CompileError> { + let mut consumers = BTreeMap::<_, BTreeSet<_>>::new(); + for (index, query) in queries.iter().enumerate() { + let states = collect_selected_materializations(&query.post_asap).map_err(|reason| { + CompileError::Query { + query_id: query.query_id.clone(), + reason, + } + })?; + for state in states { + let config = backend_plan::aggregation_config_for_materialization( + &physical_aggregation(query, &state, query.query_id.clone(), target), + )?; + consumers + .entry(config.policy_fingerprint()) + .or_default() + .insert(index); + } + } + Ok(consumers) +} + /// Collect every executable materialization leaf in the selected post-ASAP /// graph. Readout context flows through merge nodes, so a graph such as /// `Estimate(Merge(Agg(a), Agg(b)))` creates two physical bindings while the @@ -2415,6 +2569,7 @@ fn collect_selected_materializations( "SummaryAgg has no unique time-series source in post-ASAP IR".to_string() })?; selected.push(SelectedMaterialization { + node: Rc::clone(node), node_identity: Rc::as_ptr(node) as usize, metric, family: SummaryFamilyType::Sketch( @@ -2435,6 +2590,7 @@ fn collect_selected_materializations( "SummaryAgg has no unique time-series source in post-ASAP IR".to_string() })?; selected.push(SelectedMaterialization { + node: Rc::clone(node), node_identity: Rc::as_ptr(node) as usize, metric, family: SummaryFamilyType::ExactAggregate(kind.clone(), params.clone()), @@ -2611,6 +2767,100 @@ mod tests { request_with_evidence(query_id, promql, None).expect("post-ASAP selection") } + // Both production adapters preserve canonical root identity and select the + // whole evidence-free cohort, rather than independently binding roots. + #[test] + fn shared_selection_adapter_preserves_query_mapping() { + let mut workload = request("q90", "quantile_over_time(0.9, m[1m])"); + workload + .queries + .extend(request("q99", "quantile_over_time(0.99, m[1m])").queries); + let roots = workload + .queries + .iter() + .map(|query| { + Rc::new( + crate::query_parser::parse_query_expr_canonical( + &query.query_string, + query.accuracy.clone(), + ) + .unwrap(), + ) + }) + .collect(); + select_workload_roots(&mut workload.queries, roots, &workload.evidence).unwrap(); + let bundle = PhysicalCompiler + .compile(workload, environment(10000)) + .unwrap(); + assert_eq!(bundle.query_plan.entries.len(), 2); + assert_eq!(bundle.collector_plans[0].materializations.len(), 1); + assert_eq!( + bundle + .query_plan + .entries + .values() + .map(|entry| entry.query_id.as_str()) + .collect::>(), + BTreeSet::from(["q90", "q99"]) + ); + } + + // A broken input mapping must be rejected, never silently drop a root. + #[test] + fn shared_selection_rejects_incomplete_root_mapping() { + let mut workload = request("q", "quantile_over_time(0.9, m[1m])"); + assert!(select_workload_roots(&mut workload.queries, vec![], &workload.evidence).is_err()); + } + + // Adding another readout adds recurring reads, not another update stream. + #[test] + fn joint_lifecycle_charges_shared_updates_once() { + let baseline = PhysicalCompiler + .compile( + request("q90", "quantile_over_time(0.9, m[1m])"), + environment(10000), + ) + .unwrap(); + let mut workload = request("q90", "quantile_over_time(0.9, m[1m])"); + let mut second = request("q99", "quantile_over_time(0.99, m[1m])") + .queries + .remove(0); + second.lifecycle.evaluation_interval_ms = 20000; + workload.queries.push(second); + let shared = PhysicalCompiler + .compile(workload, environment(10000)) + .unwrap(); + assert_eq!(shared.lifecycle_estimates.len(), 1); + let estimate = &shared.lifecycle_estimates[0]; + assert_eq!(estimate.consumer_query_ids, vec!["q90", "q99"]); + assert_eq!(estimate.expected_reads, 45.0); + assert_eq!(estimate.expected_updates, 30000.0); + assert!((estimate.lifecycle_cost - 45.8).abs() < 1e-9); + assert_eq!( + estimate.expected_updates, + baseline.lifecycle_estimates[0].expected_updates + ); + assert!( + (estimate.lifecycle_cost - baseline.lifecycle_estimates[0].lifecycle_cost - 1.5).abs() + < 1e-9 + ); + } + + // Unknown joint provenance cannot be replaced by whichever root came first. + #[test] + fn shared_lifecycle_rejects_conflicting_source_evidence() { + let mut workload = request("q90", "quantile_over_time(0.9, m[1m])"); + let mut second = request("q99", "quantile_over_time(0.99, m[1m])") + .queries + .remove(0); + second.lifecycle.ingestion_rate_per_second = 200.0; + workload.queries.push(second); + assert!(matches!( + PhysicalCompiler.compile(workload, environment(10000)), + Err(CompileError::Lifecycle { .. }) + )); + } + #[test] fn shared_materialization_is_emitted_once_for_every_runtime() { for target in [ diff --git a/control_plane/src/planner_selection.rs b/control_plane/src/planner_selection.rs index a41ef105..c464bfed 100644 --- a/control_plane/src/planner_selection.rs +++ b/control_plane/src/planner_selection.rs @@ -158,12 +158,33 @@ pub fn select_workload( roots: Vec<(usize, Rc)>, accuracy: AccuracyTarget, cost_model: &dyn CostModel, +) -> Result)>, SelectionError> { + select_workload_with_evidence( + roots, + accuracy, + cost_model, + &asap_aware_mapping::NoAccuracyEvidence, + ) +} + +/// The entire cohort uses the same scoped accuracy certificate; callers must +/// not spread one query's evidence to unrelated workload roots. +pub fn select_workload_with_evidence( + roots: Vec<(usize, Rc)>, + accuracy: AccuracyTarget, + cost_model: &dyn CostModel, + evidence: &dyn AccuracyEvidenceProvider, ) -> Result)>, SelectionError> { // Canonical CSE still runs inside search_workload_with_targets. Do not // offer CSE's per-invocation recompute alternative: this runtime currently // provisions continuously maintained, content-addressed state only. let strategies: Vec> = vec![ - Box::new(SketchAlgorithmStrategy::new(cost_model)), + Box::new(SketchAlgorithmStrategy::with_models_and_evidence( + cost_model, + &asap_aware_mapping::DefaultAccuracyModel, + &asap_aware_mapping::EqualSplitAllocator, + evidence, + )), Box::new(asap_aware_mapping::SemanticEquivalentRewriteStrategy), ]; let space = asap_aware_mapping::search_workload_with_targets( diff --git a/data_plane/tests/backend_process_e2e.rs b/data_plane/tests/backend_process_e2e.rs index a0beba65..338f2e9e 100644 --- a/data_plane/tests/backend_process_e2e.rs +++ b/data_plane/tests/backend_process_e2e.rs @@ -67,18 +67,63 @@ async fn wait_http(client: &reqwest::Client, url: &str, child: &mut Child, name: panic!("{name} did not become ready at {url}"); } -fn ddsketch_export(metric: &str, timestamp_ns: u64, values: &[f64], alpha: f64) -> Vec { +fn ddsketch_export( + metric: &str, + timestamp_ns: u64, + values: &[f64], + alpha: f64, + plan: &serde_json::Value, + sequence: u64, +) -> Vec { let mut sketch = asap_sketchlib::DdSketch::new(alpha); for value in values { sketch.update(*value); } - let point = DdSketchDataPoint { - attributes: vec![KeyValue { - key: "service".into(), + let materialization = plan["materializations"][0]["materialization"] + .as_u64() + .unwrap(); + let mut attributes = vec![KeyValue { + key: "service".into(), + value: Some(AnyValue { + value: Some(any_value::Value::StringValue("whole-e2e".into())), + }), + }]; + for (key, value) in [ + ("identity_version", "1".into()), + ("plan_id", plan["envelope"]["plan_id"].to_string()), + ("plan_version", plan["envelope"]["plan_version"].to_string()), + ( + "backend_compat", + control_plane::backend_plan::BACKEND_COMPAT.into(), + ), + ("materialization", materialization.to_string()), + ( + "series_identity", + data_plane::drivers::ingest::canonical_attrs_fingerprint(&[("service", "whole-e2e")]), + ), + ( + "schema_id", + format!( + "{}:summary-state:v1:{materialization}", + control_plane::backend_plan::BACKEND_COMPAT + ), + ), + ("producer_id", "whole-e2e-collector".into()), + ("producer_epoch", "process-e2e".into()), + ("sequence", sequence.to_string()), + ("kind", "full".into()), + ("encoding", "sketchlib_protobuf_v1".into()), + ("checkpoint_id", format!("checkpoint-{sequence}")), + ] { + attributes.push(KeyValue { + key: format!("asap.frame.{key}"), value: Some(AnyValue { - value: Some(any_value::Value::StringValue("whole-e2e".into())), + value: Some(any_value::Value::StringValue(value)), }), - }], + }); + } + let point = DdSketchDataPoint { + attributes, start_time_unix_nano: timestamp_ns.saturating_sub(1_000_000_000), time_unix_nano: timestamp_ns, sketch: DdSketchState { @@ -178,7 +223,13 @@ async fn send_agent_message( .expect("send OpAMP AgentToServer"); } -async fn apply_next_collector_plan(address: String) -> serde_json::Value { +type CollectorSocket = + tokio_tungstenite::WebSocketStream>; + +async fn apply_next_collector_plan( + address: String, + accept: bool, +) -> (serde_json::Value, CollectorSocket) { let mut socket = connect_collector(&address).await; send_agent_message( &mut socket, @@ -190,7 +241,13 @@ async fn apply_next_collector_plan(address: String) -> serde_json::Value { }, ) .await; + respond_next_collector_plan(socket, accept).await +} +async fn respond_next_collector_plan( + mut socket: CollectorSocket, + accept: bool, +) -> (serde_json::Value, CollectorSocket) { let frame = tokio::time::timeout(Duration::from_secs(10), socket.next()) .await .expect("controller did not publish a collector plan") @@ -205,6 +262,12 @@ async fn apply_next_collector_plan(address: String) -> serde_json::Value { .expect("collector-plan custom message"); assert_eq!(custom.capability, COLLECTOR_PLAN_CAPABILITY); assert_eq!(custom.r#type, COLLECTOR_PLAN_MESSAGE); + let decoded = asap_precompute_rs::collector_plan::CollectorPlan::from_json( + &custom.data, + "whole-e2e-collector", + ) + .expect("actual Collector validator accepts the emitted plan"); + assert_eq!(decoded.to_precompute_config_set().unwrap().configs.len(), 1); let plan: serde_json::Value = serde_json::from_slice(&custom.data).expect("decode collector physical plan"); let plan_id = plan["envelope"]["plan_id"] @@ -217,8 +280,12 @@ async fn apply_next_collector_plan(address: String) -> serde_json::Value { let status = serde_json::to_vec(&CollectorPlanStatus { plan_id, plan_version, - status: CollectorPlanStatusKind::Applied, - error: None, + status: if accept { + CollectorPlanStatusKind::Staged + } else { + CollectorPlanStatusKind::Failed + }, + error: (!accept).then(|| "injected Collector staging failure".into()), }) .expect("encode applied status"); send_agent_message( @@ -233,7 +300,7 @@ async fn apply_next_collector_plan(address: String) -> serde_json::Value { }, ) .await; - plan + (plan, socket) } #[tokio::test] @@ -312,45 +379,65 @@ async fn production_control_plane_to_data_plane_otlp_to_promql() { ) .await; - let collector = tokio::spawn(apply_next_collector_plan(control_opamp.clone())); + let collector = tokio::spawn(apply_next_collector_plan(control_opamp.clone(), true)); let observed_at_ms = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .expect("system clock") .as_millis() as u64; + let mut request = serde_json::json!({ + "queries": [{ + "query_id": "whole-process-e2e-query", + "query_string": "quantile_over_time(0.99, whole_process_e2e_latency_ms[1s])", + "metric": "whole_process_e2e_latency_ms", + "window_secs": 1, + "group_by": ["service"], + "accuracy": {"Epsilon": 0.01}, + "window_implementations": [{ + "implementation_id": "collector-tumbling-v1", "framework": "tumbling", + "window_secs": 1, "pane_secs": 1, "state_layout": "anchored-pane-v1", + "cost": { + "model_version": "process-e2e-v1", "workload_fingerprint": "shared-quantiles", + "observed_at_unix_ms": observed_at_ms, "valid_for_ms": 60000, + "horizon_seconds": 300.0, "cpu_cost": 1.0, "weighted_cost": 1.0, + "peak_memory_bytes": 4096, "network_bytes": 1024, "storage_bytes": 2048, + "source_scan_bytes": 0 + } + }], + "lifecycle": { + "evaluation_interval_ms": 1000, + "ingestion_rate_per_second": 100.0, + "evidence_observed_at_unix_ms": observed_at_ms, + "evidence_valid_for_ms": 60000, + "horizon_seconds": 300.0, + "costs": { + "build": 10.0, + "maintenance_per_update": 0.001, + "read": 0.1, + "retention_per_second": 0.001, + "retirement": 1.0 + } + } + }], + "collector_ids": ["whole-e2e-collector"], + "capability_snapshot_id": "whole-e2e-capabilities", + "evidence": {}, + "planner_revision": PLANNER_REVISION, + "max_evidence_age_ms": 60000, + "plan_version": 1, + "activation_unix_ms": observed_at_ms, + "expiry_unix_ms": null, + "backend_compat": control_plane::backend_plan::BACKEND_COMPAT, + "apply_timeout_ms": 10000 + }); + let mut second = request["queries"][0].clone(); + second["query_id"] = "whole-process-e2e-median".into(); + second["query_string"] = "quantile_over_time(0.5, whole_process_e2e_latency_ms[1s])".into(); + request["queries"].as_array_mut().unwrap().push(second); let publication_response = client .post(format!( "{control_base}/api/v1/physical-plan/compile-and-publish" )) - .json(&serde_json::json!({ - "queries": [{ - "query_id": "whole-process-e2e-query", - "query_string": "quantile_over_time(0.99, whole_process_e2e_latency_ms[30s])", - "metric": "whole_process_e2e_latency_ms", - "window_secs": 1, - "group_by": ["service"], - "accuracy": {"Epsilon": 0.01}, - "lifecycle": { - "evaluation_interval_ms": 1000, - "ingestion_rate_per_second": 100.0, - "evidence_observed_at_unix_ms": observed_at_ms, - "evidence_valid_for_ms": 60000, - "horizon_seconds": 300.0, - "costs": { - "build": 10.0, - "maintenance_per_update": 0.001, - "read": 0.1, - "retention_per_second": 0.001, - "retirement": 1.0 - } - } - }], - "collector_ids": ["whole-e2e-collector"], - "capability_snapshot_id": "whole-e2e-capabilities", - "evidence": {}, - "planner_revision": PLANNER_REVISION, - "max_evidence_age_ms": 60000, - "apply_timeout_ms": 10000 - })) + .json(&request) .send() .await .expect("request physical-plan publication"); @@ -365,7 +452,7 @@ async fn production_control_plane_to_data_plane_otlp_to_promql() { ); let publication: serde_json::Value = serde_json::from_str(&publication_body).expect("decode publication response"); - let collector_plan = collector.await.expect("collector task completed"); + let (collector_plan, collector_socket) = collector.await.expect("collector task completed"); assert_eq!( publication["plan_id"], collector_plan["envelope"]["plan_id"] @@ -407,10 +494,14 @@ async fn production_control_plane_to_data_plane_otlp_to_promql() { let now = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .expect("system clock"); - let sample_ns = now.as_nanos() as u64; - let raw_values = (1..=100).map(|value| value as f64).collect::>(); + let window_ms = planned_window_secs * 1000; + let window_end_ms = (now.as_millis() as u64 / window_ms) * window_ms; + let sample_ns = window_end_ms * 1_000_000; + // These ranks are integral under both PromQL interpolation and sketch + // order-statistic readout; interpolation coverage is a separate contract. + let raw_values = (1..=101).map(|value| value as f64).collect::>(); let reference_p99 = exact_quantile(&raw_values, 0.99); - client + let ingestion = client .post(format!("http://{otlp_http}/v1/metrics")) .header("content-type", "application/x-protobuf") .body(ddsketch_export( @@ -418,23 +509,19 @@ async fn production_control_plane_to_data_plane_otlp_to_promql() { sample_ns, &raw_values, planned_alpha, + &collector_plan, + 1, )) .send() .await - .expect("POST OTLP to production data plane") - .error_for_status() - .expect("data plane accepted OTLP"); + .expect("POST OTLP to production data plane"); + let status = ingestion.status(); + let body = ingestion.text().await.unwrap(); + assert!(status.is_success(), "OTLP rejected: {status}: {body}"); - // Advance event time after the controller-selected tumbling window has - // really ended. The E2E uses no artificial future timestamp here. - let window_ms = planned_window_secs * 1000; - let now_ms = now.as_millis() as u64; - let window_end_ms = (now_ms / window_ms + 1) * window_ms; - tokio::time::sleep(Duration::from_millis(window_end_ms - now_ms + 100)).await; - let watermark_ns = std::time::SystemTime::now() - .duration_since(std::time::UNIX_EPOCH) - .expect("system clock") - .as_nanos() as u64; + // The sealed full frame already carries its exact window. A subsequent + // checkpoint exercises the same producer's next window independently. + let watermark_ns = sample_ns + window_ms * 1_000_000; client .post(format!("http://{otlp_http}/v1/metrics")) .header("content-type", "application/x-protobuf") @@ -443,6 +530,8 @@ async fn production_control_plane_to_data_plane_otlp_to_promql() { watermark_ns, &[], planned_alpha, + &collector_plan, + 2, )) .send() .await @@ -450,12 +539,15 @@ async fn production_control_plane_to_data_plane_otlp_to_promql() { .error_for_status() .expect("data plane accepted watermark"); - let query = "quantile_over_time(0.99, whole_process_e2e_latency_ms[30s])"; + let query = "quantile_over_time(0.99, whole_process_e2e_latency_ms[1s])"; let mut last_response = serde_json::Value::Null; for _ in 0..50 { let response: serde_json::Value = client .get(format!("{data_base}/api/v1/query")) - .query(&[("query", query)]) + .query(&[ + ("query", query.to_string()), + ("time", (window_end_ms as f64 / 1000.0).to_string()), + ]) .send() .await .expect("query production data plane") @@ -474,6 +566,96 @@ async fn production_control_plane_to_data_plane_otlp_to_promql() { response["data"]["result"][0]["metric"]["service"], "whole-e2e" ); + let median: serde_json::Value = client + .get(format!("{data_base}/api/v1/query")) + .query(&[ + ( + "query", + "quantile_over_time(0.5, whole_process_e2e_latency_ms[1s])".to_string(), + ), + ("time", (window_end_ms as f64 / 1000.0).to_string()), + ]) + .send() + .await + .unwrap() + .json() + .await + .unwrap(); + let median_value = first_scalar(&median).expect("second consumer is warm"); + let reference = exact_quantile(&raw_values, 0.5); + assert!( + (median_value - reference).abs() / reference <= planned_alpha * 1.05, + "{median}" + ); + // A rejected successor must not replace the active generation. + request["plan_version"] = 2.into(); + request["activation_unix_ms"] = (std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_millis() as u64 + + 1000) + .into(); + let collector = tokio::spawn(respond_next_collector_plan(collector_socket, false)); + let failed = client + .post(format!( + "{control_base}/api/v1/physical-plan/compile-and-publish" + )) + .json(&request) + .send() + .await + .unwrap(); + let failed_status = failed.status().as_u16(); + let failed_body = failed.text().await.unwrap(); + assert_eq!(failed_status, 502, "{failed_body}"); + assert!( + failed_body.contains("injected Collector staging failure"), + "{failed_body}" + ); + let (rejected_plan, _collector_socket) = collector.await.unwrap(); + let rejected_frame = client + .post(format!("http://{otlp_http}/v1/metrics")) + .header("content-type", "application/x-protobuf") + .body(ddsketch_export( + "whole_process_e2e_latency_ms", + sample_ns, + &[9999.0], + planned_alpha, + &rejected_plan, + 1, + )) + .send() + .await + .unwrap(); + assert_eq!( + rejected_frame.status().as_u16(), + 422, + "inactive generation frame was accepted" + ); + let still_active: serde_json::Value = client + .get(format!("{data_base}/api/v1/backend-plan")) + .send() + .await + .unwrap() + .json() + .await + .unwrap(); + assert_eq!( + still_active, backend_plan, + "failed rollout changed active plan" + ); + let still_warm: serde_json::Value = client + .get(format!("{data_base}/api/v1/query")) + .query(&[ + ("query", query.to_string()), + ("time", (window_end_ms as f64 / 1000.0).to_string()), + ]) + .send() + .await + .unwrap() + .json() + .await + .unwrap(); + assert_eq!(first_scalar(&still_warm), Some(value)); return; } last_response = response;