diff --git a/crates/asap-aware-mapping/src/lib.rs b/crates/asap-aware-mapping/src/lib.rs index 716c6ee4..921f504e 100644 --- a/crates/asap-aware-mapping/src/lib.rs +++ b/crates/asap-aware-mapping/src/lib.rs @@ -228,7 +228,7 @@ pub use summary_maintenance_lifecycle::{ SummaryMaintenanceLifecycleCapabilities, SummaryMaintenanceLifecycleChoiceError, SummaryMaintenanceLifecycleCostInputs, SummaryMaintenanceLifecyclePlan, SummaryMaintenanceLifecyclePlanError, SummaryMaintenanceLifecycleRejection, - SummaryMaintenanceLifecycleSelectionError, WorkloadDemand, + SummaryMaintenanceLifecycleSelectionError, SummaryMaintenanceTimingError, WorkloadDemand, }; pub use topk_reuse::TopKLimitReuseStrategy; diff --git a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs index 7e40645f..43bacfcb 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs @@ -20,10 +20,11 @@ use std::collections::{HashMap, HashSet}; use std::rc::Rc; use asap_types::post_asap::{ - compile_post_asap_dag_with_node_ids, EvaluationSchedule, ExecutionDataStateError, - OutputRepresentation, PostAsapNodeId, ResultGuarantee, SummaryExpr, - SummaryMaintenanceLifecycle, SummaryMaintenanceLifecycleGuarantee, SummaryMaintenanceMode, - SummaryNode, SummaryWindowFramework, + compile_post_asap_dag, compile_post_asap_dag_with_node_ids, EvaluationSchedule, + ExecutionDataStateError, ExecutionTiming, OutputRepresentation, PostAsapDag, + PostAsapDagValidationError, PostAsapNodeId, PostAsapOperatorPayload, ResultGuarantee, + SummaryExpr, SummaryMaintenanceLifecycle, SummaryMaintenanceLifecycleGuarantee, + SummaryMaintenanceMode, SummaryNode, SummaryWindowFramework, ValueOperation, }; use asap_types::pre_asap::QueryExpr; use asap_types::types::AccuracyTarget; @@ -211,6 +212,84 @@ pub struct SummaryMaintenanceLifecyclePlan { pub raw_recompute_total_cost: Option, } +/// Why a lifecycle plan cannot assign execution timing to its DAG. +#[derive(Debug, thiserror::Error, PartialEq)] +pub enum SummaryMaintenanceTimingError { + #[error(transparent)] + InvalidPostAsapDag(#[from] ExecutionDataStateError), + #[error("summary {0:?} has no selected lifecycle")] + UnselectedLifecycle(PostAsapNodeId), + /// Lifecycle enumeration covers `SummaryAgg` states only; timing for other + /// retained state would otherwise be guessed. + #[error("node {0:?} maintains state that has no summary-maintenance lifecycle")] + UnplannedMaintainedState(PostAsapNodeId), + #[error(transparent)] + InvalidPhases(#[from] PostAsapDagValidationError), +} + +impl SummaryMaintenanceLifecyclePlan { + /// The post-ASAP DAG of [`Self::root`] with every node's timing derived + /// from the selected lifecycles, so physical compilation places it. + /// + /// A retained (non-`Ephemeral`) state outlives one query, so it and every + /// input it consumes run at ingestion time. Every other node runs at query + /// time: readouts and consumers of retained state, and each `Ephemeral` + /// state not consumed by retained state together with its inputs, whose + /// raw data the deployment must supply as a query source. Timings already + /// on the root are ignored. + pub fn execution_timed_dag(&self) -> Result { + let dag = compile_post_asap_dag(&self.root)?; + if let Some(node) = dag.nodes.iter().find(|node| { + matches!( + node.payload, + PostAsapOperatorPayload::Value { + operation: ValueOperation::MaintainPopulation { .. } + } + ) + }) { + return Err(SummaryMaintenanceTimingError::UnplannedMaintainedState( + node.id, + )); + } + let mut pending = Vec::new(); + for deployment in &self.deployments { + let guarantee = deployment + .summary_maintenance_lifecycle_guarantee + .as_ref() + .ok_or(SummaryMaintenanceTimingError::UnselectedLifecycle( + deployment.post_asap_node_id, + ))?; + if guarantee.summary_maintenance_lifecycle != SummaryMaintenanceLifecycle::Ephemeral { + pending.push(deployment.post_asap_node_id); + } + } + let mut ingestion = HashSet::new(); + while let Some(id) = pending.pop() { + if ingestion.insert(id) { + pending.extend( + dag.edges + .iter() + .filter(|edge| edge.consumer == id) + .map(|edge| edge.producer), + ); + } + } + let phases = dag + .nodes + .iter() + .map(|node| { + let timing = if ingestion.contains(&node.id) { + ExecutionTiming::IngestionTime + } else { + ExecutionTiming::QueryTime + }; + (node.id, timing) + }) + .collect(); + Ok(dag.with_execution_phases(&phases)?) + } +} + /// Explicit association between a materialized target and the normalized /// workload entries whose demand consumes it. /// @@ -2845,4 +2924,294 @@ mod tests { .iter() .any(|deployment| Rc::ptr_eq(&deployment.summary, &shared))); } + + fn readout(state: &Rc) -> Rc { + Rc::new(SummaryNode { + expr: SummaryExpr::ValueOperation { + child: Rc::clone(state), + operation: ValueOperation::FinalizeExactAccumulator, + timing: ExecutionTiming::QueryTime, + }, + schema: SummarySchema { + fields: vec![SummaryField { + name: "value".into(), + dtype: SummaryFamilyType::Plain(DataType::Float64), + nullable: false, + }], + time_index: None, + }, + guarantee: Some(ResultGuarantee::exact("sum")), + }) + } + + fn lifecycle_matching( + alternatives: &[SummaryMaintenanceLifecycleAlternative], + kind: fn(&SummaryMaintenanceLifecycle) -> bool, + ) -> SummaryMaintenanceLifecycle { + alternatives + .iter() + .map(|alternative| &alternative.summary_maintenance_lifecycle) + .find(|lifecycle| kind(lifecycle)) + .expect("lifecycle kind is an alternative") + .clone() + } + + /// Bind the lifecycle `choose` picks for every state of `root`, then + /// derive the timed DAG. + fn timed_dag( + root: Rc, + workload: &QueryWorkload, + data: &DataWorkload, + horizon: Option, + choose: impl Fn(&SummaryMaintenanceDeployment) -> SummaryMaintenanceLifecycle, + ) -> PostAsapDag { + let candidates = enumerate_summary_maintenance_lifecycles( + root, + WorkloadDemand::new_with_data(workload, data, &[0]), + 1_000, + horizon, + SummaryMaintenanceLifecycleCapabilities::ALL, + &UnitCosts, + ) + .unwrap(); + let choice: Vec<_> = candidates + .deployments() + .iter() + .map(|deployment| (deployment.post_asap_node_id, choose(deployment))) + .collect(); + let dag = candidates + .select(&choice) + .unwrap() + .execution_timed_dag() + .unwrap(); + dag.validate().unwrap(); + dag + } + + /// Operator kinds in node-id order, each paired with its timing. + fn timings(dag: &PostAsapDag) -> Vec<(&'static str, ExecutionTiming)> { + dag.nodes + .iter() + .map(|node| { + let kind = match node.payload { + PostAsapOperatorPayload::Fallback { .. } => "raw", + PostAsapOperatorPayload::SummaryAgg { .. } => "state", + PostAsapOperatorPayload::Value { .. } => "readout", + PostAsapOperatorPayload::Binary { .. } => "binary", + _ => "other", + }; + (kind, node.output_state.timing) + }) + .collect() + } + + const INGEST: ExecutionTiming = ExecutionTiming::IngestionTime; + const QUERY: ExecutionTiming = ExecutionTiming::QueryTime; + + // Every retained lifecycle kind runs its state and inputs at ingestion + // time and its readout at query time. + #[test] + fn retained_lifecycles_time_state_and_inputs_at_ingestion() { + let mut scheduled = batch(Predictability::Predictable { + known_at: Some(TimestampMs(1_000)), + }); + scheduled.execute_at = Some(TimestampMs(11_000)); + type Case = ( + QueryWorkload, + DataWorkload, + Option, + fn(&SummaryMaintenanceLifecycle) -> bool, + ); + let cases: [Case; 3] = [ + ( + workload(vec![], vec![repeating()], continuous(1_000, 60_000)), + continuous(1_000, 60_000), + Some(Horizon(10.0)), + |lifecycle| { + matches!( + lifecycle, + SummaryMaintenanceLifecycle::ContinuouslyMaintained + ) + }, + ), + ( + workload(vec![], vec![repeating()], at_rest()), + at_rest(), + Some(Horizon(10.0)), + |lifecycle| matches!(lifecycle, SummaryMaintenanceLifecycle::Shared { .. }), + ), + ( + workload(vec![scheduled], vec![], at_rest()), + at_rest(), + None, + |lifecycle| matches!(lifecycle, SummaryMaintenanceLifecycle::Prepared { .. }), + ), + ]; + for (workload, data, horizon, kind) in cases { + let dag = timed_dag( + readout(&summary()), + &workload, + &data, + horizon, + |deployment| lifecycle_matching(&deployment.alternatives, kind), + ); + assert_eq!( + timings(&dag), + [("raw", INGEST), ("state", INGEST), ("readout", QUERY)] + ); + } + } + + // An Ephemeral state, its raw input, and its readout all run at query time. + #[test] + fn ephemeral_lifecycle_times_state_and_downstream_at_query() { + let workload = workload(vec![batch(Predictability::AdHoc)], vec![], at_rest()); + let dag = timed_dag(readout(&summary()), &workload, &at_rest(), None, |_| { + SummaryMaintenanceLifecycle::Ephemeral + }); + assert_eq!( + timings(&dag), + [("raw", QUERY), ("state", QUERY), ("readout", QUERY)] + ); + } + + // One state read by two consumers is one deployment; its timing follows + // that single choice while both consumers run at query time. + #[test] + fn shared_state_is_timed_once_for_all_consumers() { + let state = summary(); + let lhs = readout(&state); + let rhs = Rc::new(lhs.as_ref().clone()); + let root = Rc::new(SummaryNode { + expr: SummaryExpr::BinaryOp { + lhs, + rhs, + operator: asap_types::post_asap::BinaryOperator { + kind: asap_types::pre_asap::BinaryOpKind::Arithmetic( + asap_types::pre_asap::ArithmeticOpKind::Add, + ), + vector_match: None, + checked_relative_division: false, + checked_finite_division: false, + }, + timing: QUERY, + }, + schema: readout(&state).schema.clone(), + guarantee: None, + }); + let data = continuous(1_000, 60_000); + let workload = workload(vec![], vec![repeating()], data.clone()); + let dag = timed_dag(root, &workload, &data, Some(Horizon(10.0)), |deployment| { + assert!(Rc::ptr_eq(&deployment.summary, &state)); + SummaryMaintenanceLifecycle::ContinuouslyMaintained + }); + assert_eq!( + timings(&dag), + [ + ("raw", INGEST), + ("state", INGEST), + ("readout", QUERY), + ("readout", QUERY), + ("binary", QUERY), + ] + ); + } + + // An Ephemeral state consumed by retained state is built on the retained + // state's ingestion path; it is not retained, but cannot run at query time. + #[test] + fn ephemeral_state_feeding_retained_state_runs_at_ingestion() { + let mut scheduled = batch(Predictability::Predictable { + known_at: Some(TimestampMs(1_000)), + }); + scheduled.execute_at = Some(TimestampMs(11_000)); + let root = nested_summary(); + let workload = workload(vec![scheduled], vec![], at_rest()); + let dag = timed_dag( + Rc::clone(&root), + &workload, + &at_rest(), + None, + |deployment| { + if Rc::ptr_eq(&deployment.summary, &root) { + lifecycle_matching(&deployment.alternatives, |lifecycle| { + matches!(lifecycle, SummaryMaintenanceLifecycle::Prepared { .. }) + }) + } else { + SummaryMaintenanceLifecycle::Ephemeral + } + }, + ); + assert_eq!( + timings(&dag), + [("raw", INGEST), ("state", INGEST), ("state", INGEST)] + ); + } + + // Timing is not derived for a state without a selected lifecycle, and a + // raw-recompute plan runs entirely at query time. + #[test] + fn timing_requires_a_selected_lifecycle_for_every_state() { + let workload = workload(vec![batch(Predictability::AdHoc)], vec![], at_rest()); + let data = at_rest(); + let demand = WorkloadDemand::new_with_data(&workload, &data, &[0]); + let plan = plan_summary_maintenance_lifecycles( + readout(&summary()), + demand, + 1_000, + None, + SummaryMaintenanceLifecycleCapabilities::ALL, + &crate::cost_model::DefaultCostModel, + ) + .unwrap(); + assert_eq!( + plan.execution_timed_dag().unwrap_err(), + SummaryMaintenanceTimingError::UnselectedLifecycle( + plan.deployments[0].post_asap_node_id + ) + ); + let raw = plan_summary_maintenance_lifecycles( + crate::replacement::keep_pre_asap(&sum_query()).unwrap(), + demand, + 1_000, + None, + SummaryMaintenanceLifecycleCapabilities::ALL, + &UnitCosts, + ) + .unwrap(); + assert_eq!( + timings(&raw.execution_timed_dag().unwrap()), + [("raw", QUERY)] + ); + } + + // A maintained population is retained state the lifecycle plan does not + // enumerate, so its timing is refused rather than guessed. + #[test] + fn timing_refuses_state_outside_the_lifecycle_plan() { + let target = Rc::new(crate::test_support::lower_promql( + "sum(a)", + AccuracyTarget::Exact, + )); + let root = crate::maintained_population::MaintainedPopulationStrategy::new( + std::slice::from_ref(&target), + ) + .candidate(&target) + .unwrap(); + let workload = workload(vec![batch(Predictability::AdHoc)], vec![], at_rest()); + let plan = plan_summary_maintenance_lifecycles( + root, + WorkloadDemand::new_with_data(&workload, &at_rest(), &[0]), + 1_000, + None, + SummaryMaintenanceLifecycleCapabilities::ALL, + &UnitCosts, + ) + .unwrap(); + assert!(plan.deployments.is_empty()); + assert!(matches!( + plan.execution_timed_dag(), + Err(SummaryMaintenanceTimingError::UnplannedMaintainedState(_)) + )); + } } diff --git a/crates/integration-tests/tests/summary_maintenance_lifecycle_e2e.rs b/crates/integration-tests/tests/summary_maintenance_lifecycle_e2e.rs index d5406f23..77cdd3d7 100644 --- a/crates/integration-tests/tests/summary_maintenance_lifecycle_e2e.rs +++ b/crates/integration-tests/tests/summary_maintenance_lifecycle_e2e.rs @@ -426,3 +426,193 @@ fn continuous_lifecycle_compiles_and_executes_spatial_kll() { ); } } + +fn quantile_workload(query: &str) -> PlanningWorkload { + let mut workload = dashboard_workload(); + workload.query_workload.query_batch.as_mut().unwrap()[0].query = Query(query.into()); + workload.query_workload.repeating_queries.as_mut().unwrap()[0].query = Query(query.into()); + workload +} + +/// Precompute outputs implied by timing: ingestion-time nodes read by a +/// query-time node, or the root when it is itself ingestion-timed. +fn ingestion_frontier(dag: &asap_types::post_asap::PostAsapDag) -> Vec { + use asap_types::post_asap::ExecutionTiming::IngestionTime; + let timing = |id| { + dag.nodes + .iter() + .find(|node| node.id == id) + .unwrap() + .output_state + .timing + }; + let mut frontier: Vec<_> = dag + .nodes + .iter() + .filter(|node| { + node.output_state.timing == IngestionTime + && (node.id == dag.root + || dag.edges.iter().any(|edge| { + edge.producer == node.id && timing(edge.consumer) != IngestionTime + })) + }) + .map(|node| u64::from(node.id.0)) + .collect(); + frontier.sort_unstable(); + frontier +} + +/// For existing PromQL fixtures, Planner's own retained lifecycle selection +/// reproduces the timing that realization strategies assign today. +#[test] +fn planner_lifecycle_selection_reproduces_strategy_timing() { + for query in [ + "quantile_over_time(0.99, latency[5m])", + "quantile(0.99, latency)", + "sum by(job)(rate(m[1m]))", + ] { + let plan = selected_plan(&quantile_workload(query)); + assert!(!plan.selected_raw_recompute, "{query}"); + assert!(plan.deployments.iter().all(|deployment| { + deployment + .summary_maintenance_lifecycle_guarantee + .as_ref() + .is_some_and(|guarantee| { + guarantee.summary_maintenance_lifecycle + != SummaryMaintenanceLifecycle::Ephemeral + }) + })); + let strategy = asap_types::post_asap::compile_post_asap_dag(&plan.root).unwrap(); + assert_eq!(plan.execution_timed_dag().unwrap(), strategy, "{query}"); + } +} + +/// An explicitly chosen lifecycle reaches physical compilation through timing: +/// ContinuouslyMaintained puts the state in precompute, Ephemeral leaves +/// precompute empty and reads the raw source at query time; both answer alike. +#[test] +fn chosen_lifecycle_timing_decides_precompute_contents() { + use asap_aware_mapping::enumerate_summary_maintenance_lifecycles; + use asap_physical_operators::{ + physical_planner::{compile_candidate, InputContract}, + runtime::Scope, + values::{Batch, Value}, + }; + use asap_types::{ + post_asap::{PostAsapOperatorPayload, SummaryFamilyType}, + pre_asap::DataType, + }; + use std::{collections::BTreeMap, sync::Arc}; + + let workload = quantile_workload("quantile(0.99, latency)"); + let root = selected_plan(&workload).root; + let mut answers = Vec::new(); + for lifecycle in [ + SummaryMaintenanceLifecycle::ContinuouslyMaintained, + SummaryMaintenanceLifecycle::Ephemeral, + ] { + let candidates = enumerate_summary_maintenance_lifecycles( + Rc::clone(&root), + WorkloadDemand::new_with_data( + &workload.query_workload, + workload.data_workload.as_ref().unwrap(), + &[1], + ), + NOW_MS, + Some(Horizon(100.)), + SummaryMaintenanceLifecycleCapabilities::ALL, + &FullyCostedRuntime, + ) + .unwrap(); + let [deployment] = candidates.deployments() else { + panic!("one summary state"); + }; + let id = deployment.post_asap_node_id; + let state = u64::from(id.0); + let dag = candidates + .select(&[(id, lifecycle.clone())]) + .unwrap() + .execution_timed_dag() + .unwrap(); + let raw = dag + .nodes + .iter() + .find(|node| matches!(node.payload, PostAsapOperatorPayload::Fallback { .. })) + .unwrap(); + let (raw_id, schema) = (u64::from(raw.id.0), Arc::new(raw.output_schema.clone())); + let frontier = ingestion_frontier(&dag); + let candidate = compile_candidate( + &dag, + BTreeMap::from([(raw_id, InputContract::bounded(schema.clone()))]), + &[u64::from(dag.root.0)], + &frontier, + ) + .unwrap(); + let rows = (1..=100) + .map(|value| { + schema + .fields + .iter() + .map(|field| match field.dtype { + SummaryFamilyType::Plain(DataType::Float64) => { + Value::Float64(f64::from(value)) + } + SummaryFamilyType::Plain(DataType::Timestamp) => Value::Timestamp(300_000), + _ => panic!("unexpected field {field:?}"), + }) + .collect() + }) + .collect(); + let raw_batch = Batch::try_new(schema.clone(), rows).unwrap(); + let query_scope = Scope::Query { + evaluation_time_ms: 300_000, + revision: 1, + }; + let result = if lifecycle == SummaryMaintenanceLifecycle::Ephemeral { + assert!(frontier.is_empty()); + assert!(candidate.precompute.is_none()); + physical_common::execute( + &candidate.query, + BTreeMap::from([(raw_id, raw_batch)]), + query_scope, + ) + } else { + assert_eq!(frontier, [state]); + assert_eq!( + candidate + .materialized_outputs + .keys() + .copied() + .collect::>(), + [state] + ); + let stored = physical_common::execute( + candidate.precompute.as_ref().unwrap(), + BTreeMap::from([(raw_id, raw_batch)]), + Scope::Ingestion { + window_start_ms: 0, + window_end_ms: 300_000, + revision: 1, + }, + ); + physical_common::execute( + &candidate.query, + BTreeMap::from([(state, stored[0][0].clone())]), + query_scope, + ) + }; + answers.push( + result[0] + .iter() + .flat_map(|batch| batch.rows()) + .flat_map(|row| row.iter()) + .filter_map(|value| match value { + Value::Float64(value) => Some(*value), + _ => None, + }) + .collect::>(), + ); + } + assert_eq!(answers[0], answers[1]); + assert_eq!(answers[0].len(), 1); +} diff --git a/docs/design_docs/physical-planning-and-deployment.md b/docs/design_docs/physical-planning-and-deployment.md index 2aad9ca7..69e9b90a 100644 --- a/docs/design_docs/physical-planning-and-deployment.md +++ b/docs/design_docs/physical-planning-and-deployment.md @@ -32,10 +32,31 @@ The Logical Post-ASAP DAG is preceded by the Pre-ASAP DAG (`QueryExpr`), the language-independent query semantics before summary selection. Both are logical. Planning builds Post-ASAP `SummaryNode` trees; `compile_post_asap_dag` exports the selected tree as a `PostAsapDag`, which is the Physical Plan -Compiler's input. Its per-node execution phase (ingestion or query time) is an -initial placement: compilation places ingestion-time nodes in the precompute DAG, -while frontier enumeration proposes alternative materialization splits. Which -layer owns placement is an open design question, deferred to a later change. +Compiler's input. Its per-node execution phase (ingestion or query time) is +decided by the selected summary maintenance lifecycle, as the layer contract +below states. + +### Layer contract + +1. **Logical Post-ASAP** (`PlanSpace`) decides what to compute: summary + families, readouts and sharing. It does not decide placement; timing that a + realization strategy writes while building a candidate is provisional. +2. **Summary maintenance lifecycle** (Planner) lists the lifecycle choices for + each unique summary state. A chosen assignment determines every node's + `ExecutionTiming`, plus window framework and retention. + `SummaryMaintenanceLifecyclePlan::execution_timed_dag` applies it: a retained + (non-`Ephemeral`) state and all of its inputs run at ingestion time; + readouts, other consumers, and `Ephemeral` states not consumed by retained + state run at query time. +3. **Physical compile** (Planner) reads timing: ingestion-time nodes form the + precompute DAG and the rest form the query DAG, joined by typed outputs. It + does not see raw ingestion, panes, storage or stored-state readout. +4. **Backend** chooses the lifecycle assignment with its own `CostModel`: + precompute CPU (`maintenance_cost_per_update`), sketch/summary store cost + (`retention_cost_rate`), query reads (`summary_read_cost`) and per-query + builds (`build_cost`, for `Ephemeral`), counting shared state once. + `Ephemeral` requires the deployment to supply the state's raw input as a + query-time source. ### Candidate generation and deployment selection @@ -363,6 +384,9 @@ DAG. If the required behavior cannot be realized, physical compilation fails. Materialization frontiers are Planner decisions. A candidate records both the precompute Physical DAG and the query Physical DAG, with typed outputs connecting them. The deployment compiler binds those outputs; it does not move operators. +Lifecycle timing gives the frontier: ingestion-time nodes read by query-time +nodes. Moving further bounded consumers into precompute, as in Candidate B +below, is not yet expressed as a lifecycle choice. For `sum by(job)(rate(m[1m]))`, legal physical candidates can include: @@ -469,7 +493,7 @@ The complete example makes the ownership boundary explicit: | --- | --- | | **Logical Post-ASAP DAG** | Use `KLL(k=200)` with shared merge for p50/p99 | | **Summary Maintenance Candidate Generation** | Maintain 1-minute panes and reuse them for aligned five-minute queries | -| **Summary Maintenance Lifecycle** | Record pane/window/freshness/reuse requirements | +| **Summary Maintenance Lifecycle** | Record pane/window/freshness/reuse requirements and each node's execution timing | | **Physical Plan Compiler** | Lower to native KLL build, merge, and readout operators | | **Physical DAG** | Define precompute and query DAGs with typed input/output boundaries | | **Deployment Plan Compiler** | Bind raw input and KLL state slots to concrete sources/materializations | @@ -515,6 +539,8 @@ operator/runtime fixtures: | Test | Contract exercised | | --- | --- | | `summary_maintenance_lifecycle_e2e::continuous_lifecycle_compiles_and_executes_spatial_kll` | PromQL workload → selected continuous lifecycle → logical DAG → compiled precompute/query candidate → results in independent revisions; an unbounded candidate fails before pricing, and a bounded request candidate summarizes the same input samples | +| `summary_maintenance_lifecycle_e2e::chosen_lifecycle_timing_decides_precompute_contents` | PromQL workload → enumerated lifecycles → explicit choice → timed DAG → compiled candidate; ContinuouslyMaintained stores the state in precompute, Ephemeral leaves precompute empty and reads the raw source at query time; both return the same p99 | +| `summary_maintenance_lifecycle_e2e::planner_lifecycle_selection_reproduces_strategy_timing` | For PromQL fixtures, the timed DAG from Planner's retained selection equals the DAG realization strategies produce today | | `kll_pane_execution::five_panes_roundtrip_and_shared_merge_runs_once` | Explicit one-minute precompute DAGs → real MessagePack state bytes → five required query inputs → shared native merge → p50/p99; counts every sample once, checks adjacent aligned windows and instruments one merge start per run | | `kll_pane_execution::restored_panes_reject_corruption_parameters_schema_and_missing_binding` | Corrupt bytes, parameter relabelling, incompatible schemas and absent bindings fail explicitly | | `precompute_candidates::grouped_rate_can_be_materialized_before_or_after_grouped_sum` | Cost changes select different legal precompute frontiers; both selected candidates execute with the same reset-sensitive result; uncompilable candidates are not priced |