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
14 changes: 13 additions & 1 deletion crates/asap-aware-mapping/src/cost_model.rs
Original file line number Diff line number Diff line change
Expand Up @@ -52,7 +52,8 @@ use crate::exact_composition::ExactOperation;
use asap_types::ir::{ASAPOp, Operator, OperatorNode};
use asap_types::post_asap::{
FieldDataType, GroupingStrategy, HydraParams, ResultGuarantee, SketchAlgorithm, SketchParams,
SketchStatistic, SummaryMaintenanceLifecycleGuarantee, SummaryWindowFramework,
SketchStatistic, SummaryMaintenanceLifecycle, SummaryMaintenanceLifecycleGuarantee,
SummaryWindowFramework,
};
use asap_types::pre_asap::agg_intent::AggIntent;
use asap_types::pre_asap::expr_ir::ColumnRef;
Expand Down Expand Up @@ -801,6 +802,17 @@ pub trait CostModel {
SummaryMaintenanceCapabilities::default()
}

/// Estimated response latency for reading `summary` under one selected
/// physical lifecycle. `None` means the deployment has no estimate; it
/// does not make the alternative invalid.
fn summary_read_latency_ms(
&self,
_summary: &OperatorNode,
_lifecycle: &SummaryMaintenanceLifecycle,
) -> Option<f64> {
None
}

/// Replace the sum of selected per-state lifecycle costs with a complete
/// root-DAG cost. The default preserves legacy models. Evidence-strict
/// models return `None` when any root operation is unavailable; callers
Expand Down
33 changes: 32 additions & 1 deletion crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs
Original file line number Diff line number Diff line change
Expand Up @@ -138,6 +138,7 @@ pub enum SummaryMaintenanceLifecycleRejection {
SummaryDoesNotSupportIncrementalUpdates,
SummaryDoesNotSupportDeletion,
MissingCostEvidence,
ExceedsLatencyBound,
}

/// One candidate lifecycle policy for a particular summary deployment.
Expand Down Expand Up @@ -559,6 +560,7 @@ impl SummaryMaintenanceLifecycleCandidates<'_> {
#[derive(Debug)]
struct SummaryMaintenanceWorkloadFacts {
required_accuracy: Vec<AccuracyTarget>,
latency_bound_ms: Option<f64>,
/// Total one-time and recurring reads inside the horizon. `None` means a
/// recurrence or horizon was unknown, not zero reads.
reads: Option<f64>,
Expand Down Expand Up @@ -706,13 +708,35 @@ fn enumerate_with_profile<'a>(
} else {
capabilities
};
let alternatives = alternatives_for(
let mut alternatives = alternatives_for(
&facts,
horizon,
capabilities,
cost_model.summary_maintenance_capabilities(&summary),
cost_model.summary_maintenance_lifecycle_cost_inputs_for_horizon(&summary, horizon),
);
if let Some(bound_ms) = facts.latency_bound_ms {
for alternative in &mut alternatives {
match cost_model.summary_read_latency_ms(
&summary,
&alternative.summary_maintenance_lifecycle,
) {
Some(latency_ms) if latency_ms.is_finite() && latency_ms >= 0.0 => {
if latency_ms > bound_ms && alternative.rejection.is_none() {
alternative.rejection = Some(
SummaryMaintenanceLifecycleRejection::ExceedsLatencyBound,
);
alternative.assumptions.push(format!(
"estimated response latency {latency_ms} ms exceeds {bound_ms} ms bound"
));
}
}
_ => alternative
.assumptions
.push(format!("latency bound {bound_ms} ms unchecked: no estimate")),
}
}
}
SummaryMaintenanceDeployment {
post_asap_node_id: timing_memo
.timed(&summary)
Expand Down Expand Up @@ -1055,6 +1079,7 @@ fn workload_facts(
let mut prepared_eligible = true;
let mut requires_deletion = false;
let mut required_accuracy = Vec::new();
let mut latency_bound_ms: Option<f64> = None;

let entries: Vec<_> = workload.entries().collect();
if workload_entry_indices.is_empty() {
Expand All @@ -1072,6 +1097,11 @@ fn workload_facts(
},
)?;
required_accuracy.push(entry.requirements.accuracy.target());
if let asap_types::workload::LatencyRequirement::ExplicitMaxMs(bound) =
entry.requirements.response_latency
{
latency_bound_ms = Some(latency_bound_ms.map_or(bound, |current| current.min(bound)));
}
requires_deletion |= entry.time_selection.lookback.is_some()
&& entry.time_selection.as_of.is_none()
&& matches!(
Expand Down Expand Up @@ -1186,6 +1216,7 @@ fn workload_facts(
};
Ok(SummaryMaintenanceWorkloadFacts {
required_accuracy,
latency_bound_ms,
reads,
one_time_invocations,
evaluation_rate: has_evaluation_rate.then_some(EvaluationRate(evaluation_rate)),
Expand Down
86 changes: 84 additions & 2 deletions crates/planner/tests/summary_sharing.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,13 +10,13 @@ use asap_aware_mapping::pass::{PlanOutput, PlanningModels};
use asap_aware_mapping::replacement::{default_size_params, DEFAULT_DELTA};
use asap_aware_mapping::{
CostModel, CostRate, DefaultCostModel, Horizon, LifecycleInput, SummaryMaintenanceCapabilities,
SummaryMaintenanceLifecycleCostInputs,
SummaryMaintenanceLifecycleCostInputs, SummaryMaintenanceLifecycleRejection,
};
use asap_frontend_sql::SqlCatalog;
use asap_planner::{e2e_plan, FrontendInput, UserInput};
use asap_types::post_asap::{
AccuracyError, BoundExpr, CompositionOperator, ErrorMetric, ProbabilityExpr, ResultGuarantee,
SketchStatistic,
SketchStatistic, SummaryMaintenanceLifecycle,
};
use asap_types::post_asap::{FieldDataType, SketchAlgorithm, SketchParams};
use asap_types::pre_asap::agg_intent::default_quantile;
Expand All @@ -37,6 +37,7 @@ const HORIZON_S: f64 = 3_600.0;
struct FixedCosts {
build: f64,
raw_per_read: f64,
latency_estimates: bool,
}

impl CostModel for FixedCosts {
Expand Down Expand Up @@ -72,6 +73,17 @@ impl CostModel for FixedCosts {
}
}

fn summary_read_latency_ms(
&self,
_summary: &OperatorNode,
lifecycle: &SummaryMaintenanceLifecycle,
) -> Option<f64> {
self.latency_estimates.then_some(match lifecycle {
SummaryMaintenanceLifecycle::Ephemeral => 250.0,
_ => 50.0,
})
}

fn raw_query_recompute_cost(&self, _target: &OperatorNode) -> Option<Cost> {
Some(Cost(self.raw_per_read))
}
Expand All @@ -82,6 +94,7 @@ impl CostModel for FixedCosts {
const CHEAP_SUMMARY: FixedCosts = FixedCosts {
build: 1.0,
raw_per_read: 1_000.0,
latency_estimates: false,
};

fn requirements(epsilon: f64) -> QueryRequirements {
Expand Down Expand Up @@ -147,6 +160,74 @@ async fn plan_promql(queries: &[(&str, f64)], costs: &FixedCosts) -> PlanOutput
e2e_plan(input).await.expect("workload plans")
}

#[tokio::test]
async fn latency_bound_rejects_slow_ephemeral_summary_and_keeps_fast_maintained_one() {
let mut workload = promql_workload(&[("quantile_over_time(0.99, lat[5m])", 0.01)]);
workload.query_workload.repeating_queries.as_mut().unwrap()[0]
.requirements
.response_latency = LatencyRequirement::ExplicitMaxMs(100.0);
let costs = FixedCosts {
build: 1.0,
raw_per_read: 1_000.0,
latency_estimates: true,
};
let output = e2e_plan(UserInput::new(
&workload,
FrontendInput::Promql {
now_ms: NOW_MS,
histograms: None,
},
PlanningModels::builtin().with_cost(&costs),
lifecycle(),
))
.await
.expect("maintained summary satisfies the response bound");

let deployment = &output.plans[0].plan.deployments[0];
assert!(deployment.alternatives.iter().any(|alternative| {
matches!(
alternative.summary_maintenance_lifecycle,
SummaryMaintenanceLifecycle::Ephemeral
) && alternative.rejection
== Some(SummaryMaintenanceLifecycleRejection::ExceedsLatencyBound)
}));
assert!(deployment.alternatives.iter().any(|alternative| {
matches!(
alternative.summary_maintenance_lifecycle,
SummaryMaintenanceLifecycle::ContinuouslyMaintained
) && alternative.rejection.is_none()
}));
}

#[tokio::test]
async fn missing_latency_estimate_keeps_candidate_and_records_unchecked_reason() {
let mut workload = promql_workload(&[("quantile_over_time(0.99, lat[5m])", 0.01)]);
workload.query_workload.repeating_queries.as_mut().unwrap()[0]
.requirements
.response_latency = LatencyRequirement::ExplicitMaxMs(100.0);
let output = e2e_plan(UserInput::new(
&workload,
FrontendInput::Promql {
now_ms: NOW_MS,
histograms: None,
},
PlanningModels::builtin().with_cost(&CHEAP_SUMMARY),
lifecycle(),
))
.await
.expect("an unchecked bound keeps legal candidates");

assert!(output.plans[0].plan.deployments.iter().any(|deployment| {
deployment.alternatives.iter().any(|alternative| {
alternative.rejection.is_none()
&& alternative
.assumptions
.iter()
.any(|assumption| assumption.contains("latency bound 100 ms unchecked"))
})
}));
}

async fn plan_sql(queries: &[&str], costs: &FixedCosts) -> PlanOutput {
let workload = PlanningWorkload {
query_workload: QueryWorkload {
Expand Down Expand Up @@ -430,6 +511,7 @@ async fn shared_amortization_alone_can_beat_raw_recompute() {
let costs = FixedCosts {
build: 100.0,
raw_per_read: 10.0,
latency_estimates: false,
};
let p50 = ("quantile_over_time(0.5, lat[5m])", 0.01);
let p99 = ("quantile_over_time(0.99, lat[5m])", 0.01);
Expand Down
Loading