From b7e4546a40d2562f084a327b978316b54ffaca03 Mon Sep 17 00:00:00 2001 From: zz_y Date: Wed, 26 Aug 2026 13:50:37 -0600 Subject: [PATCH 1/3] feat(asap-aware-mapping): recurrence-aware CSE share-vs-recompute costing (#287) Add a generic recurrence-aware cost context that carries RepeatingEntry intervals and DataCharacteristics-derived ingest rates into the CSE share-vs-recompute decision, without adding scheduling or a runtime execution loop. New `recurrence` module: - UpdateRate/EvaluationRate/CostRate/Horizon: distinct newtypes so a steady-state cost rate (cost units/second) and a one-shot Cost can never be combined except explicitly via `total_cost(rate, horizon, one_shot)` -- enforced by the type system (no `impl Add for CostRate`). - `evaluation_rate_of`: sum(1/interval_i) over repeating consumers' intervals, validating zero/invalid intervals. - `update_rate_from_data_characteristics`: series_count * samples_per_sec_per_series proxy. - `RecurrenceProfile`: aggregated per-target evaluation_rate, one_shot_consumers, update_rate; `is_empty()` is the "no metadata" case that preserves pre-#287 behavior exactly. - `RecurrenceCostExplanation`: selected alternative, both compared cost rates, every input, units, and provenance, for a downstream consumer (e.g. issue #286's DAG-viewer annotations). CostModel trait (cost_model.rs): - New hooks `maintenance_cost_per_update`/`summary_read_cost`/ `raw_recompute_cost`, all with defaults delegating to existing hooks. - `cse_share_decision_with_recurrence`: the recurrence-aware Share/RecomputeIndependently decision. Falls back to `cse_share_decision` (structural consumer_count) when `RecurrenceProfile::is_empty()`; otherwise compares maintained_cost_rate vs recompute_cost_rate, requiring an explicit Horizon whenever one-shot and repeating work are mixed (RecurrenceError::MissingHorizon otherwise). PlanSpace (replacement.rs): - `PlanSpace::recurrence_profiles`: walks every root's reachable sub-DAG (mirroring discover_targets' own traversal) and folds each root's RootRecurrence (Repeating(interval) or OneShot) into every site reachable from it -- so a summary shared by consumers with different intervals gets one profile combining all of them. Keeps `Id` fully opaque (positional, no Eq/Hash/Clone bound needed). Tests cover: mixed intervals, a pure one-shot consumer (with and without an explicit horizon), multiple roots sharing a sub-DAG via a real PlanSpace, invalid/zero intervals, CostRate/Cost unit consistency, a deterministic CostModel showing high evaluation frequency selects Share and low frequency selects RecomputeIndependently, and that update_rate affects only maintained_cost_rate while evaluation_rate affects both. Co-Authored-By: Claude Sonnet 5 --- crates/asap-aware-mapping/src/cost_model.rs | 76 ++ crates/asap-aware-mapping/src/lib.rs | 12 +- crates/asap-aware-mapping/src/recurrence.rs | 1023 ++++++++++++++++++ crates/asap-aware-mapping/src/replacement.rs | 132 +++ 4 files changed, 1240 insertions(+), 3 deletions(-) create mode 100644 crates/asap-aware-mapping/src/recurrence.rs diff --git a/crates/asap-aware-mapping/src/cost_model.rs b/crates/asap-aware-mapping/src/cost_model.rs index 82a3e97..393094a 100644 --- a/crates/asap-aware-mapping/src/cost_model.rs +++ b/crates/asap-aware-mapping/src/cost_model.rs @@ -56,6 +56,9 @@ use asap_types::pre_asap::agg_intent::AggIntent; use asap_types::pre_asap::expr_ir::ColumnRef; use asap_types::pre_asap::query_expr::QueryExpr; +use crate::recurrence::{ + self, Horizon, RecurrenceCostExplanation, RecurrenceError, RecurrenceProfile, +}; use crate::replacement::{ realize_child, Implementation, Replacement, ReplacementSubDAG, TargetSubDAG, }; @@ -340,6 +343,79 @@ pub trait CostModel { } } + // ── Recurrence-aware costing (issue #287) ─────────────────────────── + // + // See `crate::recurrence`'s module docs for the full cost model + // (`maintained_cost_rate`/`recompute_cost_rate` formulas, units, + // provenance of every new input). The three hooks below are the + // per-update-event/per-read/per-recomputation cost primitives that + // formula is built from; `cse_share_decision_with_recurrence` is the + // composed decision, mirroring how `cse_share_decision` above composes + // `cse_recompute_cost`/`cse_shared_maintenance_cost`. + + /// Cost of maintaining `candidate`'s bound summary for a single ingest + /// update event. Units: cost units per update — the + /// `maintenance_cost_per_update` term of `maintained_cost_rate` + /// (`crate::recurrence`). Default: delegates to + /// [`cse_shared_maintenance_cost`](Self::cse_shared_maintenance_cost)'s + /// per-family weight table, reinterpreted as a per-update charge — a + /// deployment with a real measured per-update cost (e.g. observed + /// sketch-insert latency) should override this instead. + fn maintenance_cost_per_update(&self, candidate: &CseCandidate) -> Cost { + self.cse_shared_maintenance_cost(candidate) + } + + /// Cost of one read against `candidate`'s already-maintained summary. + /// Units: cost units per read — the `summary_read_cost` term of + /// `maintained_cost_rate`. Default: `Cost(1.0)`, a nominal unit read — + /// illustrative, like every other numeric default in this trait; a + /// deployment with a real read-path cost should override this. + fn summary_read_cost(&self, _candidate: &CseCandidate) -> Cost { + Cost(1.0) + } + + /// Cost of recomputing `candidate.subtree` once, from the pre-ASAP/raw + /// path. Units: cost units per recomputation — the `raw_recompute_cost` + /// term of `recompute_cost_rate`. Default: delegates to + /// [`cse_recompute_cost`](Self::cse_recompute_cost) (the same + /// structural-size proxy `cse_share_decision` already uses). + fn raw_recompute_cost(&self, candidate: &CseCandidate) -> Cost { + self.cse_recompute_cost(candidate) + } + + /// The recurrence-aware counterpart to + /// [`cse_share_decision`](Self::cse_share_decision): the same + /// `Share`/`RecomputeIndependently` choice, weighted by how *often* + /// `candidate`'s consumers actually run (`recurrence`) instead of only + /// how many structurally exist (`candidate.consumer_count`). See + /// `crate::recurrence`'s module docs for the full design. + /// + /// - `recurrence.is_empty()` (no [`RepeatingEntry`]/[`DataCharacteristics`]-derived + /// metadata available): delegates to + /// [`cse_share_decision`](Self::cse_share_decision), preserving + /// today's structural-consumer-count behavior exactly — issue #287's + /// "preserve existing behavior when recurrence metadata is + /// unavailable" requirement. + /// - Otherwise: compares `maintained_cost_rate` against + /// `recompute_cost_rate` (both cost units/second). If + /// `recurrence.one_shot_consumers > 0` alongside any recurring rate + /// (mixed one-shot + repeating work), `horizon` MUST be `Some` — + /// `Err(RecurrenceError::MissingHorizon)` otherwise, per "the cost + /// model must not silently combine rate-valued and one-shot costs". + /// With no one-shot consumers, `horizon` is optional (comparing bare + /// rates is equivalent to comparing `rate * H` for any fixed `H > 0`). + /// + /// [`RepeatingEntry`]: asap_types::workload::RepeatingEntry + /// [`DataCharacteristics`]: asap_types::workload::DataCharacteristics + fn cse_share_decision_with_recurrence( + &self, + candidate: &CseCandidate, + recurrence: &RecurrenceProfile, + horizon: Option, + ) -> Result { + recurrence::decide(self, candidate, recurrence, horizon) + } + /// Estimate a comparable, numeric cost for one already-constructed /// [`ReplacementSubDAG`] candidate at `target` — a real `f64`, not just a /// relative rank, meant for a caller that wants to *display* "candidate A diff --git a/crates/asap-aware-mapping/src/lib.rs b/crates/asap-aware-mapping/src/lib.rs index cd52612..f12bac5 100644 --- a/crates/asap-aware-mapping/src/lib.rs +++ b/crates/asap-aware-mapping/src/lib.rs @@ -172,6 +172,7 @@ pub mod cost_model; pub mod explanation; pub mod grouping; +pub mod recurrence; pub mod replacement; pub mod rewrite; pub mod rollup; @@ -182,12 +183,17 @@ pub use explanation::{ explain_replacements, explain_replacements_with, ExplanationKind, ReplacementExplanation, }; pub use grouping::{has_subpopulations, HydraGroupingStrategy}; +pub use recurrence::{ + evaluation_rate_of, total_cost, update_rate_from_data_characteristics, CostRate, + EvaluationRate, Horizon, RecurrenceCostExplanation, RecurrenceError, RecurrenceProfile, + RootRecurrence, UpdateRate, +}; pub use replacement::{ default_strategies, default_strategies_with, search_workload, search_workload_with, summary_candidates, GlobalSelection, ImplementError, Implementation, Matcher, MemoGroup, - PlanSpace, RankedGroup, Replacement, ReplacementProvenance, ReplacementStrategy, - ReplacementSubDAG, SelectedGroup, SharedSubtreeStrategy, SketchAlgorithmStrategy, TargetSubDAG, - MAX_SEARCH_ITERATIONS, + PlanSpace, RankedGroup, RecurrenceProfileMap, Replacement, ReplacementProvenance, + ReplacementStrategy, ReplacementSubDAG, SelectedGroup, SharedSubtreeStrategy, + SketchAlgorithmStrategy, TargetSubDAG, MAX_SEARCH_ITERATIONS, }; pub use rewrite::AvgToSumOverCountStrategy; pub use topk_reuse::TopKLimitReuseStrategy; diff --git a/crates/asap-aware-mapping/src/recurrence.rs b/crates/asap-aware-mapping/src/recurrence.rs new file mode 100644 index 0000000..75a9851 --- /dev/null +++ b/crates/asap-aware-mapping/src/recurrence.rs @@ -0,0 +1,1023 @@ +//! Recurrence-aware cost context (issue #287). +//! +//! ASAPPlanner already models recurring-workload metadata +//! ([`asap_types::workload::RepeatingEntry`]) and ingest-rate metadata +//! ([`asap_types::workload::DataCharacteristics`]), but until this module +//! neither reached [`CostModel`]'s CSE share-vs-recompute decision +//! ([`CostModel::cse_share_decision`]): that decision only ever compared a +//! *structural* consumer count (how many workload locations reference a +//! shared subtree) against a flat per-family maintenance weight — it had no +//! notion of how *often* those consumers actually run. +//! +//! This module adds that notion as a generic cost context, not a scheduler: +//! no Prometheus rule-group semantics, no execution loop, no temporal-pane +//! boundary reasoning (see the module's own "Out of scope" list mirrored +//! from the issue). +//! +//! ## Units — every new cost input names its own unit explicitly +//! +//! | Type | Unit | Meaning | +//! |---|---|---| +//! | [`UpdateRate`] | Hz (updates/second) | how often the *raw* data underlying a maintained summary changes (ingest rate) | +//! | [`EvaluationRate`] | Hz (evaluations/second) | how often a target is *read* — `sum(1 / query_interval_i)` over every repeating consumer | +//! | [`CostRate`] | cost units / second | a steady-state cost rate — never comparable to a bare [`Cost`](crate::cost_model::Cost) without going through [`total_cost`] | +//! | [`Horizon`] | seconds | the explicit evaluation window a caller supplies to compare a rate-valued cost against a one-shot cost | +//! +//! `Cost` (bare, from [`crate::cost_model`]) stays a one-time, unitless +//! magnitude — exactly what it was before this module existed, preserved +//! for [`CostModel::cse_share_decision`] and everything else that already +//! uses it. `CostRate` is a *new*, distinct type specifically so a rate and +//! a one-shot cost can never be added directly (no `impl Add for +//! CostRate`, and vice versa) — the compiler enforces the issue's "must not +//! silently combine rate-valued and one-shot costs" requirement; [`total_cost`] +//! is the one sanctioned way to combine them, and it takes an explicit +//! [`Horizon`] to do it. +//! +//! ## Cost semantics (from the issue) +//! +//! For a maintained summary: +//! +//! ```text +//! maintained_cost_rate = +//! update_rate * maintenance_cost_per_update +//! + evaluation_rate * summary_read_cost +//! ``` +//! +//! For recomputation from the pre-ASAP/raw path: +//! +//! ```text +//! recompute_cost_rate = evaluation_rate * raw_recompute_cost +//! ``` +//! +//! For a summary shared by multiple repeating consumers with intervals +//! `t1..tn`: +//! +//! ```text +//! evaluation_rate = sum(1 / query_interval_i) +//! ``` +//! ([`evaluation_rate_of`].) +//! +//! Repetition amortizes a maintained summary across more reads; it does +//! *not* reduce the physical maintenance work caused by ingest updates — +//! that's why `update_rate` and `evaluation_rate` are two separate terms in +//! `maintained_cost_rate` above, never folded into one. +//! +//! One-shot consumers are represented separately from a steady-state rate +//! ([`RecurrenceProfile::one_shot_consumers`]). Comparing a mix of one-shot +//! and repeating work requires an explicit evaluation horizon `H`: +//! +//! ```text +//! total_cost(H) = recurring_cost_rate * H + one_shot_cost +//! ``` +//! ([`total_cost`].) +//! +//! ## Provenance of each new cost input +//! +//! - [`EvaluationRate`]: derived from [`asap_types::workload::RepeatingEntry::interval`] +//! values of every repeating consumer reaching a target (via +//! [`evaluation_rate_of`], or [`crate::replacement::PlanSpace::recurrence_profiles`] +//! for a whole workload). A one-shot ([`asap_types::workload::BatchEntry`]) +//! consumer contributes to [`RecurrenceProfile::one_shot_consumers`] +//! instead, never to this rate. +//! - [`UpdateRate`]: derived from workload-level +//! [`asap_types::workload::DataCharacteristics`] via +//! [`update_rate_from_data_characteristics`] (`series_count * +//! samples_per_sec_per_series`) — a deployment with a more precise +//! per-target ingest measurement should compute its own `UpdateRate` +//! instead of relying on this proxy. +//! - `maintenance_cost_per_update` / `summary_read_cost` / +//! `raw_recompute_cost`: [`CostModel`] hooks (defaults documented on the +//! trait itself, in `cost_model.rs`) — illustrative placeholders, like +//! every other default in that trait; a deployment with real numbers +//! overrides them. +//! +//! ## Preserving existing behavior when recurrence metadata is unavailable +//! +//! [`RecurrenceProfile::is_empty`] is `true` exactly when a caller supplied +//! no [`RepeatingEntry`](asap_types::workload::RepeatingEntry)/ +//! [`DataCharacteristics`](asap_types::workload::DataCharacteristics)-derived +//! information at all (no evaluation rate, no update rate, zero recorded +//! one-shot consumers — [`RecurrenceProfile::EMPTY`], its `Default`). +//! [`CostModel::cse_share_decision_with_recurrence`]'s default body checks +//! this first and, when true, delegates to +//! [`CostModel::cse_share_decision`] byte-for-byte — the existing, +//! structural-consumer-count decision this module never has to touch or +//! second-guess when there's nothing new to feed it. + +use std::fmt; + +use asap_types::workload::{DataCharacteristics, RepetitionInterval}; + +use crate::cost_model::{Cost, CostModel, CseCandidate, ShareDecision}; + +// ── Units ──────────────────────────────────────────────────────────────── + +/// How often the *raw* data underlying a maintained summary changes — +/// ingest/update events per second (Hz). See the module docs' provenance +/// table: normally derived from [`DataCharacteristics`] via +/// [`update_rate_from_data_characteristics`]. +#[derive(Debug, Clone, Copy, PartialEq, PartialOrd)] +pub struct UpdateRate(pub f64); + +/// How often a target is *evaluated* by its consumers — evaluations per +/// second (Hz). For a summary shared by repeating consumers with intervals +/// `t1..tn`, `sum(1 / t_i)` ([`evaluation_rate_of`]); a one-shot consumer +/// never contributes to this rate (see [`RecurrenceProfile::one_shot_consumers`]). +#[derive(Debug, Clone, Copy, PartialEq, PartialOrd)] +pub struct EvaluationRate(pub f64); + +/// A steady-state cost rate: cost units per second. Deliberately a +/// different type from [`Cost`] (a one-time, unitless magnitude) — there is +/// no `impl Add for CostRate` on purpose, so a rate and a one-shot +/// cost can never be combined except explicitly, through [`total_cost`]. +#[derive(Debug, Clone, Copy, PartialEq, PartialOrd)] +pub struct CostRate(pub f64); + +impl CostRate { + /// A cost rate of exactly zero (no ongoing cost at all). + pub const ZERO: CostRate = CostRate(0.0); +} + +impl fmt::Display for CostRate { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!(f, "{}/s", self.0) + } +} + +impl std::ops::Add for CostRate { + type Output = CostRate; + fn add(self, rhs: CostRate) -> CostRate { + CostRate(self.0 + rhs.0) + } +} + +/// An explicit evaluation horizon, in seconds — the only input that lets a +/// [`CostRate`] be combined with a one-shot [`Cost`] (via [`total_cost`]). +#[derive(Debug, Clone, Copy, PartialEq, PartialOrd)] +pub struct Horizon(pub f64); + +/// The one sanctioned way to combine a steady-state [`CostRate`] with a +/// one-shot [`Cost`]: `rate * horizon + one_shot`. There is no other path +/// in this module (or in [`CostModel`]) that adds a `CostRate` to a `Cost` +/// — every other cost value stays in exactly one of the two units, +/// enforced by the type system, satisfying issue #287's "the cost model +/// must not silently combine rate-valued and one-shot costs" requirement. +pub fn total_cost(rate: CostRate, horizon: Horizon, one_shot: Cost) -> Cost { + Cost(rate.0 * horizon.0 + one_shot.0) +} + +// ── Errors ─────────────────────────────────────────────────────────────── + +/// Errors from building or applying recurrence-aware cost inputs. +#[derive(Debug, Clone, Copy, PartialEq, thiserror::Error)] +pub enum RecurrenceError { + /// A [`RepetitionInterval`] of zero (or, in principle, any non-positive + /// value — `RepetitionInterval` is `u32`-backed so only zero is + /// representable) was supplied. A zero interval has no finite rate + /// (`1 / 0`), so it cannot contribute to an [`EvaluationRate`]. + #[error( + "invalid RepetitionInterval({0:?}ms): a repeating query's interval must be > 0 to \ + contribute a finite evaluation rate" + )] + InvalidInterval(RepetitionInterval), + /// A comparison mixed one-shot and repeating work + /// ([`RecurrenceProfile::one_shot_consumers`] > 0 alongside a non-empty + /// [`RecurrenceProfile::evaluation_rate`] or [`RecurrenceProfile::update_rate`]) + /// without an explicit [`Horizon`] to combine them — see [`total_cost`] + /// and the module docs' "Cost semantics" section. + #[error( + "comparison mixes one-shot and repeating work but no evaluation horizon was supplied; \ + an explicit Horizon is required to combine a CostRate with a one-shot Cost (see \ + recurrence::total_cost)" + )] + MissingHorizon, +} + +// ── Aggregation ────────────────────────────────────────────────────────── + +/// `sum(1 / interval_i)`, converted from milliseconds +/// ([`RepetitionInterval`]'s own unit) to Hz, over every repeating +/// consumer's interval. `Ok(None)` when `intervals` is empty — "no +/// repeating consumers observed", distinct from "observed consumers whose +/// rate happens to be zero" (impossible: every valid interval contributes a +/// strictly positive rate). `Err` on the first zero interval encountered. +pub fn evaluation_rate_of(intervals: I) -> Result, RecurrenceError> +where + I: IntoIterator, +{ + let mut total_hz = 0.0; + let mut any = false; + for interval in intervals { + if interval.0 == 0 { + return Err(RecurrenceError::InvalidInterval(interval)); + } + any = true; + // RepetitionInterval is in milliseconds; Hz = 1000 / ms. + total_hz += 1000.0 / f64::from(interval.0); + } + Ok(any.then_some(EvaluationRate(total_hz))) +} + +/// Derive an [`UpdateRate`] from workload-level [`DataCharacteristics`]: +/// `series_count * samples_per_sec_per_series` — the total number of raw +/// ingest samples per second across every series this characteristics +/// value describes. A proxy, not a measurement: a deployment with a more +/// precise per-target ingest rate should compute its own `UpdateRate` +/// rather than rely on this conversion. +pub fn update_rate_from_data_characteristics(dc: &DataCharacteristics) -> UpdateRate { + UpdateRate(dc.series_count as f64 * dc.samples_per_sec_per_series) +} + +// ── RecurrenceProfile ──────────────────────────────────────────────────── + +/// Aggregated recurrence context for one [`CseCandidate`]'s shared target: +/// how fast it's evaluated, how many one-shot consumers reference it +/// separately, and how fast its underlying raw data updates. +/// +/// [`RecurrenceProfile::EMPTY`] (also its `Default`) represents "no +/// recurrence metadata available at all" — the case +/// [`CostModel::cse_share_decision_with_recurrence`]'s default body +/// recognizes via [`is_empty`](Self::is_empty) and treats as "preserve +/// existing (structural) behavior". +#[derive(Debug, Clone, Copy, PartialEq, Default)] +pub struct RecurrenceProfile { + /// `sum(1 / interval_i)` over every repeating consumer of this target, + /// in Hz. `None` when no repeating consumer references it. + pub evaluation_rate: Option, + /// How many one-shot (batch) consumers reference this target, + /// independent of `evaluation_rate` — see the module docs' "Cost + /// semantics" section on why these can't be merged into one rate. + pub one_shot_consumers: usize, + /// The ingest/update rate of the raw data this target (if maintained) + /// would be kept up to date against. `None` when no + /// [`DataCharacteristics`] were available. + pub update_rate: Option, +} + +impl RecurrenceProfile { + /// No recurrence metadata at all: no evaluation rate, no one-shot + /// consumers, no update rate. + pub const EMPTY: RecurrenceProfile = RecurrenceProfile { + evaluation_rate: None, + one_shot_consumers: 0, + update_rate: None, + }; + + /// Whether this profile carries no recurrence information at all — the + /// "missing metadata" case [`CostModel::cse_share_decision_with_recurrence`] + /// falls back on. + pub fn is_empty(&self) -> bool { + self.evaluation_rate.is_none() && self.one_shot_consumers == 0 && self.update_rate.is_none() + } + + /// Build a profile purely from a set of repeating consumers' intervals + /// (no one-shot consumers, no update rate — attach those with + /// [`with_one_shot_consumers`](Self::with_one_shot_consumers)/ + /// [`with_update_rate`](Self::with_update_rate)). + pub fn from_repeating_intervals( + intervals: impl IntoIterator, + ) -> Result { + Ok(Self { + evaluation_rate: evaluation_rate_of(intervals)?, + ..Self::EMPTY + }) + } + + /// Attach a count of one-shot consumers. + pub fn with_one_shot_consumers(mut self, one_shot_consumers: usize) -> Self { + self.one_shot_consumers = one_shot_consumers; + self + } + + /// Attach an ingest/update rate (from [`update_rate_from_data_characteristics`] + /// or a deployment-specific measurement). + pub fn with_update_rate(mut self, update_rate: UpdateRate) -> Self { + self.update_rate = Some(update_rate); + self + } +} + +/// How one workload root recurs — the opaque per-root tag +/// [`crate::replacement::PlanSpace::recurrence_profiles`] threads down to +/// every target reachable from that root. Mirrors +/// [`asap_types::workload::QueryWorkload`]'s own `query_batch` (one-shot) +/// vs. `repeating_queries` (an interval each) split, but at the +/// already-opaque `Id` granularity `search_workload`'s callers already use +/// — this crate needs no more of a caller's own query identity than "which +/// of these two recurrence kinds is this root". +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum RootRecurrence { + /// A one-shot (batch) root — contributes to a reached target's + /// [`RecurrenceProfile::one_shot_consumers`], never to its + /// `evaluation_rate`. + OneShot, + /// A repeating root firing every `RepetitionInterval` — contributes to + /// a reached target's `evaluation_rate` (`1 / interval`, aggregated via + /// [`evaluation_rate_of`]). + Repeating(RepetitionInterval), +} + +// ── Explanation ────────────────────────────────────────────────────────── + +/// The full readout [`CostModel::cse_share_decision_with_recurrence`] +/// returns: which alternative was selected, both compared cost rates +/// (and, when a [`Horizon`] was supplied, both compared totals), every +/// input that went into them, their units, and provenance — meant to be +/// both machine-consumable (e.g. by issue #286's DAG-viewer cost/benefit +/// annotations) and human-readable (via its [`fmt::Display`] impl). +#[derive(Debug, Clone, PartialEq)] +pub struct RecurrenceCostExplanation { + /// The selected alternative. + pub decision: ShareDecision, + /// `maintained_cost_rate` — cost units/second — as defined in the + /// module docs' "Cost semantics" section. + pub maintained_cost_rate: CostRate, + /// `recompute_cost_rate` — cost units/second. + pub recompute_cost_rate: CostRate, + /// `total_cost(horizon)` for the maintained alternative, when `horizon` + /// is `Some`. + pub maintained_total: Option, + /// `total_cost(horizon)` for the recompute alternative, when `horizon` + /// is `Some`. + pub recompute_total: Option, + /// The horizon the totals above were computed over, if any. + pub horizon: Option, + /// The [`UpdateRate`] input used, if any. + pub update_rate: Option, + /// The [`EvaluationRate`] input used, if any. + pub evaluation_rate: Option, + /// The one-shot consumer count input used. + pub one_shot_consumers: usize, + /// `maintenance_cost_per_update` — cost units/update. + pub maintenance_cost_per_update: Cost, + /// `summary_read_cost` — cost units/read. + pub summary_read_cost: Cost, + /// `raw_recompute_cost` — cost units/recomputation. + pub raw_recompute_cost: Cost, + /// Human-readable provenance: which model/path produced this + /// explanation (e.g. "recurrence-aware: " or the + /// structural fallback note when `RecurrenceProfile::is_empty()`). + pub provenance: String, +} + +impl fmt::Display for RecurrenceCostExplanation { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + writeln!(f, "decision: {:?} ({})", self.decision, self.provenance)?; + writeln!( + f, + " maintained_cost_rate = {} (update_rate={:?}Hz * maintenance_cost_per_update={} + \ + evaluation_rate={:?}Hz * summary_read_cost={})", + self.maintained_cost_rate, + self.update_rate.map(|r| r.0), + self.maintenance_cost_per_update, + self.evaluation_rate.map(|r| r.0), + self.summary_read_cost, + )?; + writeln!( + f, + " recompute_cost_rate = {} (evaluation_rate={:?}Hz * raw_recompute_cost={})", + self.recompute_cost_rate, + self.evaluation_rate.map(|r| r.0), + self.raw_recompute_cost, + )?; + writeln!(f, " one_shot_consumers = {}", self.one_shot_consumers)?; + if let Some(h) = self.horizon { + writeln!( + f, + " horizon = {}s; maintained_total = {:?}, recompute_total = {:?}", + h.0, self.maintained_total, self.recompute_total + )?; + } + Ok(()) + } +} + +/// [`CostModel::cse_share_decision_with_recurrence`]'s default body — see +/// that method's own doc for the decision rule; kept as a free function so +/// the logic exists exactly once regardless of how many `CostModel` +/// implementors inherit the default. +pub(crate) fn decide( + cost_model: &C, + candidate: &CseCandidate, + recurrence: &RecurrenceProfile, + horizon: Option, +) -> Result { + if recurrence.is_empty() { + let decision = cost_model.cse_share_decision(candidate); + return Ok(RecurrenceCostExplanation { + decision, + maintained_cost_rate: CostRate::ZERO, + recompute_cost_rate: CostRate::ZERO, + maintained_total: None, + recompute_total: None, + horizon: None, + update_rate: None, + evaluation_rate: None, + one_shot_consumers: 0, + maintenance_cost_per_update: Cost::ZERO, + summary_read_cost: Cost::ZERO, + raw_recompute_cost: Cost::ZERO, + provenance: "structural fallback: no recurrence metadata supplied — delegated to \ + CostModel::cse_share_decision (consumer_count-based), preserving \ + pre-#287 behavior exactly" + .to_string(), + }); + } + + let has_recurring = recurrence.evaluation_rate.is_some() || recurrence.update_rate.is_some(); + if recurrence.one_shot_consumers > 0 && has_recurring && horizon.is_none() { + return Err(RecurrenceError::MissingHorizon); + } + + let update_rate = recurrence.update_rate.map_or(0.0, |r| r.0); + let evaluation_rate = recurrence.evaluation_rate.map_or(0.0, |r| r.0); + + let maintenance_cost_per_update = cost_model.maintenance_cost_per_update(candidate); + let summary_read_cost = cost_model.summary_read_cost(candidate); + let raw_recompute_cost = cost_model.raw_recompute_cost(candidate); + + let maintained_cost_rate = CostRate( + update_rate * maintenance_cost_per_update.0 + evaluation_rate * summary_read_cost.0, + ); + let recompute_cost_rate = CostRate(evaluation_rate * raw_recompute_cost.0); + + // A pure one-shot comparison (no recurring rate at all — `has_recurring` + // is false, so the `MissingHorizon` gate above never fired even though + // `one_shot_consumers > 0`) still needs *some* horizon to run + // `total_cost` through, or its one-shot costs would never enter the + // decision at all. Any positive horizon gives the same ordering here, + // since `maintained_cost_rate`/`recompute_cost_rate` are both exactly + // zero in this case (`rate * H` contributes nothing regardless of `H`) + // — so an implicit `Horizon(1.0)` is exact, not approximate, and this + // is never reached for the genuinely mixed case (that already required + // an explicit `horizon` above). + let effective_horizon = horizon + .or_else(|| (recurrence.one_shot_consumers > 0 && !has_recurring).then_some(Horizon(1.0))); + + let (maintained_total, recompute_total, decision) = if let Some(h) = effective_horizon { + let one_shot_maintained = Cost(summary_read_cost.0 * recurrence.one_shot_consumers as f64); + let one_shot_recompute = Cost(raw_recompute_cost.0 * recurrence.one_shot_consumers as f64); + let maintained = total_cost(maintained_cost_rate, h, one_shot_maintained); + let recompute = total_cost(recompute_cost_rate, h, one_shot_recompute); + let decision = if maintained.0 <= recompute.0 { + ShareDecision::Share + } else { + ShareDecision::RecomputeIndependently + }; + (Some(maintained), Some(recompute), decision) + } else { + // Pure recurring, no one-shot consumers at all: comparing the bare + // rates is exactly equivalent to comparing `rate * H` for any fixed + // `H > 0`, so no horizon is needed to get a correct decision. + let decision = if maintained_cost_rate.0 <= recompute_cost_rate.0 { + ShareDecision::Share + } else { + ShareDecision::RecomputeIndependently + }; + (None, None, decision) + }; + + Ok(RecurrenceCostExplanation { + decision, + maintained_cost_rate, + recompute_cost_rate, + maintained_total, + recompute_total, + horizon: effective_horizon, + update_rate: recurrence.update_rate, + evaluation_rate: recurrence.evaluation_rate, + one_shot_consumers: recurrence.one_shot_consumers, + maintenance_cost_per_update, + summary_read_cost, + raw_recompute_cost, + provenance: "recurrence-aware: maintained_cost_rate vs recompute_cost_rate (cost \ + units/second), per issue #287" + .to_string(), + }) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::cost_model::DefaultCostModel; + + fn interval(ms: u32) -> RepetitionInterval { + RepetitionInterval(ms) + } + + // ── evaluation_rate_of ────────────────────────────────────────────── + + #[test] + fn evaluation_rate_of_empty_is_none() { + assert_eq!(evaluation_rate_of(vec![]).unwrap(), None); + } + + #[test] + fn evaluation_rate_of_single_interval() { + // 1000ms interval => 1 Hz. + let rate = evaluation_rate_of(vec![interval(1000)]).unwrap().unwrap(); + assert!((rate.0 - 1.0).abs() < 1e-9); + } + + #[test] + fn evaluation_rate_of_mixed_intervals_sums_reciprocals() { + // 1s, 10s, 100s intervals => 1 + 0.1 + 0.01 Hz. + let rate = evaluation_rate_of(vec![interval(1_000), interval(10_000), interval(100_000)]) + .unwrap() + .unwrap(); + assert!((rate.0 - 1.11).abs() < 1e-9, "rate={}", rate.0); + } + + #[test] + fn evaluation_rate_of_rejects_zero_interval() { + let err = evaluation_rate_of(vec![interval(1000), interval(0)]).unwrap_err(); + assert_eq!(err, RecurrenceError::InvalidInterval(interval(0))); + } + + // ── update_rate_from_data_characteristics ──────────────────────────── + + #[test] + fn update_rate_from_data_characteristics_multiplies_series_by_sample_rate() { + let dc = DataCharacteristics { + series_count: 1_000, + samples_per_sec_per_series: 0.1, + bytes_per_raw_sample: 80, + distinct_keys_per_window: None, + data_distribution: Default::default(), + }; + let rate = update_rate_from_data_characteristics(&dc); + assert!((rate.0 - 100.0).abs() < 1e-9); + } + + // ── RecurrenceProfile ───────────────────────────────────────────────── + + #[test] + fn empty_profile_is_empty() { + assert!(RecurrenceProfile::EMPTY.is_empty()); + assert!(RecurrenceProfile::default().is_empty()); + } + + #[test] + fn profile_with_any_field_set_is_not_empty() { + assert!(!RecurrenceProfile::EMPTY + .with_one_shot_consumers(1) + .is_empty()); + assert!(!RecurrenceProfile::EMPTY + .with_update_rate(UpdateRate(1.0)) + .is_empty()); + assert!( + !RecurrenceProfile::from_repeating_intervals(vec![interval(1000)]) + .unwrap() + .is_empty() + ); + } + + #[test] + fn from_repeating_intervals_rejects_zero_interval() { + let err = RecurrenceProfile::from_repeating_intervals(vec![interval(0)]).unwrap_err(); + assert_eq!(err, RecurrenceError::InvalidInterval(interval(0))); + } + + // ── total_cost / unit consistency ──────────────────────────────────── + + #[test] + fn total_cost_combines_rate_and_one_shot_explicitly() { + let rate = CostRate(2.0); + let horizon = Horizon(10.0); + let one_shot = Cost(5.0); + assert_eq!(total_cost(rate, horizon, one_shot), Cost(25.0)); + } + + #[test] + fn cost_rate_and_cost_are_distinct_types() { + // Compile-time unit-consistency check: `CostRate` has no `Add` + // or `From` — the only way shown here to combine them is + // `total_cost`, which forces an explicit `Horizon`. (This test's + // real assertion is that the crate compiles at all with no such + // impls; the runtime assertion below just exercises the intended + // combination path.) + let combined = total_cost(CostRate(1.0), Horizon(1.0), Cost(1.0)); + assert_eq!(combined, Cost(2.0)); + } + + // ── decide (structural fallback) ───────────────────────────────────── + + use crate::cost_model::CseCandidate; + use asap_types::post_asap::{ + ExactKind, ExactParams, GroupingStrategy, SummaryExpr, SummaryFamilyType, SummaryField, + SummaryNode, SummarySchema, + }; + use asap_types::pre_asap::expr_ir::ColumnRef; + use asap_types::pre_asap::query_expr::{QueryExpr, Reduction, Source}; + use asap_types::pre_asap::schema::{Column, DataType, Schema}; + use std::rc::Rc; + + fn scan() -> QueryExpr { + QueryExpr::Scan { + source: Source::TimeSeries { metric: "m".into() }, + predicates: vec![], + schema: Schema::with_time_index( + vec![ + Column::new("ts", DataType::Timestamp, false), + Column::new("value", DataType::Float64, false), + ], + 0, + vec![], + ), + } + } + + fn summary_node(family: SummaryFamilyType) -> SummaryNode { + SummaryNode { + expr: SummaryExpr::SummaryAgg { + child: Rc::new(SummaryNode { + expr: SummaryExpr::KeepPreAsap(Rc::new(scan())), + schema: SummarySchema { + fields: vec![], + time_index: None, + }, + }), + family: family.clone(), + col: ColumnRef::Named("value".into()), + reduction: Reduction::by(vec![]), + grouping: GroupingStrategy::default(), + }, + schema: SummarySchema { + fields: vec![SummaryField { + name: "state".into(), + dtype: family, + nullable: false, + }], + time_index: None, + }, + } + } + + #[test] + fn decide_falls_back_to_structural_decision_when_profile_is_empty() { + let subtree = scan(); + let bound = summary_node(SummaryFamilyType::ExactAggregate( + ExactKind::Sum, + ExactParams::Sum, + )); + let candidate = CseCandidate { + subtree: &subtree, + bound_summary: &bound, + consumer_count: 1000, + }; + let explanation = decide( + &DefaultCostModel, + &candidate, + &RecurrenceProfile::EMPTY, + None, + ) + .unwrap(); + assert_eq!( + explanation.decision, + DefaultCostModel.cse_share_decision(&candidate) + ); + assert!(explanation.provenance.contains("structural fallback")); + } + + #[test] + fn decide_rejects_mixed_one_shot_and_repeating_without_horizon() { + let subtree = scan(); + let bound = summary_node(SummaryFamilyType::ExactAggregate( + ExactKind::Sum, + ExactParams::Sum, + )); + let candidate = CseCandidate { + subtree: &subtree, + bound_summary: &bound, + consumer_count: 2, + }; + let profile = RecurrenceProfile::from_repeating_intervals(vec![interval(1000)]) + .unwrap() + .with_one_shot_consumers(1); + let err = decide(&DefaultCostModel, &candidate, &profile, None).unwrap_err(); + assert_eq!(err, RecurrenceError::MissingHorizon); + } + + #[test] + fn decide_accepts_mixed_one_shot_and_repeating_with_an_explicit_horizon() { + let subtree = scan(); + let bound = summary_node(SummaryFamilyType::ExactAggregate( + ExactKind::Sum, + ExactParams::Sum, + )); + let candidate = CseCandidate { + subtree: &subtree, + bound_summary: &bound, + consumer_count: 2, + }; + let profile = RecurrenceProfile::from_repeating_intervals(vec![interval(1000)]) + .unwrap() + .with_one_shot_consumers(1); + let explanation = decide( + &DefaultCostModel, + &candidate, + &profile, + Some(Horizon(3600.0)), + ) + .unwrap(); + assert!(explanation.maintained_total.is_some()); + assert!(explanation.recompute_total.is_some()); + assert_eq!(explanation.horizon, Some(Horizon(3600.0))); + } + + // ── A deterministic cost model exercising the trait hooks directly ── + + struct DeterministicUnitCostModel; + impl CostModel for DeterministicUnitCostModel { + fn rank_candidates( + &self, + _intent: &asap_types::pre_asap::agg_intent::AggIntent, + candidates: &[asap_types::post_asap::SketchAlgorithm], + ) -> Vec { + candidates.to_vec() + } + fn maintenance_cost_per_update(&self, _candidate: &CseCandidate) -> Cost { + Cost(1.0) + } + fn summary_read_cost(&self, _candidate: &CseCandidate) -> Cost { + Cost(1.0) + } + fn raw_recompute_cost(&self, _candidate: &CseCandidate) -> Cost { + Cost(50.0) + } + } + + /// Issue #287's headline acceptance criterion: with a deterministic test + /// cost model and identical IR, a high evaluation frequency selects + /// maintained/shared state while a sufficiently infrequent workload + /// selects recomputation — driven purely by `evaluation_rate`, with a + /// fixed, nonzero `update_rate` representing continuous ingest that + /// keeps a maintained summary's floor cost independent of how often + /// it's read. + #[test] + fn high_frequency_selects_maintained_low_frequency_selects_recompute() { + let subtree = scan(); + let bound = summary_node(SummaryFamilyType::ExactAggregate( + ExactKind::Sum, + ExactParams::Sum, + )); + let candidate = CseCandidate { + subtree: &subtree, + bound_summary: &bound, + consumer_count: 1, + }; + + // A steady 10Hz ingest rate underlies both scenarios — the physical + // maintenance cost is unaffected by how often the summary is read + // (issue #287: "repetition ... does not reduce the physical + // maintenance work caused by ingest updates"). + let update_rate = UpdateRate(10.0); + + // High frequency: a consumer firing every 10ms => 100Hz. + let frequent = RecurrenceProfile::from_repeating_intervals(vec![interval(10)]) + .unwrap() + .with_update_rate(update_rate); + let frequent_explanation = + decide(&DeterministicUnitCostModel, &candidate, &frequent, None).unwrap(); + assert_eq!(frequent_explanation.decision, ShareDecision::Share); + assert!( + frequent_explanation.maintained_cost_rate.0 + < frequent_explanation.recompute_cost_rate.0 + ); + + // Low frequency: a consumer firing every 100s => 0.01Hz. + let infrequent = RecurrenceProfile::from_repeating_intervals(vec![interval(100_000)]) + .unwrap() + .with_update_rate(update_rate); + let infrequent_explanation = + decide(&DeterministicUnitCostModel, &candidate, &infrequent, None).unwrap(); + assert_eq!( + infrequent_explanation.decision, + ShareDecision::RecomputeIndependently + ); + assert!( + infrequent_explanation.maintained_cost_rate.0 + > infrequent_explanation.recompute_cost_rate.0 + ); + } + + /// Update rate feeds only `maintained_cost_rate` (via + /// `maintenance_cost_per_update`), never `recompute_cost_rate` — and + /// evaluation rate feeds both `summary_read_cost` (maintained) and + /// `raw_recompute_cost` (recompute), never bypassing either. Pins the + /// issue's "update rate affects maintained-summary cost but not read + /// frequency; evaluation rate affects summary-read and recomputation + /// cost" acceptance criterion directly against the trait hooks. + #[test] + fn update_rate_only_affects_maintained_cost_evaluation_rate_affects_both() { + let subtree = scan(); + let bound = summary_node(SummaryFamilyType::ExactAggregate( + ExactKind::Sum, + ExactParams::Sum, + )); + let candidate = CseCandidate { + subtree: &subtree, + bound_summary: &bound, + consumer_count: 1, + }; + + let base = RecurrenceProfile::from_repeating_intervals(vec![interval(1000)]).unwrap(); + let with_update = base.with_update_rate(UpdateRate(1000.0)); + + let base_explanation = + decide(&DeterministicUnitCostModel, &candidate, &base, None).unwrap(); + let with_update_explanation = + decide(&DeterministicUnitCostModel, &candidate, &with_update, None).unwrap(); + + // Bumping update_rate alone raises maintained_cost_rate... + assert!( + with_update_explanation.maintained_cost_rate.0 + > base_explanation.maintained_cost_rate.0 + ); + // ...but leaves recompute_cost_rate (a pure function of + // evaluation_rate, unchanged between the two profiles) untouched. + assert_eq!( + with_update_explanation.recompute_cost_rate, + base_explanation.recompute_cost_rate + ); + } + + /// One-shot consumers alone (no repeating consumer, no update rate) — + /// still produce a real decision with no explicit `Horizon` required, + /// by comparing the one-shot costs directly (see `decide`'s + /// `effective_horizon` fallback): a cheap-to-recompute, expensive-to-read + /// target should recompute; an expensive-to-recompute, cheap-to-read one + /// should share. + #[test] + fn one_shot_only_consumer_decides_without_an_explicit_horizon() { + let subtree = scan(); + let bound = summary_node(SummaryFamilyType::ExactAggregate( + ExactKind::Sum, + ExactParams::Sum, + )); + let candidate = CseCandidate { + subtree: &subtree, + bound_summary: &bound, + consumer_count: 1, + }; + let profile = RecurrenceProfile::EMPTY.with_one_shot_consumers(3); + + // raw_recompute_cost=50 >> summary_read_cost=1: sharing wins even + // for purely one-shot consumers, since materializing once and + // reading it 3 times beats recomputing all 3 times independently. + let explanation = decide(&DeterministicUnitCostModel, &candidate, &profile, None).unwrap(); + assert_eq!(explanation.decision, ShareDecision::Share); + assert_eq!(explanation.maintained_total, Some(Cost(3.0))); + assert_eq!(explanation.recompute_total, Some(Cost(150.0))); + } + + // ── multiple roots sharing a sub-DAG, via PlanSpace ────────────────── + + use crate::replacement::search_workload; + use asap_types::pre_asap::agg_intent::AggIntent; + use asap_types::pre_asap::expr_ir::ScalarValue; + use asap_types::pre_asap::query_expr::{Predicate, Reduction as QueryReduction}; + + /// Like `scan()`, plus a "job" label column to group by — CSE's + /// sharing legality gate requires a provable unique key + /// (`Schema::has_unique_key`), and an *ungrouped* aggregate's empty + /// `by` reports none (see `asap_types::pre_asap::cse`'s own "Legality" + /// module docs); grouping by a label column gives `sum_agg()` below a + /// real one, matching the pattern + /// `replacement.rs`'s own CSE fixtures already use (`metric_scan`/`agg` + /// grouped by a label column). + fn labeled_scan() -> QueryExpr { + QueryExpr::Scan { + source: Source::TimeSeries { metric: "m".into() }, + predicates: vec![], + schema: Schema::with_time_index( + vec![ + Column::new("ts", DataType::Timestamp, false), + Column::new("value", DataType::Float64, false), + Column::new("job", DataType::Utf8, true), + ], + 0, + vec![], + ), + } + } + + fn sum_agg() -> QueryExpr { + QueryExpr::Aggregate { + reduction: QueryReduction::by(vec![2]), + measures: vec![AggIntent::Sum { col: Some(1) }], + output_names: vec![], + having: None, + child: Rc::new(labeled_scan()), + } + } + + /// A root wrapping a fresh, independently-built (but structurally + /// identical to every other call's) `sum_agg()` in a `Filter` whose + /// literal predicate is unique per root — keeps the three roots + /// themselves structurally distinct (so they don't collapse into one + /// root the way whole-root-identical fixtures do — see + /// `shared_aggregate_across_two_roots_gets_both_strategies_candidates`'s + /// own doc) while letting `share_common_subtrees` unify their + /// identical `sum_agg()` children onto one shared `Rc`. + fn filtered_root(distinguishing_literal: i64) -> QueryExpr { + QueryExpr::Filter { + pred: Predicate(Rc::new(QueryExpr::Literal(ScalarValue::Int64( + distinguishing_literal, + )))), + child: Rc::new(sum_agg()), + } + } + + /// Three workload roots share one underlying `sum_agg()` sub-DAG: two + /// repeating consumers with different intervals, one one-shot batch + /// consumer. `PlanSpace::recurrence_profiles` must aggregate all three + /// onto the shared sub-DAG's own profile: `evaluation_rate = 1/t1 + + /// 1/t2`, `one_shot_consumers = 1` — issue #287's "support a shared + /// sub-DAG consumed by queries with different intervals" and "multiple + /// roots sharing a sub-DAG" acceptance criteria. + #[test] + fn recurrence_profiles_aggregates_mixed_intervals_across_roots_sharing_a_subdag() { + let roots: Vec<(&str, Rc)> = vec![ + ("root_a", Rc::new(filtered_root(1))), + ("root_b", Rc::new(filtered_root(2))), + ("root_c", Rc::new(filtered_root(3))), + ]; + let space = search_workload(roots); + + // Fixture sanity: the three roots stayed distinct (different + // literal predicates), but their `sum_agg()` children merged onto + // one shared `Rc` (consumer_count 3) — this is the "shared sub-DAG" + // under test. Its own child Scan collapses along with it (all 3 + // Filters' aggregates now point at the *same* Aggregate Rc, so + // there is only ever one Scan Rc underneath, directly referenced + // from exactly one place — the shared Aggregate's own `child`). + assert_eq!( + space.len(), + 5, + "3 distinct Filters + 1 shared Aggregate + 1 Scan underneath it" + ); + let shared_group = space + .groups() + .find(|g| matches!(g.target.as_ref(), QueryExpr::Aggregate { .. })) + .expect("the shared sum_agg() is a discovered target"); + assert_eq!(shared_group.consumer_count, 3, "shared by all 3 roots"); + + let root_recurrence = vec![ + RootRecurrence::Repeating(interval(1_000)), // 1 Hz + RootRecurrence::Repeating(interval(10_000)), // 0.1 Hz + RootRecurrence::OneShot, + ]; + let profiles = space + .recurrence_profiles(&root_recurrence, Some(UpdateRate(5.0))) + .unwrap(); + + let profile = profiles.for_target(&shared_group.target); + let expected_rate = 1.0 / 1.0 + 1.0 / 10.0; // Hz + assert!( + (profile.evaluation_rate.unwrap().0 - expected_rate).abs() < 1e-9, + "evaluation_rate={:?}", + profile.evaluation_rate + ); + assert_eq!(profile.one_shot_consumers, 1); + assert_eq!(profile.update_rate, Some(UpdateRate(5.0))); + + // Each root's own unshared Filter node sees only its own + // contribution — no cross-contamination between sibling roots: the + // 1Hz root's own Filter carries only that 1Hz, not the combined + // rate the shared Aggregate beneath all three carries. + let root_a_profile = profiles.for_target(&space.roots[0].1); + assert!( + (root_a_profile.evaluation_rate.unwrap().0 - 1.0).abs() < 1e-9, + "root_a's own Filter should see only its own 1Hz, not the combined rate: {:?}", + root_a_profile.evaluation_rate + ); + assert_eq!(root_a_profile.one_shot_consumers, 0); + + // The one-shot root's own Filter sees only its one-shot + // contribution, no evaluation rate at all. + let root_c_profile = profiles.for_target(&space.roots[2].1); + assert_eq!(root_c_profile.evaluation_rate, None); + assert_eq!(root_c_profile.one_shot_consumers, 1); + } + + #[test] + fn recurrence_profiles_propagates_an_invalid_interval_error() { + let root = Rc::new(scan()); + let roots: Vec<(&str, Rc)> = vec![("only", root)]; + let space = search_workload(roots); + let err = space + .recurrence_profiles(&[RootRecurrence::Repeating(interval(0))], None) + .unwrap_err(); + assert_eq!(err, RecurrenceError::InvalidInterval(interval(0))); + } + + #[test] + #[should_panic(expected = "must have one entry per root")] + fn recurrence_profiles_asserts_slice_length_matches_root_count() { + let root = Rc::new(scan()); + let roots: Vec<(&str, Rc)> = vec![("only", root)]; + let space = search_workload(roots); + let _ = space.recurrence_profiles(&[], None); + } +} diff --git a/crates/asap-aware-mapping/src/replacement.rs b/crates/asap-aware-mapping/src/replacement.rs index e18a9a5..d33f4e5 100644 --- a/crates/asap-aware-mapping/src/replacement.rs +++ b/crates/asap-aware-mapping/src/replacement.rs @@ -359,11 +359,13 @@ use asap_types::pre_asap::expr_ir::ColumnRef; use asap_types::pre_asap::query_expr::{QueryExpr, QueryExprError, Reduction}; use asap_types::pre_asap::schema::Schema; use asap_types::types::AccuracyTarget; +use asap_types::workload::RepetitionInterval; use std::rc::Rc; use thiserror::Error; use crate::cost_model::{CostModel, CseCandidate, DefaultCostModel, ShareDecision}; use crate::grouping::HydraGroupingStrategy; +use crate::recurrence::{evaluation_rate_of, RecurrenceProfile, RootRecurrence, UpdateRate}; use crate::rollup::RollupStrategy; use crate::topk_reuse::TopKLimitReuseStrategy; @@ -1774,6 +1776,136 @@ impl PlanSpace { } } +// ── Recurrence-aware cost context (issue #287) ────────────────────────── + +/// One [`RecurrenceProfile`] per discovered [`MemoGroup`] target, built by +/// [`PlanSpace::recurrence_profiles`] — the "carry `RepeatingEntry.interval` +/// and relevant `DataCharacteristics` into ASAP-aware search/cost context" +/// half of issue #287. Looked up by `Rc` pointer identity, the same +/// currency [`PlanSpace::group_for`]/[`GlobalSelection::for_target`] already +/// use. +#[derive(Debug, Clone)] +pub struct RecurrenceProfileMap { + profiles: HashMap<*const QueryExpr, RecurrenceProfile>, +} + +impl RecurrenceProfileMap { + /// The [`RecurrenceProfile`] for `target`, or + /// [`RecurrenceProfile::EMPTY`] when `target` wasn't a discovered site + /// in the [`PlanSpace`] this map was built from (or carried no + /// recurring/one-shot/update-rate metadata at all) — always a valid, + /// "no metadata" answer, never a panic. + pub fn for_target(&self, target: &Rc) -> RecurrenceProfile { + self.profiles + .get(&Rc::as_ptr(target)) + .copied() + .unwrap_or(RecurrenceProfile::EMPTY) + } +} + +impl PlanSpace { + /// Build one [`RecurrenceProfile`] per discovered site, by walking every + /// root's whole reachable sub-DAG (the same relational-skeleton + /// traversal [`discover_targets`] itself used to discover those sites) + /// and folding each root's own recurrence tag + /// ([`RootRecurrence::Repeating`]'s interval, or + /// [`RootRecurrence::OneShot`]) into every site reachable from it. + /// + /// `root_recurrence` is positional: `root_recurrence[i]` describes + /// `self.roots[i]` — the same order [`search_workload`]/ + /// [`search_workload_with`] were originally called with (post-CSE + /// dedup preserves both root count and order — see + /// `asap_types::pre_asap::cse::share_common_subtrees`'s own + /// `.map(...).collect()` body). This keeps `Id` fully opaque (no `Eq`/ + /// `Hash`/`Clone` bound needed on it at all — issue #287's "keep + /// caller/query identifiers opaque" requirement) at the cost of the + /// caller keeping the two slices in step; `root_recurrence.len()` must + /// equal `self.roots.len()`. + /// + /// A shared sub-DAG reachable from more than one root aggregates every + /// reaching root's contribution — repeating roots' intervals combine via + /// [`evaluation_rate_of`]'s `sum(1 / interval_i)`, one-shot roots + /// increment [`RecurrenceProfile::one_shot_consumers`] — so a summary + /// consumed by queries with different intervals gets one profile + /// reflecting all of them, per issue #287's "support a shared sub-DAG + /// consumed by queries with different intervals". + /// + /// `update_rate` is applied uniformly to every discovered site: today's + /// [`asap_types::workload::DataCharacteristics`] is a single + /// workload-level value (applies to every query in a `QueryWorkload`), + /// not per-target, so there is no finer-grained source to attach + /// instead. `None` when no `DataCharacteristics` were available — + /// preserves "missing metadata" behavior for the update-rate term alone + /// even when repeating/one-shot consumer information is present. + /// + /// Returns [`RecurrenceError::InvalidInterval`] if any + /// `RootRecurrence::Repeating` interval is zero. + pub fn recurrence_profiles( + &self, + root_recurrence: &[RootRecurrence], + update_rate: Option, + ) -> Result { + assert_eq!( + root_recurrence.len(), + self.roots.len(), + "recurrence_profiles: root_recurrence must have one entry per root, in the same \ + order self.roots is in (got {} entries for {} roots)", + root_recurrence.len(), + self.roots.len() + ); + + let mut intervals: HashMap<*const QueryExpr, Vec> = HashMap::new(); + let mut one_shot_counts: HashMap<*const QueryExpr, usize> = HashMap::new(); + + for ((_, root), recurrence) in self.roots.iter().zip(root_recurrence) { + let mut seen: HashSet<*const QueryExpr> = HashSet::new(); + let mut queue: VecDeque<*const QueryExpr> = VecDeque::new(); + let root_ptr = Rc::as_ptr(root); + seen.insert(root_ptr); + queue.push_back(root_ptr); + + while let Some(ptr) = queue.pop_front() { + match recurrence { + RootRecurrence::Repeating(interval) => { + intervals.entry(ptr).or_default().push(*interval); + } + RootRecurrence::OneShot => { + *one_shot_counts.entry(ptr).or_insert(0) += 1; + } + } + // Every reachable node was itself discovered as its own + // `MemoGroup` (`discover_targets` walks the identical + // relational-skeleton scope) — its own `target` is the + // canonical `Rc` to read children off. + if let Some(group) = self.groups.get(&ptr) { + for (child, _) in direct_child_counts(&group.target) { + if seen.insert(child) { + queue.push_back(child); + } + } + } + } + } + + let mut profiles = HashMap::with_capacity(self.order.len()); + for ptr in &self.order { + let evaluation_rate = + evaluation_rate_of(intervals.get(ptr).cloned().unwrap_or_default())?; + let one_shot_consumers = one_shot_counts.get(ptr).copied().unwrap_or(0); + profiles.insert( + *ptr, + RecurrenceProfile { + evaluation_rate, + one_shot_consumers, + update_rate, + }, + ); + } + + Ok(RecurrenceProfileMap { profiles }) + } +} + /// One [`MemoGroup`]'s candidates, ranked best-first by /// [`PlanSpace::cost_sorted`]. #[derive(Debug)] From 644ca4ccb54615b59254983279275b639fc7d1d9 Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 27 Aug 2026 09:37:00 -0600 Subject: [PATCH 2/3] fix(asap-aware-mapping): correctness fixes from #287 review MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Addresses a code review of PR #295 that found real correctness bugs in the recurrence-aware cost context. All fixes on the same branch. 1. Batch-only workloads no longer silently prefer Share (recurrence.rs): `decide` now charges a one-time `summary_build_cost` (new CostModel hook, defaults to `raw_recompute_cost`) whenever a maintained total is computed, not just `summary_read_cost * one_shot_consumers`. A single one-shot consumer now strictly prefers RecomputeIndependently (build + one read costs more than one direct recompute); many one-shot consumers still amortize the build cost correctly. 2. `update_rate` no longer stamped on sites unreachable from any root (replacement.rs): `PlanSpace::recurrence_profiles` now tracks which sites its BFS actually reaches from a real root and only attaches the caller-supplied `update_rate` to those — a site only ever produced by a Replacement::Rewrite candidate (e.g. AvgToSumOverCountStrategy's invented sum/count sub-DAG) now falls back to RecurrenceProfile::EMPTY instead of getting an ingest-driven maintenance cost charged against a zero evaluation signal. 3. Fixed `DefaultCostModel::maintenance_cost_per_update`'s unit mismatch (cost_model.rs): it no longer reuses `cse_shared_maintenance_cost`'s life-of-the-workload weight (~1-6) as a per-update-event rate multiplier, which made maintained_cost_rate blow up at any realistic update rate. Default is now a small nominal `Cost(0.01)`, independent of that table. 4. `RecurrenceProfileMap` now holds owned `Rc` clones alongside each profile (replacement.rs), not just raw pointers, so it's safe to outlive the `PlanSpace` it was built from — no more risk of a stale profile matching an unrelated node whose allocation reused a freed address. 5. `UpdateRate` is now validated (recurrence.rs): new `validate_update_rate`/`RecurrenceError::InvalidUpdateRate`, applied in `RecurrenceProfile::with_update_rate`, `update_rate_from_data_characteristics`, `PlanSpace::recurrence_profiles`'s own parameter, and as a backstop inside `decide` itself (so a profile built via a direct struct literal can't bypass it). NaN/negative/infinite rates are now rejected instead of silently corrupting the comparison. 6. `PlanSpace::recurrence_profiles`'s root-count mismatch is now a `RecurrenceError::RootCountMismatch`, not a panic — consistent with the method's own `Result`-returning signature. 7. `RecurrenceCostExplanation`'s rate/cost fields (maintained_cost_rate, recompute_cost_rate, maintenance_cost_per_update, summary_read_cost, raw_recompute_cost) are now `Option<...>`, so the structural-fallback path reports "not computed" (`None`) instead of a literal zero that a downstream consumer (e.g. #286's DAG-viewer annotations) could misread as "free". Also fixed (lower priority, time permitting): - Edge multiplicity: a root referencing the same shared subtree twice from one parent (e.g. `BinaryOp{lhs: X, rhs: X}`) now credits X with 2 contributions per repeating root, matching how `MemoGroup::consumer_count` already counts that occurrence — the per-root BFS in `recurrence_profiles` now mirrors `discover_targets`' own "count every occurrence, recurse into children once" pattern instead of a plain reachability set. - `Horizon` is now validated the same way (new `RecurrenceError::InvalidHorizon`): a zero, negative, or non-finite horizon is rejected in `decide` instead of silently distorting (or inverting) the recurring-cost-rate term. - Removed an unnecessary `Vec` clone in the final per-site aggregation loop (iterates `&Vec` directly instead). Not fixed (documented, not attempted): the full O(roots * nodes) per-root BFS in `recurrence_profiles` still walks each root independently rather than a single O(N+E) topological pass — flagged as a real but lower-priority performance concern in the review; a correctness-focused fix pass was judged not the right place to also restructure the traversal's complexity class, and the current approach is still correct and adequately fast for realistic workload sizes. New/updated tests: batch-only-workload (both DeterministicUnitCostModel and DefaultCostModel), single-one-shot-consumer regression, UpdateRate NaN/negative/infinite rejection (builder, update_rate_from_data_characteristics, recurrence_profiles, and the decide() backstop), Horizon zero/negative/NaN/infinite rejection, RootCountMismatch as Err not panic, update_rate-not-stamped-on-an-unreachable-site (via a real AvgToSumOverCountStrategy rewrite fixture), and edge-multiplicity credit via a BinaryOp{lhs: X, rhs: X} fixture. cargo build --workspace, cargo test --workspace, cargo clippy --workspace --all-targets -- -D warnings, and cargo fmt --all -- --check all pass. Co-Authored-By: Claude Sonnet 5 --- crates/asap-aware-mapping/src/cost_model.rs | 50 +- crates/asap-aware-mapping/src/recurrence.rs | 511 ++++++++++++++++--- crates/asap-aware-mapping/src/replacement.rs | 169 ++++-- 3 files changed, 628 insertions(+), 102 deletions(-) diff --git a/crates/asap-aware-mapping/src/cost_model.rs b/crates/asap-aware-mapping/src/cost_model.rs index 393094a..e6c0ae0 100644 --- a/crates/asap-aware-mapping/src/cost_model.rs +++ b/crates/asap-aware-mapping/src/cost_model.rs @@ -356,13 +356,28 @@ pub trait CostModel { /// Cost of maintaining `candidate`'s bound summary for a single ingest /// update event. Units: cost units per update — the /// `maintenance_cost_per_update` term of `maintained_cost_rate` - /// (`crate::recurrence`). Default: delegates to + /// (`crate::recurrence`), where it is multiplied by an `UpdateRate` in + /// **Hz** (`update_rate * maintenance_cost_per_update`). + /// + /// Default: a small nominal constant, `Cost(0.01)` — deliberately + /// **not** derived from /// [`cse_shared_maintenance_cost`](Self::cse_shared_maintenance_cost)'s - /// per-family weight table, reinterpreted as a per-update charge — a - /// deployment with a real measured per-update cost (e.g. observed - /// sketch-insert latency) should override this instead. - fn maintenance_cost_per_update(&self, candidate: &CseCandidate) -> Cost { - self.cse_shared_maintenance_cost(candidate) + /// per-family weight table. That table's values (~1-6) are calibrated + /// against [`cse_recompute_cost`](Self::cse_recompute_cost)'s + /// structural-size proxy for a *life-of-the-workload*, one-time + /// maintenance magnitude — multiplying them by a real ingest rate (even + /// a modest one, e.g. 100 events/s) inflates `maintained_cost_rate` far + /// past any realistic `recompute_cost_rate`, making `Share` + /// unreachable regardless of how infrequently the summary is actually + /// read (issue #287 review). `Cost(0.01)` — one order of magnitude + /// below [`summary_read_cost`](Self::summary_read_cost)'s own nominal + /// default — reflects only that an incremental per-event update is + /// normally far cheaper than a full read or recompute, not a measured + /// ratio; a deployment with a real per-update cost (e.g. observed + /// sketch-insert latency) should override this instead of relying on + /// this placeholder. + fn maintenance_cost_per_update(&self, _candidate: &CseCandidate) -> Cost { + Cost(0.01) } /// Cost of one read against `candidate`'s already-maintained summary. @@ -383,6 +398,29 @@ pub trait CostModel { self.cse_recompute_cost(candidate) } + /// The one-time cost of materializing `candidate`'s bound summary for + /// the *first* time — before any read or ingest-driven update charges + /// anything. Units: cost units (a one-time [`Cost`], not a rate). + /// + /// This is what makes a purely (or mostly) one-shot comparison + /// economically sound: without a build cost, "maintained" looked free + /// to construct, so `Share` won unconditionally for any number of + /// one-shot consumers, no matter how few (issue #287 review, bug 1). + /// With it, a single one-shot consumer never benefits from sharing + /// (build + one read costs more than one direct recompute), while many + /// one-shot consumers still amortize the fixed build cost across their + /// reads, same as before. + /// + /// Default: delegates to + /// [`raw_recompute_cost`](Self::raw_recompute_cost) — materializing a + /// summary for the first time costs about as much as computing its + /// answer once from raw, since there's no delta history yet to apply + /// incrementally. A deployment with a distinct measured "cold build" + /// cost should override this instead. + fn summary_build_cost(&self, candidate: &CseCandidate) -> Cost { + self.raw_recompute_cost(candidate) + } + /// The recurrence-aware counterpart to /// [`cse_share_decision`](Self::cse_share_decision): the same /// `Share`/`RecomputeIndependently` choice, weighted by how *often* diff --git a/crates/asap-aware-mapping/src/recurrence.rs b/crates/asap-aware-mapping/src/recurrence.rs index 75a9851..bf1af7b 100644 --- a/crates/asap-aware-mapping/src/recurrence.rs +++ b/crates/asap-aware-mapping/src/recurrence.rs @@ -191,6 +191,72 @@ pub enum RecurrenceError { recurrence::total_cost)" )] MissingHorizon, + /// An [`UpdateRate`] that isn't finite and non-negative (NaN, infinite, + /// or negative) was supplied — such a value would silently corrupt + /// every downstream comparison (a NaN rate makes every `<=`/`>` + /// comparison `false`, which [`decide`] would otherwise read as "always + /// recompute" with no diagnostic at all). Validated the same way + /// [`evaluation_rate_of`] validates a zero [`RepetitionInterval`]. + #[error( + "invalid UpdateRate({0:?}Hz): an update rate must be finite and >= 0 to contribute a \ + well-defined maintained_cost_rate" + )] + InvalidUpdateRate(UpdateRate), + /// A [`Horizon`] that isn't finite and strictly positive (NaN, + /// infinite, zero, or negative) was supplied — a non-positive or + /// infinite horizon would silently drop or invert the recurring + /// `CostRate` term in [`total_cost`], exactly the "silently combine" + /// outcome [`MissingHorizon`](Self::MissingHorizon) exists to prevent. + #[error( + "invalid Horizon({0:?}s): an evaluation horizon must be finite and > 0 to combine a \ + CostRate with a one-shot Cost without distorting the comparison" + )] + InvalidHorizon(Horizon), + /// [`crate::replacement::PlanSpace::recurrence_profiles`] was called + /// with a `root_recurrence` slice whose length doesn't match the + /// `PlanSpace`'s own root count — a caller error, but recoverable + /// (this method's whole signature promises a `Result`, so this is + /// reported the same way every other input-validation failure is, + /// never a panic). + #[error( + "recurrence_profiles: root_recurrence must have one entry per root, in the same order \ + PlanSpace::roots is in (got {got} entries for {expected} roots)" + )] + RootCountMismatch { + /// `PlanSpace::roots.len()`. + expected: usize, + /// `root_recurrence.len()`. + got: usize, + }, +} + +/// Reject a non-finite or negative [`UpdateRate`] — the same validation +/// discipline [`evaluation_rate_of`] applies to each [`RepetitionInterval`], +/// applied at every point an `UpdateRate` enters a [`RecurrenceProfile`] +/// ([`RecurrenceProfile::with_update_rate`], +/// [`update_rate_from_data_characteristics`], +/// [`crate::replacement::PlanSpace::recurrence_profiles`]'s own parameter) +/// *and*, as a backstop that can't be bypassed by constructing a +/// `RecurrenceProfile` via its public fields directly, inside [`decide`] +/// itself before any comparison uses it. +pub fn validate_update_rate(rate: UpdateRate) -> Result { + if rate.0.is_finite() && rate.0 >= 0.0 { + Ok(rate) + } else { + Err(RecurrenceError::InvalidUpdateRate(rate)) + } +} + +/// Reject a non-finite or non-positive [`Horizon`] — validated inside +/// [`decide`] wherever a caller-supplied `horizon` is used, so `Horizon(0.0)` +/// or a negative horizon can't silently zero out or invert the recurring +/// `CostRate` term in [`total_cost`]. +pub fn validate_horizon(horizon: Horizon) -> Result { + if horizon.0.is_finite() && horizon.0 > 0.0 { + Ok(horizon) + } else { + Err(RecurrenceError::InvalidHorizon(horizon)) + } } // ── Aggregation ────────────────────────────────────────────────────────── @@ -224,8 +290,18 @@ where /// value describes. A proxy, not a measurement: a deployment with a more /// precise per-target ingest rate should compute its own `UpdateRate` /// rather than rely on this conversion. -pub fn update_rate_from_data_characteristics(dc: &DataCharacteristics) -> UpdateRate { - UpdateRate(dc.series_count as f64 * dc.samples_per_sec_per_series) +/// +/// Validated via [`validate_update_rate`]: `samples_per_sec_per_series` is +/// caller-supplied `f64` with no type-level guarantee of being finite or +/// non-negative, so a garbage `DataCharacteristics` value (NaN, infinite, +/// or negative) is rejected here rather than silently propagating into a +/// [`RecurrenceProfile`]. +pub fn update_rate_from_data_characteristics( + dc: &DataCharacteristics, +) -> Result { + validate_update_rate(UpdateRate( + dc.series_count as f64 * dc.samples_per_sec_per_series, + )) } // ── RecurrenceProfile ──────────────────────────────────────────────────── @@ -290,10 +366,12 @@ impl RecurrenceProfile { } /// Attach an ingest/update rate (from [`update_rate_from_data_characteristics`] - /// or a deployment-specific measurement). - pub fn with_update_rate(mut self, update_rate: UpdateRate) -> Self { - self.update_rate = Some(update_rate); - self + /// or a deployment-specific measurement). Validated via + /// [`validate_update_rate`] — rejects a NaN, infinite, or negative rate + /// rather than silently storing it. + pub fn with_update_rate(mut self, update_rate: UpdateRate) -> Result { + self.update_rate = Some(validate_update_rate(update_rate)?); + Ok(self) } } @@ -330,10 +408,16 @@ pub struct RecurrenceCostExplanation { /// The selected alternative. pub decision: ShareDecision, /// `maintained_cost_rate` — cost units/second — as defined in the - /// module docs' "Cost semantics" section. - pub maintained_cost_rate: CostRate, - /// `recompute_cost_rate` — cost units/second. - pub recompute_cost_rate: CostRate, + /// module docs' "Cost semantics" section. `None` on the structural + /// fallback path (`RecurrenceProfile::is_empty()`): no rate was + /// computed at all there (the decision came from + /// [`CostModel::cse_share_decision`] instead), so this is "not + /// computed", deliberately distinct from a real, computed rate of + /// exactly zero. + pub maintained_cost_rate: Option, + /// `recompute_cost_rate` — cost units/second. `None` on the same + /// structural fallback path, for the same reason. + pub recompute_cost_rate: Option, /// `total_cost(horizon)` for the maintained alternative, when `horizon` /// is `Some`. pub maintained_total: Option, @@ -348,12 +432,15 @@ pub struct RecurrenceCostExplanation { pub evaluation_rate: Option, /// The one-shot consumer count input used. pub one_shot_consumers: usize, - /// `maintenance_cost_per_update` — cost units/update. - pub maintenance_cost_per_update: Cost, - /// `summary_read_cost` — cost units/read. - pub summary_read_cost: Cost, - /// `raw_recompute_cost` — cost units/recomputation. - pub raw_recompute_cost: Cost, + /// `maintenance_cost_per_update` — cost units/update. `None` on the + /// structural fallback path (not computed there). + pub maintenance_cost_per_update: Option, + /// `summary_read_cost` — cost units/read. `None` on the structural + /// fallback path. + pub summary_read_cost: Option, + /// `raw_recompute_cost` — cost units/recomputation. `None` on the + /// structural fallback path. + pub raw_recompute_cost: Option, /// Human-readable provenance: which model/path produced this /// explanation (e.g. "recurrence-aware: " or the /// structural fallback note when `RecurrenceProfile::is_empty()`). @@ -363,23 +450,36 @@ pub struct RecurrenceCostExplanation { impl fmt::Display for RecurrenceCostExplanation { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { writeln!(f, "decision: {:?} ({})", self.decision, self.provenance)?; - writeln!( - f, - " maintained_cost_rate = {} (update_rate={:?}Hz * maintenance_cost_per_update={} + \ - evaluation_rate={:?}Hz * summary_read_cost={})", - self.maintained_cost_rate, - self.update_rate.map(|r| r.0), - self.maintenance_cost_per_update, - self.evaluation_rate.map(|r| r.0), - self.summary_read_cost, - )?; - writeln!( - f, - " recompute_cost_rate = {} (evaluation_rate={:?}Hz * raw_recompute_cost={})", - self.recompute_cost_rate, - self.evaluation_rate.map(|r| r.0), - self.raw_recompute_cost, - )?; + match (self.maintained_cost_rate, self.recompute_cost_rate) { + (Some(maintained), Some(recompute)) => { + writeln!( + f, + " maintained_cost_rate = {} (update_rate={:?}Hz * \ + maintenance_cost_per_update={:?} + evaluation_rate={:?}Hz * \ + summary_read_cost={:?})", + maintained, + self.update_rate.map(|r| r.0), + self.maintenance_cost_per_update, + self.evaluation_rate.map(|r| r.0), + self.summary_read_cost, + )?; + writeln!( + f, + " recompute_cost_rate = {} (evaluation_rate={:?}Hz * \ + raw_recompute_cost={:?})", + recompute, + self.evaluation_rate.map(|r| r.0), + self.raw_recompute_cost, + )?; + } + _ => { + writeln!( + f, + " maintained_cost_rate / recompute_cost_rate: not computed (structural \ + fallback — no recurrence metadata was supplied)" + )?; + } + } writeln!(f, " one_shot_consumers = {}", self.one_shot_consumers)?; if let Some(h) = self.horizon { writeln!( @@ -402,21 +502,33 @@ pub(crate) fn decide( recurrence: &RecurrenceProfile, horizon: Option, ) -> Result { + // Backstop validation: a `RecurrenceProfile`/`Horizon` may have reached + // here via a direct struct literal (bypassing `with_update_rate`'s own + // check) or a caller-supplied `horizon` argument — this is the one + // choke point every path funnels through before a value actually enters + // a comparison, so it's validated here regardless of how it arrived. + if let Some(rate) = recurrence.update_rate { + validate_update_rate(rate)?; + } + if let Some(h) = horizon { + validate_horizon(h)?; + } + if recurrence.is_empty() { let decision = cost_model.cse_share_decision(candidate); return Ok(RecurrenceCostExplanation { decision, - maintained_cost_rate: CostRate::ZERO, - recompute_cost_rate: CostRate::ZERO, + maintained_cost_rate: None, + recompute_cost_rate: None, maintained_total: None, recompute_total: None, horizon: None, update_rate: None, evaluation_rate: None, one_shot_consumers: 0, - maintenance_cost_per_update: Cost::ZERO, - summary_read_cost: Cost::ZERO, - raw_recompute_cost: Cost::ZERO, + maintenance_cost_per_update: None, + summary_read_cost: None, + raw_recompute_cost: None, provenance: "structural fallback: no recurrence metadata supplied — delegated to \ CostModel::cse_share_decision (consumer_count-based), preserving \ pre-#287 behavior exactly" @@ -435,6 +547,15 @@ pub(crate) fn decide( let maintenance_cost_per_update = cost_model.maintenance_cost_per_update(candidate); let summary_read_cost = cost_model.summary_read_cost(candidate); let raw_recompute_cost = cost_model.raw_recompute_cost(candidate); + // The one-time cost of materializing the shared summary at all, before + // any read or update — see `CostModel::summary_build_cost`'s own doc. + // Without this term, a purely (or mostly) one-shot comparison modeled + // "maintained" as free to construct, so `Share` won unconditionally + // for *any* number of one-shot consumers (issue #287 review, bug 1) — + // this term is what makes materializing-and-reading actually cost more + // than a single direct recompute for a lone consumer, while still + // amortizing correctly across many. + let summary_build_cost = cost_model.summary_build_cost(candidate); let maintained_cost_rate = CostRate( update_rate * maintenance_cost_per_update.0 + evaluation_rate * summary_read_cost.0, @@ -455,7 +576,11 @@ pub(crate) fn decide( .or_else(|| (recurrence.one_shot_consumers > 0 && !has_recurring).then_some(Horizon(1.0))); let (maintained_total, recompute_total, decision) = if let Some(h) = effective_horizon { - let one_shot_maintained = Cost(summary_read_cost.0 * recurrence.one_shot_consumers as f64); + // `summary_build_cost` is paid exactly once — whether the summary + // is ever read again by a repeating consumer or not — never scaled + // by `one_shot_consumers`. + let one_shot_maintained = + Cost(summary_read_cost.0 * recurrence.one_shot_consumers as f64) + summary_build_cost; let one_shot_recompute = Cost(raw_recompute_cost.0 * recurrence.one_shot_consumers as f64); let maintained = total_cost(maintained_cost_rate, h, one_shot_maintained); let recompute = total_cost(recompute_cost_rate, h, one_shot_recompute); @@ -468,7 +593,12 @@ pub(crate) fn decide( } else { // Pure recurring, no one-shot consumers at all: comparing the bare // rates is exactly equivalent to comparing `rate * H` for any fixed - // `H > 0`, so no horizon is needed to get a correct decision. + // `H > 0` in the limit of a long-lived, continuously-maintained + // summary, so no horizon is needed to get a correct decision — the + // one-time `summary_build_cost` is asymptotically negligible next + // to an ongoing rate term and is deliberately not charged here (it + // only enters the comparison when a caller actually needs an + // absolute total over a finite horizon, via the branch above). let decision = if maintained_cost_rate.0 <= recompute_cost_rate.0 { ShareDecision::Share } else { @@ -479,17 +609,17 @@ pub(crate) fn decide( Ok(RecurrenceCostExplanation { decision, - maintained_cost_rate, - recompute_cost_rate, + maintained_cost_rate: Some(maintained_cost_rate), + recompute_cost_rate: Some(recompute_cost_rate), maintained_total, recompute_total, horizon: effective_horizon, update_rate: recurrence.update_rate, evaluation_rate: recurrence.evaluation_rate, one_shot_consumers: recurrence.one_shot_consumers, - maintenance_cost_per_update, - summary_read_cost, - raw_recompute_cost, + maintenance_cost_per_update: Some(maintenance_cost_per_update), + summary_read_cost: Some(summary_read_cost), + raw_recompute_cost: Some(raw_recompute_cost), provenance: "recurrence-aware: maintained_cost_rate vs recompute_cost_rate (cost \ units/second), per issue #287" .to_string(), @@ -545,10 +675,23 @@ mod tests { distinct_keys_per_window: None, data_distribution: Default::default(), }; - let rate = update_rate_from_data_characteristics(&dc); + let rate = update_rate_from_data_characteristics(&dc).unwrap(); assert!((rate.0 - 100.0).abs() < 1e-9); } + #[test] + fn update_rate_from_data_characteristics_rejects_a_negative_sample_rate() { + let dc = DataCharacteristics { + series_count: 1_000, + samples_per_sec_per_series: -0.1, + bytes_per_raw_sample: 80, + distinct_keys_per_window: None, + data_distribution: Default::default(), + }; + let err = update_rate_from_data_characteristics(&dc).unwrap_err(); + assert!(matches!(err, RecurrenceError::InvalidUpdateRate(_))); + } + // ── RecurrenceProfile ───────────────────────────────────────────────── #[test] @@ -564,6 +707,7 @@ mod tests { .is_empty()); assert!(!RecurrenceProfile::EMPTY .with_update_rate(UpdateRate(1.0)) + .unwrap() .is_empty()); assert!( !RecurrenceProfile::from_repeating_intervals(vec![interval(1000)]) @@ -572,6 +716,23 @@ mod tests { ); } + #[test] + fn with_update_rate_rejects_nan_infinite_and_negative() { + for bad in [f64::NAN, f64::INFINITY, f64::NEG_INFINITY, -1.0] { + let err = RecurrenceProfile::EMPTY + .with_update_rate(UpdateRate(bad)) + .unwrap_err(); + assert!( + matches!(err, RecurrenceError::InvalidUpdateRate(_)), + "bad={bad}, err={err:?}" + ); + } + // Zero is a legitimate (if unusual) update rate: no ingest at all. + assert!(RecurrenceProfile::EMPTY + .with_update_rate(UpdateRate(0.0)) + .is_ok()); + } + #[test] fn from_repeating_intervals_rejects_zero_interval() { let err = RecurrenceProfile::from_repeating_intervals(vec![interval(0)]).unwrap_err(); @@ -776,19 +937,21 @@ mod tests { // High frequency: a consumer firing every 10ms => 100Hz. let frequent = RecurrenceProfile::from_repeating_intervals(vec![interval(10)]) .unwrap() - .with_update_rate(update_rate); + .with_update_rate(update_rate) + .unwrap(); let frequent_explanation = decide(&DeterministicUnitCostModel, &candidate, &frequent, None).unwrap(); assert_eq!(frequent_explanation.decision, ShareDecision::Share); assert!( - frequent_explanation.maintained_cost_rate.0 - < frequent_explanation.recompute_cost_rate.0 + frequent_explanation.maintained_cost_rate.unwrap().0 + < frequent_explanation.recompute_cost_rate.unwrap().0 ); // Low frequency: a consumer firing every 100s => 0.01Hz. let infrequent = RecurrenceProfile::from_repeating_intervals(vec![interval(100_000)]) .unwrap() - .with_update_rate(update_rate); + .with_update_rate(update_rate) + .unwrap(); let infrequent_explanation = decide(&DeterministicUnitCostModel, &candidate, &infrequent, None).unwrap(); assert_eq!( @@ -796,8 +959,8 @@ mod tests { ShareDecision::RecomputeIndependently ); assert!( - infrequent_explanation.maintained_cost_rate.0 - > infrequent_explanation.recompute_cost_rate.0 + infrequent_explanation.maintained_cost_rate.unwrap().0 + > infrequent_explanation.recompute_cost_rate.unwrap().0 ); } @@ -822,7 +985,7 @@ mod tests { }; let base = RecurrenceProfile::from_repeating_intervals(vec![interval(1000)]).unwrap(); - let with_update = base.with_update_rate(UpdateRate(1000.0)); + let with_update = base.with_update_rate(UpdateRate(1000.0)).unwrap(); let base_explanation = decide(&DeterministicUnitCostModel, &candidate, &base, None).unwrap(); @@ -831,8 +994,8 @@ mod tests { // Bumping update_rate alone raises maintained_cost_rate... assert!( - with_update_explanation.maintained_cost_rate.0 - > base_explanation.maintained_cost_rate.0 + with_update_explanation.maintained_cost_rate.unwrap().0 + > base_explanation.maintained_cost_rate.unwrap().0 ); // ...but leaves recompute_cost_rate (a pure function of // evaluation_rate, unchanged between the two profiles) untouched. @@ -845,9 +1008,9 @@ mod tests { /// One-shot consumers alone (no repeating consumer, no update rate) — /// still produce a real decision with no explicit `Horizon` required, /// by comparing the one-shot costs directly (see `decide`'s - /// `effective_horizon` fallback): a cheap-to-recompute, expensive-to-read - /// target should recompute; an expensive-to-recompute, cheap-to-read one - /// should share. + /// `effective_horizon` fallback): with enough one-shot consumers, the + /// fixed `summary_build_cost` amortizes and sharing wins even though + /// `raw_recompute_cost` is expensive. #[test] fn one_shot_only_consumer_decides_without_an_explicit_horizon() { let subtree = scan(); @@ -863,14 +1026,82 @@ mod tests { let profile = RecurrenceProfile::EMPTY.with_one_shot_consumers(3); // raw_recompute_cost=50 >> summary_read_cost=1: sharing wins even - // for purely one-shot consumers, since materializing once and - // reading it 3 times beats recomputing all 3 times independently. + // for purely one-shot consumers, since materializing once (paying + // summary_build_cost=50 exactly once — DeterministicUnitCostModel's + // default delegates build cost to raw_recompute_cost) and reading it + // 3 times (3 * 1 = 3) totals 53, cheaper than recomputing + // independently 3 times (50 * 3 = 150). let explanation = decide(&DeterministicUnitCostModel, &candidate, &profile, None).unwrap(); assert_eq!(explanation.decision, ShareDecision::Share); - assert_eq!(explanation.maintained_total, Some(Cost(3.0))); + assert_eq!(explanation.maintained_total, Some(Cost(53.0))); assert_eq!(explanation.recompute_total, Some(Cost(150.0))); } + /// Regression for issue #287 review bug 1: without a `summary_build_cost` + /// term, a purely one-shot comparison modeled "maintained" as free to + /// construct, so `Share` won unconditionally for *any* number of + /// one-shot consumers — including exactly one, where sharing can never + /// make sense (you always pay at least as much to build-then-read once + /// as you would to just recompute once directly). With the fix, a + /// single one-shot consumer strictly prefers `RecomputeIndependently`. + #[test] + fn one_shot_only_single_consumer_does_not_unconditionally_prefer_share() { + let subtree = scan(); + let bound = summary_node(SummaryFamilyType::ExactAggregate( + ExactKind::Sum, + ExactParams::Sum, + )); + let candidate = CseCandidate { + subtree: &subtree, + bound_summary: &bound, + consumer_count: 1, + }; + let profile = RecurrenceProfile::EMPTY.with_one_shot_consumers(1); + + let explanation = decide(&DeterministicUnitCostModel, &candidate, &profile, None).unwrap(); + // build(50) + read(1) = 51 > recompute(50) * 1 = 50. + assert_eq!(explanation.decision, ShareDecision::RecomputeIndependently); + assert_eq!(explanation.maintained_total, Some(Cost(51.0))); + assert_eq!(explanation.recompute_total, Some(Cost(50.0))); + } + + /// A batch-only workload (no `DataCharacteristics`, only one-shot + /// consumers) with the *default* `DefaultCostModel` must not + /// unconditionally prefer `Share` regardless of how many one-shot + /// consumers there are — issue #287 review bug 1's original repro, + /// pinned against the real default cost model rather than the test's + /// own `DeterministicUnitCostModel`. + #[test] + fn batch_only_workload_does_not_unconditionally_prefer_share_under_default_cost_model() { + let subtree = scan(); + let bound = summary_node(SummaryFamilyType::ExactAggregate( + ExactKind::Sum, + ExactParams::Sum, + )); + let candidate = CseCandidate { + subtree: &subtree, + bound_summary: &bound, + consumer_count: 1, + }; + // `DefaultCostModel`'s `raw_recompute_cost`/`summary_build_cost` + // both delegate to the same structural-size proxy + // (`cse_recompute_cost`), and `summary_read_cost` defaults to a + // nominal `1.0` — build and recompute cost the same, so reading a + // materialized copy even once more than a bare recompute can never + // pay off: for every one-shot-only consumer count, the maintained + // total (`build + read * n`) must be strictly greater than the + // recompute total (`recompute * n`), i.e. `Share` must never win. + for n in [1usize, 2, 10, 1_000] { + let profile = RecurrenceProfile::EMPTY.with_one_shot_consumers(n); + let explanation = decide(&DefaultCostModel, &candidate, &profile, None).unwrap(); + assert_eq!( + explanation.decision, + ShareDecision::RecomputeIndependently, + "n={n}, explanation={explanation:?}" + ); + } + } + // ── multiple roots sharing a sub-DAG, via PlanSpace ────────────────── use crate::replacement::search_workload; @@ -1012,12 +1243,164 @@ mod tests { assert_eq!(err, RecurrenceError::InvalidInterval(interval(0))); } + /// Issue #287 review bug 6: a length mismatch is a recoverable + /// `RecurrenceError`, not a panic — `recurrence_profiles`'s whole + /// signature promises a `Result`. + #[test] + fn recurrence_profiles_reports_a_root_count_mismatch_as_an_error_not_a_panic() { + let root = Rc::new(scan()); + let roots: Vec<(&str, Rc)> = vec![("only", root)]; + let space = search_workload(roots); + let err = space.recurrence_profiles(&[], None).unwrap_err(); + assert_eq!( + err, + RecurrenceError::RootCountMismatch { + expected: 1, + got: 0, + } + ); + } + #[test] - #[should_panic(expected = "must have one entry per root")] - fn recurrence_profiles_asserts_slice_length_matches_root_count() { + fn recurrence_profiles_rejects_an_invalid_update_rate() { let root = Rc::new(scan()); let roots: Vec<(&str, Rc)> = vec![("only", root)]; let space = search_workload(roots); - let _ = space.recurrence_profiles(&[], None); + let err = space + .recurrence_profiles(&[RootRecurrence::OneShot], Some(UpdateRate(f64::NAN))) + .unwrap_err(); + assert!(matches!(err, RecurrenceError::InvalidUpdateRate(_))); + } + + /// Issue #287 review bug 2: a site no root's own structural tree + /// actually reaches must not have the caller-supplied `update_rate` + /// stamped onto it. `AvgToSumOverCountStrategy` (part of + /// `default_strategies`, so included by `search_workload`) is a real, + /// already-shipped source of exactly this shape: it rewrites a bare + /// `avg` `Aggregate` into a *brand new* `Project(sum, count)` sub-DAG — + /// `sum`/`count` are genuinely new `Rc`s, discovered via + /// `discover_new_descendant_targets` from the *candidate's* own + /// children, never reachable by walking the original `avg` root's own + /// structural children (which is just the raw scan). Before the fix, + /// this `count` site would get `{evaluation_rate: None, + /// one_shot_consumers: 0, update_rate: Some(rate)}` — `maintained_cost_rate + /// > 0` against a `recompute_cost_rate` of exactly `0` — unconditionally + /// `RecomputeIndependently`, regardless of the site's own real + /// `consumer_count`. + #[test] + fn recurrence_profiles_does_not_stamp_update_rate_on_a_site_unreachable_from_any_root() { + let avg_root = QueryExpr::Aggregate { + reduction: QueryReduction::by(vec![]), + measures: vec![AggIntent::Avg { col: None }], + output_names: vec![], + having: None, + child: Rc::new(scan()), + }; + let roots: Vec<(&str, Rc)> = vec![("q", Rc::new(avg_root))]; + let space = search_workload(roots); + + let count_group = space + .groups() + .find(|g| { + matches!( + g.target.as_ref(), + QueryExpr::Aggregate { measures, .. } + if measures.iter().any(|m| matches!(m, AggIntent::Count { .. })) + ) + }) + .expect( + "AvgToSumOverCountStrategy should have introduced a new Count aggregate \ + target, unreachable from the original avg root's own structural children", + ); + + let root_recurrence = vec![RootRecurrence::Repeating(interval(1_000))]; + let profiles = space + .recurrence_profiles(&root_recurrence, Some(UpdateRate(5.0))) + .unwrap(); + + let count_profile = profiles.for_target(&count_group.target); + assert_eq!( + count_profile, + RecurrenceProfile::EMPTY, + "a site unreachable from any root's own structural tree must fall back to \ + RecurrenceProfile::EMPTY (no update_rate, no evaluation_rate, no one-shot \ + consumers), not just an evaluation-rate-free profile that still carries the \ + caller's update_rate" + ); + + // The root itself (and the raw scan directly beneath it, which the + // walk *does* reach) still get the real update_rate. + let root_profile = profiles.for_target(&space.roots[0].1); + assert_eq!(root_profile.update_rate, Some(UpdateRate(5.0))); + } + + /// Issue #287 review (lower-priority item): a parent referencing the + /// same shared child twice (`BinaryOp{lhs: X, rhs: X}`, the same shape + /// `pre_asap::cse`'s own within-one-query sharing collapses onto one + /// `Rc`) must credit that child with 2 contributions per repeating + /// root, matching how `MemoGroup::consumer_count` already counts that + /// exact structural occurrence twice — not 1, which a plain + /// reachability-set walk would (wrongly) collapse it to. + #[test] + fn recurrence_profiles_credits_a_direct_repeated_reference_by_its_multiplicity() { + let root = QueryExpr::BinaryOp { + op: asap_types::pre_asap::query_expr::BinaryOpKind::Compare( + asap_types::pre_asap::expr_ir::CompareOpKind::Eq, + ), + lhs: Rc::new(sum_agg()), + rhs: Rc::new(sum_agg()), + vector_match: None, + }; + let space = search_workload(vec![("q", Rc::new(root))]); + + let shared_group = space + .groups() + .find(|g| matches!(g.target.as_ref(), QueryExpr::Aggregate { .. })) + .expect("sum_agg() should merge onto one shared Rc, referenced twice from BinaryOp"); + assert_eq!( + shared_group.consumer_count, 2, + "fixture sanity: referenced twice from the same BinaryOp parent" + ); + + let root_recurrence = vec![RootRecurrence::Repeating(interval(1_000))]; // 1 Hz + let profiles = space.recurrence_profiles(&root_recurrence, None).unwrap(); + let profile = profiles.for_target(&shared_group.target); + + // Referenced twice from the one root: evaluation_rate should be + // 2 * 1Hz = 2Hz, matching consumer_count's own multiplicity — not + // 1Hz, which would undercount by treating "reachable at all" as + // the whole story. + assert!( + (profile.evaluation_rate.unwrap().0 - 2.0).abs() < 1e-9, + "evaluation_rate={:?}", + profile.evaluation_rate + ); + } + + // ── Horizon validation ──────────────────────────────────────────────── + + #[test] + fn decide_rejects_a_zero_or_negative_horizon() { + let subtree = scan(); + let bound = summary_node(SummaryFamilyType::ExactAggregate( + ExactKind::Sum, + ExactParams::Sum, + )); + let candidate = CseCandidate { + subtree: &subtree, + bound_summary: &bound, + consumer_count: 2, + }; + let profile = RecurrenceProfile::from_repeating_intervals(vec![interval(1000)]) + .unwrap() + .with_one_shot_consumers(1); + for bad in [0.0, -1.0, f64::NAN, f64::INFINITY] { + let err = + decide(&DefaultCostModel, &candidate, &profile, Some(Horizon(bad))).unwrap_err(); + assert!( + matches!(err, RecurrenceError::InvalidHorizon(_)), + "bad={bad}, err={err:?}" + ); + } } } diff --git a/crates/asap-aware-mapping/src/replacement.rs b/crates/asap-aware-mapping/src/replacement.rs index d33f4e5..950f7f6 100644 --- a/crates/asap-aware-mapping/src/replacement.rs +++ b/crates/asap-aware-mapping/src/replacement.rs @@ -1784,9 +1784,17 @@ impl PlanSpace { /// half of issue #287. Looked up by `Rc` pointer identity, the same /// currency [`PlanSpace::group_for`]/[`GlobalSelection::for_target`] already /// use. +/// Holds an owned `Rc` clone alongside each profile (not just its +/// raw pointer) so this map keeps every node it describes alive for as long +/// as the map itself lives — a `RecurrenceProfileMap` is safe to outlive the +/// `PlanSpace` it was built from. Without this, a raw `*const QueryExpr` key +/// could, after the originating `PlanSpace` (the only other owner of those +/// `Rc`s) is dropped, collide with an unrelated, later allocation that +/// happens to reuse the same freed address — silently returning a stale +/// profile for the wrong node (issue #287 review, bug 4). #[derive(Debug, Clone)] pub struct RecurrenceProfileMap { - profiles: HashMap<*const QueryExpr, RecurrenceProfile>, + profiles: HashMap<*const QueryExpr, (Rc, RecurrenceProfile)>, } impl RecurrenceProfileMap { @@ -1798,7 +1806,7 @@ impl RecurrenceProfileMap { pub fn for_target(&self, target: &Rc) -> RecurrenceProfile { self.profiles .get(&Rc::as_ptr(target)) - .copied() + .map(|(_, profile)| *profile) .unwrap_or(RecurrenceProfile::EMPTY) } } @@ -1830,7 +1838,9 @@ impl PlanSpace { /// reflecting all of them, per issue #287's "support a shared sub-DAG /// consumed by queries with different intervals". /// - /// `update_rate` is applied uniformly to every discovered site: today's + /// `update_rate` is applied uniformly to every discovered site *that + /// this walk actually reached from some root* (see the "unreachable + /// sites" note below): today's /// [`asap_types::workload::DataCharacteristics`] is a single /// workload-level value (applies to every query in a `QueryWorkload`), /// not per-target, so there is no finer-grained source to attach @@ -1838,48 +1848,100 @@ impl PlanSpace { /// preserves "missing metadata" behavior for the update-rate term alone /// even when repeating/one-shot consumer information is present. /// + /// A parent that structurally references the same child more than once + /// (e.g. `BinaryOp{lhs: X, rhs: X}`) credits that child with one + /// contribution per reference, not one contribution per distinct node — + /// matching how [`MemoGroup::consumer_count`] already counts that exact + /// occurrence twice (issue #287 review: fixes an evaluation-rate + /// undercount for repeated in-query references). This does not + /// propagate multiplicities transitively past one level, the same raw, + /// non-effective scope `MemoGroup::consumer_count` itself has — + /// composing contributions across an *ancestor chain* is + /// [`PlanSpace::global_selection`]'s job, deliberately not wired to + /// this recurrence-aware cost context yet (see this file's own + /// "Whole-plan (cross-group) selection" module docs). + /// + /// **Unreachable sites**: [`PlanSpace`] can contain a site no root's own + /// structural tree actually reaches — e.g. one only ever produced by a + /// [`Replacement::Rewrite`] candidate a [`ReplacementStrategy`] invented + /// (this walk only follows [`MemoGroup::target`]'s own structural + /// children, the same scope [`discover_targets`] uses for the original + /// roots, never a candidate's rewritten value). Such a site gets + /// [`RecurrenceProfile::EMPTY`] — in particular, `update_rate` is + /// **not** stamped onto it — so it falls back to the ordinary + /// structural decision instead of being charged an ingest-driven + /// maintenance cost against a real evaluation/one-shot signal of + /// exactly zero, which previously made `RecomputeIndependently` win + /// there unconditionally, regardless of the site's actual + /// `consumer_count` (issue #287 review, bug 2). + /// /// Returns [`RecurrenceError::InvalidInterval`] if any - /// `RootRecurrence::Repeating` interval is zero. + /// `RootRecurrence::Repeating` interval is zero, + /// [`RecurrenceError::InvalidUpdateRate`] if `update_rate` is non-finite + /// or negative, or [`RecurrenceError::RootCountMismatch`] if + /// `root_recurrence.len() != self.roots.len()`. pub fn recurrence_profiles( &self, root_recurrence: &[RootRecurrence], update_rate: Option, ) -> Result { - assert_eq!( - root_recurrence.len(), - self.roots.len(), - "recurrence_profiles: root_recurrence must have one entry per root, in the same \ - order self.roots is in (got {} entries for {} roots)", - root_recurrence.len(), - self.roots.len() - ); + if root_recurrence.len() != self.roots.len() { + return Err(crate::recurrence::RecurrenceError::RootCountMismatch { + expected: self.roots.len(), + got: root_recurrence.len(), + }); + } + if let Some(rate) = update_rate { + crate::recurrence::validate_update_rate(rate)?; + } let mut intervals: HashMap<*const QueryExpr, Vec> = HashMap::new(); let mut one_shot_counts: HashMap<*const QueryExpr, usize> = HashMap::new(); + // Sites actually reached by at least one root's own recurrence tag + // during the walk below — see this method's own "Unreachable + // sites" doc. + let mut reached: HashSet<*const QueryExpr> = HashSet::new(); for ((_, root), recurrence) in self.roots.iter().zip(root_recurrence) { - let mut seen: HashSet<*const QueryExpr> = HashSet::new(); - let mut queue: VecDeque<*const QueryExpr> = VecDeque::new(); + let recurrence = *recurrence; let root_ptr = Rc::as_ptr(root); - seen.insert(root_ptr); + // Nodes whose children have already been enqueued once for this + // root — mirrors `discover_targets`' own `walk` (count every + // occurrence, recurse into children only the first time) rather + // than a plain reachability set, so `contribute`'s `times` + // (from `direct_child_counts`' own edge multiplicity) is honored + // without also re-expanding — and so double-contributing beyond + // what one root can produce — an already-queued child. + let mut expanded: HashSet<*const QueryExpr> = HashSet::new(); + let mut queue: VecDeque<*const QueryExpr> = VecDeque::new(); + + contribute( + root_ptr, + 1, + recurrence, + &mut intervals, + &mut one_shot_counts, + &mut reached, + ); + expanded.insert(root_ptr); queue.push_back(root_ptr); while let Some(ptr) = queue.pop_front() { - match recurrence { - RootRecurrence::Repeating(interval) => { - intervals.entry(ptr).or_default().push(*interval); - } - RootRecurrence::OneShot => { - *one_shot_counts.entry(ptr).or_insert(0) += 1; - } - } // Every reachable node was itself discovered as its own // `MemoGroup` (`discover_targets` walks the identical // relational-skeleton scope) — its own `target` is the // canonical `Rc` to read children off. if let Some(group) = self.groups.get(&ptr) { - for (child, _) in direct_child_counts(&group.target) { - if seen.insert(child) { + for (child, edge_count) in direct_child_counts(&group.target) { + contribute( + child, + edge_count, + recurrence, + &mut intervals, + &mut one_shot_counts, + &mut reached, + ); + if expanded.insert(child) { queue.push_back(child); } } @@ -1887,18 +1949,30 @@ impl PlanSpace { } } + let empty_intervals: Vec = Vec::new(); let mut profiles = HashMap::with_capacity(self.order.len()); for ptr in &self.order { - let evaluation_rate = - evaluation_rate_of(intervals.get(ptr).cloned().unwrap_or_default())?; + let site_intervals = intervals.get(ptr).unwrap_or(&empty_intervals); + let evaluation_rate = evaluation_rate_of(site_intervals.iter().copied())?; let one_shot_consumers = one_shot_counts.get(ptr).copied().unwrap_or(0); + // Bug 2 fix (see "Unreachable sites" above): only a reached + // site carries the caller-supplied `update_rate`. + let site_update_rate = if reached.contains(ptr) { + update_rate + } else { + None + }; + let node = Rc::clone(&self.groups[ptr].target); profiles.insert( *ptr, - RecurrenceProfile { - evaluation_rate, - one_shot_consumers, - update_rate, - }, + ( + node, + RecurrenceProfile { + evaluation_rate, + one_shot_consumers, + update_rate: site_update_rate, + }, + ), ); } @@ -1906,6 +1980,37 @@ impl PlanSpace { } } +/// Record `times` occurrences of `recurrence` against `ptr` — `times > 1` +/// when a single parent structurally references `ptr` more than once (see +/// [`PlanSpace::recurrence_profiles`]'s own doc on edge multiplicity). +/// A no-op for `times == 0` (an `Rc` returned as a `direct_child_counts` +/// child always has `edge_count >= 1` in practice, but this keeps the +/// helper correct regardless). +fn contribute( + ptr: *const QueryExpr, + times: usize, + recurrence: RootRecurrence, + intervals: &mut HashMap<*const QueryExpr, Vec>, + one_shot_counts: &mut HashMap<*const QueryExpr, usize>, + reached: &mut HashSet<*const QueryExpr>, +) { + if times == 0 { + return; + } + reached.insert(ptr); + match recurrence { + RootRecurrence::Repeating(interval) => { + intervals + .entry(ptr) + .or_default() + .extend(std::iter::repeat_n(interval, times)); + } + RootRecurrence::OneShot => { + *one_shot_counts.entry(ptr).or_insert(0) += times; + } + } +} + /// One [`MemoGroup`]'s candidates, ranked best-first by /// [`PlanSpace::cost_sorted`]. #[derive(Debug)] From b8cb8963d8a911c1df6cfacd206c3fa47f49c936 Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 27 Aug 2026 13:44:25 -0600 Subject: [PATCH 3/3] fix(asap-aware-mapping): wire recurrence into plan selection --- crates/asap-aware-mapping/src/recurrence.rs | 100 +++++++++- crates/asap-aware-mapping/src/replacement.rs | 193 ++++++++++++++----- 2 files changed, 246 insertions(+), 47 deletions(-) diff --git a/crates/asap-aware-mapping/src/recurrence.rs b/crates/asap-aware-mapping/src/recurrence.rs index bf1af7b..f0df734 100644 --- a/crates/asap-aware-mapping/src/recurrence.rs +++ b/crates/asap-aware-mapping/src/recurrence.rs @@ -441,6 +441,9 @@ pub struct RecurrenceCostExplanation { /// `raw_recompute_cost` — cost units/recomputation. `None` on the /// structural fallback path. pub raw_recompute_cost: Option, + /// `summary_build_cost` — one-time cost of materializing the maintained + /// summary. `None` on the structural fallback path. + pub summary_build_cost: Option, /// Human-readable provenance: which model/path produced this /// explanation (e.g. "recurrence-aware: " or the /// structural fallback note when `RecurrenceProfile::is_empty()`). @@ -484,8 +487,9 @@ impl fmt::Display for RecurrenceCostExplanation { if let Some(h) = self.horizon { writeln!( f, - " horizon = {}s; maintained_total = {:?}, recompute_total = {:?}", - h.0, self.maintained_total, self.recompute_total + " horizon = {}s; summary_build_cost = {:?}; maintained_total = {:?}, \ + recompute_total = {:?}", + h.0, self.summary_build_cost, self.maintained_total, self.recompute_total )?; } Ok(()) @@ -529,6 +533,7 @@ pub(crate) fn decide( maintenance_cost_per_update: None, summary_read_cost: None, raw_recompute_cost: None, + summary_build_cost: None, provenance: "structural fallback: no recurrence metadata supplied — delegated to \ CostModel::cse_share_decision (consumer_count-based), preserving \ pre-#287 behavior exactly" @@ -620,6 +625,7 @@ pub(crate) fn decide( maintenance_cost_per_update: Some(maintenance_cost_per_update), summary_read_cost: Some(summary_read_cost), raw_recompute_cost: Some(raw_recompute_cost), + summary_build_cost: Some(summary_build_cost), provenance: "recurrence-aware: maintained_cost_rate vs recompute_cost_rate (cost \ units/second), per issue #287" .to_string(), @@ -883,6 +889,7 @@ mod tests { .unwrap(); assert!(explanation.maintained_total.is_some()); assert!(explanation.recompute_total.is_some()); + assert!(explanation.summary_build_cost.is_some()); assert_eq!(explanation.horizon, Some(Horizon(3600.0))); } @@ -1232,6 +1239,82 @@ mod tests { assert_eq!(root_c_profile.one_shot_consumers, 1); } + #[test] + fn plan_selection_uses_recurrence_profiles_for_cse_choices() { + let roots = vec![ + ("a", Rc::new(filtered_root(1))), + ("b", Rc::new(filtered_root(2))), + ]; + let space = search_workload(roots); + let shared = space + .groups() + .find(|group| matches!(group.target.as_ref(), QueryExpr::Aggregate { .. })) + .expect("the aggregate is shared by both roots"); + let update_rate = Some(UpdateRate(10.0)); + + let frequent = space + .recurrence_profiles( + &[ + RootRecurrence::Repeating(interval(10)), + RootRecurrence::Repeating(interval(10)), + ], + update_rate, + ) + .unwrap(); + let infrequent = space + .recurrence_profiles( + &[ + RootRecurrence::Repeating(interval(100_000)), + RootRecurrence::Repeating(interval(100_000)), + ], + update_rate, + ) + .unwrap(); + + let frequent_ranked = space + .cost_sorted_with_recurrence(&DeterministicUnitCostModel, &frequent, None) + .unwrap(); + let infrequent_ranked = space + .cost_sorted_with_recurrence(&DeterministicUnitCostModel, &infrequent, None) + .unwrap(); + let first_provenance = |ranked: &[crate::replacement::RankedGroup<'_>]| { + ranked + .iter() + .find(|group| Rc::ptr_eq(group.target, &shared.target)) + .and_then(|group| group.candidates.first()) + .map(|candidate| candidate.provenance) + }; + assert_eq!( + first_provenance(&frequent_ranked), + Some(crate::replacement::ReplacementProvenance::CseShare) + ); + assert_eq!( + first_provenance(&infrequent_ranked), + Some(crate::replacement::ReplacementProvenance::CseRecompute) + ); + + let frequent_selected = space + .global_selection_with_recurrence(&DeterministicUnitCostModel, &frequent, None) + .unwrap(); + let infrequent_selected = space + .global_selection_with_recurrence(&DeterministicUnitCostModel, &infrequent, None) + .unwrap(); + assert_eq!( + frequent_selected + .for_target(&shared.target) + .and_then(|group| group.chosen) + .map(|candidate| candidate.provenance), + Some(crate::replacement::ReplacementProvenance::CseShare) + ); + assert_eq!( + infrequent_selected + .for_target(&shared.target) + .and_then(|group| group.chosen) + .map(|candidate| candidate.provenance), + Some(crate::replacement::ReplacementProvenance::CseRecompute) + ); + } + #[test] fn recurrence_profiles_propagates_an_invalid_interval_error() { let root = Rc::new(scan()); @@ -1375,6 +1458,19 @@ mod tests { "evaluation_rate={:?}", profile.evaluation_rate ); + + let scan_group = space + .groups() + .find(|group| matches!(group.target.as_ref(), QueryExpr::Scan { .. })) + .expect("the shared aggregate has a scan descendant"); + assert_eq!( + profiles + .for_target(&scan_group.target) + .evaluation_rate + .unwrap(), + EvaluationRate(2.0), + "ancestor multiplicity must propagate transitively to descendants" + ); } // ── Horizon validation ──────────────────────────────────────────────── diff --git a/crates/asap-aware-mapping/src/replacement.rs b/crates/asap-aware-mapping/src/replacement.rs index 950f7f6..59af4d5 100644 --- a/crates/asap-aware-mapping/src/replacement.rs +++ b/crates/asap-aware-mapping/src/replacement.rs @@ -365,7 +365,9 @@ use thiserror::Error; use crate::cost_model::{CostModel, CseCandidate, DefaultCostModel, ShareDecision}; use crate::grouping::HydraGroupingStrategy; -use crate::recurrence::{evaluation_rate_of, RecurrenceProfile, RootRecurrence, UpdateRate}; +use crate::recurrence::{ + evaluation_rate_of, Horizon, RecurrenceError, RecurrenceProfile, RootRecurrence, UpdateRate, +}; use crate::rollup::RollupStrategy; use crate::topk_reuse::TopKLimitReuseStrategy; @@ -1774,6 +1776,65 @@ impl PlanSpace { }) .collect() } + + /// Recurrence-aware counterpart to [`Self::cost_sorted`]. CSE + /// share/recompute pairs are ordered with the target's recurrence + /// profile; all other candidate shapes retain their existing ranking. + pub fn cost_sorted_with_recurrence( + &self, + cost_model: &dyn CostModel, + profiles: &RecurrenceProfileMap, + horizon: Option, + ) -> Result>, RecurrenceError> { + self.order + .iter() + .map(|ptr| { + let group = &self.groups[ptr]; + let mut candidates = rank_group(group, cost_model); + if cse_candidate_pair(group).is_some() { + if let Some(decision) = decide_group_with_recurrence( + group, + group.consumer_count, + profiles.for_target(&group.target), + horizon, + cost_model, + )? { + candidates.sort_by_key(|candidate| match candidate.provenance { + ReplacementProvenance::CseShare if decision == ShareDecision::Share => { + 0 + } + ReplacementProvenance::CseRecompute + if decision == ShareDecision::RecomputeIndependently => + { + 0 + } + ReplacementProvenance::CseShare + | ReplacementProvenance::CseRecompute => 2, + _ => 1, + }); + } + } + let target = TargetSubDAG::with_consumer_count(&group.target, group.consumer_count); + let costs = candidates + .iter() + .map(|candidate| { + cost_model + .grouping_state_cost(candidate, &target) + .map_or_else( + || cost_model.estimate_cost(candidate, &target), + |cost| cost.0, + ) + }) + .collect(); + Ok(RankedGroup { + target: &group.target, + consumer_count: group.consumer_count, + candidates, + costs, + }) + }) + .collect() + } } // ── Recurrence-aware cost context (issue #287) ────────────────────────── @@ -1851,15 +1912,11 @@ impl PlanSpace { /// A parent that structurally references the same child more than once /// (e.g. `BinaryOp{lhs: X, rhs: X}`) credits that child with one /// contribution per reference, not one contribution per distinct node — - /// matching how [`MemoGroup::consumer_count`] already counts that exact - /// occurrence twice (issue #287 review: fixes an evaluation-rate - /// undercount for repeated in-query references). This does not - /// propagate multiplicities transitively past one level, the same raw, - /// non-effective scope `MemoGroup::consumer_count` itself has — - /// composing contributions across an *ancestor chain* is - /// [`PlanSpace::global_selection`]'s job, deliberately not wired to - /// this recurrence-aware cost context yet (see this file's own - /// "Whole-plan (cross-group) selection" module docs). + /// matching how [`MemoGroup::consumer_count`] counts that occurrence. + /// Multiplicity is propagated through the full descendant path: if the + /// repeated parent is independently evaluated twice, its child is also + /// evaluated twice. This supplies recurrence-aware selection with the + /// effective structural execution rate rather than mere reachability. /// /// **Unreachable sites**: [`PlanSpace`] can contain a site no root's own /// structural tree actually reaches — e.g. one only ever produced by a @@ -1905,45 +1962,35 @@ impl PlanSpace { for ((_, root), recurrence) in self.roots.iter().zip(root_recurrence) { let recurrence = *recurrence; let root_ptr = Rc::as_ptr(root); - // Nodes whose children have already been enqueued once for this - // root — mirrors `discover_targets`' own `walk` (count every - // occurrence, recurse into children only the first time) rather - // than a plain reachability set, so `contribute`'s `times` - // (from `direct_child_counts`' own edge multiplicity) is honored - // without also re-expanding — and so double-contributing beyond - // what one root can produce — an already-queued child. - let mut expanded: HashSet<*const QueryExpr> = HashSet::new(); - let mut queue: VecDeque<*const QueryExpr> = VecDeque::new(); - - contribute( - root_ptr, - 1, - recurrence, - &mut intervals, - &mut one_shot_counts, - &mut reached, - ); - expanded.insert(root_ptr); - queue.push_back(root_ptr); - - while let Some(ptr) = queue.pop_front() { + // Carry path multiplicity transitively. If a shared ancestor is + // referenced twice, every descendant below an independently + // recomputed occurrence is evaluated twice as well; stopping + // expansion after the first pointer visit undercounts exactly + // the effective-consumer rate recurrence-aware costing needs. + let mut queue: VecDeque<(*const QueryExpr, usize)> = VecDeque::new(); + queue.push_back((root_ptr, 1)); + + while let Some((ptr, path_count)) = queue.pop_front() { + contribute( + ptr, + path_count, + recurrence, + &mut intervals, + &mut one_shot_counts, + &mut reached, + ); // Every reachable node was itself discovered as its own // `MemoGroup` (`discover_targets` walks the identical // relational-skeleton scope) — its own `target` is the // canonical `Rc` to read children off. if let Some(group) = self.groups.get(&ptr) { for (child, edge_count) in direct_child_counts(&group.target) { - contribute( + queue.push_back(( child, - edge_count, - recurrence, - &mut intervals, - &mut one_shot_counts, - &mut reached, - ); - if expanded.insert(child) { - queue.push_back(child); - } + path_count + .checked_mul(edge_count) + .expect("query DAG path multiplicity overflowed usize"), + )); } } } @@ -2271,6 +2318,29 @@ impl PlanSpace { /// [`Self::cost_sorted`], whose per-group ranking only ever sees a /// group's own raw [`MemoGroup::consumer_count`]. pub fn global_selection(&self, cost_model: &dyn CostModel) -> GlobalSelection<'_> { + self.global_selection_impl(cost_model, None, None) + .expect("structural global selection cannot produce a recurrence error") + } + + /// Recurrence-aware counterpart to [`Self::global_selection`]. The same + /// whole-plan traversal and effective structural consumer counts are + /// retained, while every CSE share/recompute choice is made from the + /// corresponding recurrence profile. + pub fn global_selection_with_recurrence( + &self, + cost_model: &dyn CostModel, + profiles: &RecurrenceProfileMap, + horizon: Option, + ) -> Result, RecurrenceError> { + self.global_selection_impl(cost_model, Some(profiles), horizon) + } + + fn global_selection_impl( + &self, + cost_model: &dyn CostModel, + profiles: Option<&RecurrenceProfileMap>, + horizon: Option, + ) -> Result, RecurrenceError> { let graph = reference_graph(self); let topo = topological_order(&self.order, &graph); @@ -2285,7 +2355,18 @@ impl PlanSpace { effective_uses.insert(*ptr, effective); let chosen = if effective >= 2 && cse_candidate_pair(group).is_some() { - match decide_with_effective_count(group, effective, cost_model) { + let decision = if let Some(profiles) = profiles { + decide_group_with_recurrence( + group, + effective, + profiles.for_target(&group.target), + horizon, + cost_model, + )? + } else { + decide_with_effective_count(group, effective, cost_model) + }; + match decision { Some(decision) => { let cse = pick_shared_subtree_candidate(group, decision); let effective_target = @@ -2353,10 +2434,10 @@ impl PlanSpace { ); } - GlobalSelection { + Ok(GlobalSelection { order: self.order.clone(), groups, - } + }) } } @@ -2456,6 +2537,28 @@ fn decide_with_effective_count( Some(cost_model.cse_share_decision(&candidate)) } +fn decide_group_with_recurrence( + group: &MemoGroup, + effective_consumer_count: usize, + recurrence: RecurrenceProfile, + horizon: Option, + cost_model: &dyn CostModel, +) -> Result, RecurrenceError> { + let Some(bound) = realize_child(&group.target, cost_model).ok() else { + return Ok(None); + }; + let candidate = CseCandidate { + subtree: &group.target, + bound_summary: &bound, + consumer_count: effective_consumer_count, + }; + Ok(Some( + cost_model + .cse_share_decision_with_recurrence(&candidate, &recurrence, horizon)? + .decision, + )) +} + /// The [`SharedSubtreeStrategy`] candidate matching `decision`: the one /// that shares `group.target`'s own `Rc` for [`ShareDecision::Share`], the /// freshly-allocated one for [`ShareDecision::RecomputeIndependently`] —