diff --git a/docs/rqe_optimizer_TODO.md b/docs/rqe_optimizer_TODO.md index f2ae36c0..13f42e94 100644 --- a/docs/rqe_optimizer_TODO.md +++ b/docs/rqe_optimizer_TODO.md @@ -31,6 +31,11 @@ implemented, what is next, and what is intentionally deferred. 1-hour, 6-hour, and 1-day quantile RQEs to exercise those constraints. - A regression test confirms the MILP matches exhaustive minimum-TCO search on a tiny workload with shared deployment activation costs. +- Retained memory (`(x + max S) / y` instances per active deployment) is + scored, and `milp::minimize_cost` minimizes the hourly price on one EC2 + machine family in fractional instances. Prices are a committed snapshot from + `scripts/fetch_ec2_pricing.py`; `small_problem --milp --machine-family NAME` + uses it. ## Next @@ -43,7 +48,7 @@ implemented, what is next, and what is intentionally deferred. ## Explicitly deferred - Query-result sharing across RQEs. -- Merge buffers, retained-storage capacity, and concurrent-query memory. +- Merge buffers and concurrent-query memory. - RQE churn, replanning, and migration cost. - Precomputed rollups; v1 merges selected base instances at query time. - A policy for choosing one mapping from the reported Pareto frontier. diff --git a/docs/rqe_sketch_deployment_v1.md b/docs/rqe_sketch_deployment_v1.md index 251e8901..cd1b8457 100644 --- a/docs/rqe_sketch_deployment_v1.md +++ b/docs/rqe_sketch_deployment_v1.md @@ -130,14 +130,14 @@ explicit on the RQE; it is never inferred from a metric name. A missing metric does not pass. Before mapping, prune a candidate only if another candidate can serve every -RQE it can and is no worse in query memory, ingest CPU, and latency for each -such RQE. This is safe because any mapping using the removed candidate can +RQE it can and is no worse in query memory, ingest CPU, latency, and retained +memory for each such RQE. This is safe because any mapping using the removed candidate can substitute the remaining one without weakening a modeled objective. Historical instances must be retained long enough to answer the RQEs assigned -to a deployment. Retention is an execution/storage detail in v1, not a -candidate parameter or a scored objective. After selecting a mapping, retain -the history needed by the largest assigned query window. +to a deployment. Retention is not a candidate parameter: a deployment retains +the history needed by the largest assigned query window. Its memory is scored +as retained memory (below) and priced by `minimize_cost`. ## Workload mapping @@ -185,6 +185,20 @@ TCO_cpu = Σ_D ingest_cpu_D × u_D latency_i = Σ_{D: (i,D)∈E} latency_{i,D} × z_{i,D} ``` +`minimize_cost` instead prices the plan on one EC2 machine family `f` (vCPUs +`vcpu_f`, memory `gib_f`, hourly price `price_f`) in fractional instances +`n_f`, with a continuous `R_D ≥ 0` for each candidate's retained GiB: + +```text +minimize price_f × n_f +R_D ≥ retained_gib_{i,D} × z_{i,D} for (i, D) ∈ E +vcpu_f × n_f ≥ TCO_cpu +gib_f × n_f ≥ Σ_D R_D +``` + +Prices come from `rqe-optimizer/data/ec2-pricing-.json`, written by +`scripts/fetch_ec2_pricing.py`. Disk is not priced. + One solve needs a scalar objective, such as minimum `TCO_cpu` subject to `M ≤ memory_budget` and optional `latency_i ≤ latency_budget_i`. To sample the Pareto frontier, repeat solves over memory and latency budgets, or use a @@ -216,6 +230,16 @@ The mapping-level memory objective is the worst individual query: peak_query_memory = max_i query_memory_i ``` +### Retained memory + +A deployment holds `x / y` open instances plus the closed instances still +inside the longest lookback it serves: + +```text +retained_memory = Σ_{D: u_D=1} card(D.labels) × mem_bytes(D.configuration) + × (D.x + max_{i: D(i)=D} S_i) / D.y +``` + ### CPU Ingest CPU is paid once for every active deployment: @@ -281,8 +305,7 @@ breakdown of `TCO_cpu`, together with the selected deployment mapping. - **Latency SLAs:** v1 reports per-RQE latency but does not reject a mapping for exceeding a target. Add optional per-RQE maximum latency as a hard constraint when workloads supply targets. -- **Memory model:** merge buffers, retained-storage capacity, and concurrent - queries are not modeled. +- **Memory model:** merge buffers and concurrent queries are not modeled. - **Static planning:** no RQE churn, replanning, or migration cost. - **Rollups:** v1 does not precompute merged rollups. Queries merge their selected base instances when they run. diff --git a/rqe-optimizer/data/ec2-pricing-2026-10-04.json b/rqe-optimizer/data/ec2-pricing-2026-10-04.json new file mode 100644 index 00000000..49741bbd --- /dev/null +++ b/rqe-optimizer/data/ec2-pricing-2026-10-04.json @@ -0,0 +1,34 @@ +{ + "source": "https://b0.p.awsstatic.com/pricing/2.0/meteredUnitMaps/ec2/USD/current/ec2-ondemand-without-sec-sel/US%20East%20(N.%20Virginia)/Linux/index.json", + "region": "US East (N. Virginia)", + "operating_system": "Linux", + "pricing": "on-demand", + "publication_date": "2026-09-25T17:45:21Z", + "fetched": "2026-10-04", + "families": [ + { + "family": "compute_optimized", + "instance_type": "c7i.xlarge", + "vcpu": 4, + "memory_gib": 8.0, + "usd_per_hour": 0.1785, + "rate_code": "D8WVCVBE32N227KA.JRTCKXETXF.6YS6EN2CT7" + }, + { + "family": "general_purpose", + "instance_type": "m7i.xlarge", + "vcpu": 4, + "memory_gib": 16.0, + "usd_per_hour": 0.2016, + "rate_code": "W8MD4QK5G3J3QYWR.JRTCKXETXF.6YS6EN2CT7" + }, + { + "family": "memory_optimized", + "instance_type": "r7i.xlarge", + "vcpu": 4, + "memory_gib": 32.0, + "usd_per_hour": 0.2646, + "rate_code": "PRGSWU6294SJ3Z93.JRTCKXETXF.6YS6EN2CT7" + } + ] +} diff --git a/rqe-optimizer/examples/small_problem.rs b/rqe-optimizer/examples/small_problem.rs index aeb7cb66..30564fce 100644 --- a/rqe-optimizer/examples/small_problem.rs +++ b/rqe-optimizer/examples/small_problem.rs @@ -22,7 +22,9 @@ //! mappings by default; pass `--progress-every N` to change that interval or //! `--print-first N` to display example mappings. `--milp` solves the //! minimum-TCO model without enumerating mappings. Repeat -//! `--latency-limit RQE_ID=SECONDS` to impose MILP latency bounds. +//! `--latency-limit RQE_ID=SECONDS` to impose MILP latency bounds. Add +//! `--machine-family NAME` (`compute_optimized`, `general_purpose`, +//! `memory_optimized`) to minimize that EC2 family's hourly price instead. //! `--sample-mappings N` prints N feasible mappings and exits. use std::collections::BTreeSet; @@ -32,8 +34,8 @@ use rqe_optimizer::candidates::{ build_all_candidates, build_all_candidates_unpruned, eligible_deployments_for, }; use rqe_optimizer::enumerate::{brute_force, for_each_mapping, for_each_mapping_while, unservable}; -use rqe_optimizer::milp::{minimize_tco, MilpBounds}; -use rqe_optimizer::objectives::score; +use rqe_optimizer::milp::{minimize_cost, minimize_tco, MilpBounds}; +use rqe_optimizer::objectives::{score, MachineFamily}; use rqe_optimizer::pareto::{pareto_front, ParetoFront}; use rqe_optimizer::{ AccuracyDirection, AtomicCostTable, Capability, LabelSet, LabelSetInfo, LabelSetTable, Rqe, @@ -61,6 +63,22 @@ const CARDINALITY_ERR: &str = "relative_error"; const TOPK_PRECISION: &str = "precision_at_k"; const COST_TABLE_PATH: &str = "out/rqe_atomic_costs.json"; +const EC2_PRICING: &str = include_str!("../data/ec2-pricing-2026-10-04.json"); + +fn machine_family() -> Option { + let args: Vec<_> = std::env::args().collect(); + let name = args + .windows(2) + .find(|pair| pair[0] == "--machine-family") + .map(|pair| pair[1].clone())?; + let families = MachineFamily::from_pricing_json(EC2_PRICING).expect("committed snapshot"); + Some( + families + .into_iter() + .find(|family| family.family == name) + .unwrap_or_else(|| panic!("unknown --machine-family: {name}")), + ) +} fn label_set(names: &[&str]) -> LabelSet { names.iter().map(|s| s.to_string()).collect() @@ -329,19 +347,28 @@ fn main() { } if std::env::args().any(|arg| arg == "--milp") { - let latency_bounds = latency_bounds(&rqes); - let solution = minimize_tco( - &rqes, - &deployments, - &label_sets, - &MilpBounds { - max_peak_query_memory_bytes: None, - max_query_latency_secs: latency_bounds, - }, - ) + let bounds = MilpBounds { + max_peak_query_memory_bytes: None, + max_query_latency_secs: latency_bounds(&rqes), + }; + let family = machine_family(); + let solution = match &family { + Some(family) => minimize_cost(&rqes, &deployments, &label_sets, &bounds, family), + None => minimize_tco(&rqes, &deployments, &label_sets, &bounds), + } .expect("small_problem MILP should be feasible"); + if let Some(family) = &family { + println!( + "MILP minimum-cost solution on {}: ${:.4}/hour, {:.4} instances, \ + retained_mem={:.0}MB", + family.family, + family.usd_per_hour(&solution.objectives), + family.instances(&solution.objectives), + solution.objectives.retained_memory_bytes / 1e6, + ); + } println!( - "MILP minimum-TCO solution: peak_query_mem={:.0}MB, ingest={:.3e}, \ + "MILP solution: peak_query_mem={:.0}MB, ingest={:.3e}, \ merge={:.3e}, query={:.3e}, total={:.3e} cpu-sec/sec", solution.objectives.peak_query_memory_bytes / 1e6, solution.objectives.ingest_cpu_secs_per_sec, diff --git a/rqe-optimizer/src/candidates.rs b/rqe-optimizer/src/candidates.rs index 22278a55..f93a66e9 100644 --- a/rqe-optimizer/src/candidates.rs +++ b/rqe-optimizer/src/candidates.rs @@ -101,8 +101,8 @@ pub fn build_all_candidates_unpruned(rqes: &[Rqe], costs: &[AtomicCostEntry]) -> /// This comparison is deliberately local to a capability/label-set group. /// Within such a group, label cardinality and arrival rate are common /// multipliers, so comparing per-instance query memory and -/// `active_instances * insert_cost` is sufficient. Query latency is checked -/// for each RQE the dominated candidate can serve. +/// `active_instances * insert_cost` is sufficient. Query latency and retained +/// memory are checked for each RQE the dominated candidate can serve. pub fn prune_dominated_candidates(rqes: &[Rqe], candidates: Vec) -> Vec { let eligibility: Vec> = candidates .iter() @@ -166,8 +166,10 @@ fn candidate_dominates( .enumerate() .all(|(rqe_index, &serves)| { !serves - || query_latency(replacement, &rqes[rqe_index]) + || (query_latency(replacement, &rqes[rqe_index]) <= query_latency(original, &rqes[rqe_index]) + && retained_memory(replacement, &rqes[rqe_index]) + <= retained_memory(original, &rqes[rqe_index])) }) } @@ -193,8 +195,10 @@ fn strictly_better( .enumerate() .any(|(rqe_index, &serves)| { serves - && query_latency(replacement, &rqes[rqe_index]) + && (query_latency(replacement, &rqes[rqe_index]) < query_latency(original, &rqes[rqe_index]) + || retained_memory(replacement, &rqes[rqe_index]) + < retained_memory(original, &rqes[rqe_index])) }) } @@ -206,6 +210,15 @@ fn query_latency(deployment: &Deployment, rqe: &Rqe) -> f64 { + instances.saturating_sub(1) as f64 * deployment.config.merge_cpu_secs } +/// Retained bytes per label group when serving `rqe`. Label cardinality is a +/// common multiplier within a group, so it is left out. +fn retained_memory(deployment: &Deployment, rqe: &Rqe) -> f64 { + deployment + .retained_instance_count(rqe.lookback_secs) + .expect("candidate coverage only contains exactly tiled RQEs") as f64 + * deployment.config.mem_bytes_per_instance +} + pub fn eligible_deployments_for(r: &Rqe, deployments: &[Deployment]) -> Vec { deployments .iter() @@ -334,4 +347,36 @@ mod tests { assert_eq!(retained, vec![large_window]); } + + #[test] + fn retains_candidate_with_lower_retained_memory() { + let r = rqe("r", 60, 60); + // Holds (60 + 60) / 60 = 2 instances of 10 bytes. + let whole_window = Deployment { + capability: Capability::Freq, + labels: LabelSet::new(), + config: AtomicCostEntry { + mem_bytes_per_instance: 10.0, + ..cost() + }, + window_secs: 60, + slide_secs: 60, + }; + // Smaller, faster sketches, but (20 + 60) / 20 = 4 of them: 24 bytes. + let panes = Deployment { + config: AtomicCostEntry { + mem_bytes_per_instance: 6.0, + query_cpu_secs: 0.5, + merge_cpu_secs: 0.1, + ..cost() + }, + window_secs: 20, + slide_secs: 20, + ..whole_window.clone() + }; + + let retained = prune_dominated_candidates(&[r], vec![whole_window.clone(), panes.clone()]); + + assert_eq!(retained, vec![whole_window, panes]); + } } diff --git a/rqe-optimizer/src/lib.rs b/rqe-optimizer/src/lib.rs index 09045520..7910c7be 100644 --- a/rqe-optimizer/src/lib.rs +++ b/rqe-optimizer/src/lib.rs @@ -152,6 +152,14 @@ impl Deployment { } Some(lookback_secs / self.window_secs) } + + /// Instances held to serve a lookback: `x / y` open ones plus the closed + /// ones still inside the lookback, `(x + S) / y`. + pub fn retained_instance_count(&self, lookback_secs: Seconds) -> Option { + self.query_instance_count(lookback_secs)?; + self.active_instance_count()?; + Some((self.window_secs + lookback_secs) / self.slide_secs) + } } /// A full mapping: one deployment index (into the `deployments` slice passed diff --git a/rqe-optimizer/src/milp.rs b/rqe-optimizer/src/milp.rs index 9ac6cd99..b07eabf1 100644 --- a/rqe-optimizer/src/milp.rs +++ b/rqe-optimizer/src/milp.rs @@ -4,7 +4,7 @@ //! enumerating the Cartesian product of eligible deployment choices. use crate::candidates::eligible_deployments_for; -use crate::objectives::{score, Objectives}; +use crate::objectives::{score, MachineFamily, Objectives, BYTES_PER_GIB}; use crate::{Deployment, LabelSetTable, Mapping, Rqe}; use good_lp::{ default_solver, variable, Expression, ProblemVariables, ResolutionError, Solution, SolverModel, @@ -33,6 +33,29 @@ pub fn minimize_tco( deployments: &[Deployment], label_sets: &LabelSetTable, bounds: &MilpBounds, +) -> Result { + solve(rqes, deployments, label_sets, bounds, None) +} + +/// Minimize the hourly price of running the plan on `family`, in fractional +/// instances: `n ≥ CPU / vCPU` and `n ≥ retained memory / GiB`. Same inputs +/// and bounds as [`minimize_tco`]. +pub fn minimize_cost( + rqes: &[Rqe], + deployments: &[Deployment], + label_sets: &LabelSetTable, + bounds: &MilpBounds, + family: &MachineFamily, +) -> Result { + solve(rqes, deployments, label_sets, bounds, Some(family)) +} + +fn solve( + rqes: &[Rqe], + deployments: &[Deployment], + label_sets: &LabelSetTable, + bounds: &MilpBounds, + family: Option<&MachineFamily>, ) -> Result { assert!( bounds.max_query_latency_secs.is_empty() @@ -54,7 +77,6 @@ pub fn minimize_tco( .iter() .map(|_| variables.add(variable().binary())) .collect(); - let peak_memory = variables.add(variable().min(0)); let assignments: Vec> = eligible .iter() .map(|choices| { @@ -65,33 +87,93 @@ pub fn minimize_tco( }) .collect(); - let mut objective = Expression::with_capacity(deployments.len() + rqes.len()); - for (deployment, &active_variable) in deployments.iter().zip(&active) { - let labels = &label_sets[&deployment.labels]; + // Real plans cost ~1e-6 $/hour and read in microseconds, below HiGHS's + // absolute gap (1e-6) and feasibility tolerance (1e-7). So every row is + // scaled to O(1) by `reference`, a plan's resource demand without sharing + // (each RQE on its cheapest deployment); per-RQE bounds become exclusions. Results are re-scored on + // real costs below, so the scaling never leaks out. + let ingest = |deployment: &Deployment| { let fanout = deployment .active_instance_count() .expect("candidate window and slide align") as f64; - objective.add_mul( - labels.arrival_rate_per_sec * fanout * deployment.config.insert_cpu_secs, - active_variable, - ); + label_sets[&deployment.labels].arrival_rate_per_sec + * fanout + * deployment.config.insert_cpu_secs + }; + let retained = |rqe: &Rqe, deployment: &Deployment| { + label_sets[&deployment.labels].cardinality as f64 + * deployment.config.mem_bytes_per_instance + * deployment + .retained_instance_count(rqe.lookback_secs) + .expect("assignment only contains eligible deployments") as f64 + / BYTES_PER_GIB + }; + // Per resource, in the units the solve minimizes: CPU for TCO, and for a + // family the vCPU and GiB fractions of one instance. + let (cpu_unit, gib_unit) = match family { + Some(f) => (f.vcpu, f.memory_gib), + None => (1.0, f64::INFINITY), + }; + let reference: f64 = rqes + .iter() + .zip(&eligible) + .map(|(rqe, choices)| { + choices + .iter() + .map(|&d| { + let deployment = &deployments[d]; + let cpu = ingest(deployment) + + query_and_merge_cpu_per_sec(rqe, deployment, label_sets); + cpu / cpu_unit + retained(rqe, deployment) / gib_unit + }) + .fold(f64::INFINITY, f64::min) + }) + .sum(); + let reference = if reference.is_normal() { + reference + } else { + 1.0 + }; + + let mut cpu = Expression::with_capacity(deployments.len() + rqes.len()); + for (deployment, &active_variable) in deployments.iter().zip(&active) { + cpu.add_mul(ingest(deployment) / reference, active_variable); } for (rqe_index, choices) in assignments.iter().enumerate() { for &(deployment_index, assignment) in choices { - objective.add_mul( + cpu.add_mul( query_and_merge_cpu_per_sec( &rqes[rqe_index], &deployments[deployment_index], label_sets, - ), + ) / reference, assignment, ); } } - let mut model = variables.minimise(objective).using(default_solver); - if let Some(memory_limit) = bounds.max_peak_query_memory_bytes { - model.add_constraint(Expression::from(peak_memory).leq(memory_limit)); + // Retained GiB per deployment, and fractional instances, priced only when + // a machine family is given. Both are in units of `reference`. + let retained_gib: Vec = match family { + Some(_) => deployments + .iter() + .map(|_| variables.add(variable().min(0))) + .collect(), + None => Vec::new(), + }; + let instances = variables.add(variable().min(0)); + // The price per instance is a positive constant, so minimizing instances + // minimizes the price. + let goal = match family { + Some(_) => Expression::from(instances), + None => cpu.clone(), + }; + + let mut model = variables.minimise(goal).using(default_solver); + if let Some(f) = family { + model.add_constraint((cpu / f.vcpu - instances).leq(0)); + let total_gib: Expression = retained_gib.iter().sum(); + model.add_constraint((total_gib / f.memory_gib - instances).leq(0)); } for (rqe_index, choices) in assignments.iter().enumerate() { @@ -101,20 +183,31 @@ pub fn minimize_tco( for &(deployment_index, assignment) in choices { let deployment = &deployments[deployment_index]; model.add_constraint((assignment - active[deployment_index]).leq(0)); + // Each RQE takes exactly one deployment, so a per-RQE bound just + // forbids the choices over it. Comparing in f64 here, not in the + // solver, keeps µs latencies clear of its feasibility tolerance. let memory = label_sets[&deployment.labels].cardinality as f64 * deployment.config.mem_bytes_per_instance; - model.add_constraint((memory * assignment - peak_memory).leq(0)); - } - - if let Some(Some(latency_limit)) = bounds.max_query_latency_secs.get(rqe_index) { - let latency: Expression = choices - .iter() - .map(|&(deployment_index, assignment)| { - query_latency_secs(&rqes[rqe_index], &deployments[deployment_index], label_sets) - * assignment - }) - .sum(); - model.add_constraint(latency.leq(*latency_limit)); + let latency = query_latency_secs(&rqes[rqe_index], deployment, label_sets); + if bounds + .max_peak_query_memory_bytes + .is_some_and(|limit| memory > limit) + || bounds + .max_query_latency_secs + .get(rqe_index) + .copied() + .flatten() + .is_some_and(|limit| latency > limit) + { + model.add_constraint(Expression::from(assignment).leq(0)); + } + if family.is_some() { + // Linearized max over the RQEs a deployment serves. + let retained = retained(&rqes[rqe_index], deployment) / reference; + model.add_constraint( + (retained * assignment - retained_gib[deployment_index]).leq(0), + ); + } } } @@ -267,4 +360,156 @@ mod tests { (milp.objectives.tco_cpu_secs_per_sec - expected.tco_cpu_secs_per_sec).abs() < 1e-12 ); } + + fn family(memory_gib: f64) -> MachineFamily { + MachineFamily { + family: format!("{memory_gib} GiB"), + vcpu: 4.0, + memory_gib, + usd_per_hour: 1.0, + } + } + + #[test] + fn minimize_cost_matches_brute_force_minimum() { + let mut frequent = rqe(); + frequent.id = "frequent".into(); + let mut long = rqe(); + long.id = "long".into(); + long.lookback_secs = 600; + let rqes = vec![frequent, long]; + let deployments = vec![ + deployment(1.0, 0.5 * BYTES_PER_GIB, 400.0), + deployment(3.0, 0.1 * BYTES_PER_GIB, 1.0), + deployment(0.2, 2.0 * BYTES_PER_GIB, 1.0), + ]; + let label_sets = BTreeMap::from([( + LabelSet::new(), + LabelSetInfo { + cardinality: 1, + arrival_rate_per_sec: 1.0, + }, + )]); + + for family in [family(8.0), family(32.0)] { + let best = brute_force(&rqes, &deployments) + .iter() + .map(|mapping| { + family.usd_per_hour(&score(&rqes, &deployments, mapping, &label_sets)) + }) + .min_by(f64::total_cmp) + .expect("test workload is servable"); + let milp = minimize_cost( + &rqes, + &deployments, + &label_sets, + &MilpBounds::default(), + &family, + ) + .expect("feasible MILP"); + assert!((family.usd_per_hour(&milp.objectives) - best).abs() < 1e-9); + } + } + + #[test] + fn binding_resource_decides_the_plan_per_family() { + let rqes = vec![rqe()]; + // Holds 2 instances: CPU-heavy with 2 GiB, or CPU-light with 4 GiB. + let deployments = vec![ + deployment(1.0, 1.0 * BYTES_PER_GIB, 0.0), + deployment(0.05, 2.0 * BYTES_PER_GIB, 0.0), + ]; + let label_sets = BTreeMap::from([( + LabelSet::new(), + LabelSetInfo { + cardinality: 1, + arrival_rate_per_sec: 1.0, + }, + )]); + let solve = |memory_gib| { + minimize_cost( + &rqes, + &deployments, + &label_sets, + &MilpBounds::default(), + &family(memory_gib), + ) + .expect("feasible MILP") + .mapping + }; + + assert_eq!(solve(8.0), vec![0]); // memory is scarce: 0.25 vs 0.5 instances + assert_eq!(solve(32.0), vec![1]); // CPU is scarce: 0.25 vs 0.125 instances + } + + #[test] + fn tiny_magnitudes_match_brute_force_and_respect_latency_bounds() { + // Real W2 plans cost ~1e-6 $/hour with µs latencies: below HiGHS's + // absolute gap and feasibility tolerances unless the model is scaled. + let mut frequent = rqe(); + frequent.id = "frequent".into(); + let mut long = rqe(); + long.id = "long".into(); + long.lookback_secs = 600; + let rqes = vec![frequent, long]; + let tiny = 1e-9; + let deployments = vec![ + deployment(1.0 * tiny, 0.5 * tiny * BYTES_PER_GIB, 400.0 * tiny), + deployment(3.0 * tiny, 0.1 * tiny * BYTES_PER_GIB, 1.0 * tiny), + deployment(0.2 * tiny, 2.0 * tiny * BYTES_PER_GIB, 1.0 * tiny), + ] + .into_iter() + .map(|mut d| { + d.config.merge_cpu_secs = tiny; + d + }) + .collect::>(); + let label_sets = BTreeMap::from([( + LabelSet::new(), + LabelSetInfo { + cardinality: 1, + arrival_rate_per_sec: 1.0, + }, + )]); + for family in [family(8.0), family(32.0)] { + let best = brute_force(&rqes, &deployments) + .iter() + .map(|mapping| { + family.usd_per_hour(&score(&rqes, &deployments, mapping, &label_sets)) + }) + .min_by(f64::total_cmp) + .expect("test workload is servable"); + let milp = minimize_cost( + &rqes, + &deployments, + &label_sets, + &MilpBounds::default(), + &family, + ) + .expect("feasible MILP"); + let got = family.usd_per_hour(&milp.objectives); + assert!((got - best).abs() <= 1e-9 * best, "{got} vs {best}"); + } + + // The cheap deployment's latency is 1.5e-7 s, over a 1e-7 s bound by + // less than HiGHS's absolute feasibility tolerance. + let rqes = vec![rqe()]; + let deployments = vec![ + deployment(tiny, tiny, 1.5e-7), + deployment(10.0 * tiny, tiny, 0.5e-7), + ]; + let milp = minimize_cost( + &rqes, + &deployments, + &label_sets, + &MilpBounds { + max_peak_query_memory_bytes: None, + max_query_latency_secs: vec![Some(1e-7)], + }, + &family(8.0), + ) + .expect("feasible MILP"); + assert_eq!(milp.mapping, vec![1]); + assert!(milp.objectives.query_latency_secs[0] <= 1e-7); + } } diff --git a/rqe-optimizer/src/objectives.rs b/rqe-optimizer/src/objectives.rs index 567569f6..3bd4a134 100644 --- a/rqe-optimizer/src/objectives.rs +++ b/rqe-optimizer/src/objectives.rs @@ -1,7 +1,9 @@ //! Analytical scoring built from measured per-operation costs (§5). use crate::{Deployment, LabelSetTable, Mapping, Rqe}; -use std::collections::BTreeSet; +use std::collections::{BTreeMap, BTreeSet}; + +pub(crate) const BYTES_PER_GIB: f64 = (1u64 << 30) as f64; #[derive(Debug, Clone)] pub struct Objectives { @@ -10,6 +12,9 @@ pub struct Objectives { pub merge_cpu_secs_per_sec: f64, pub query_cpu_secs_per_sec: f64, pub tco_cpu_secs_per_sec: f64, + /// State held by all active deployments: each keeps the instances needed + /// by the longest lookback it serves. + pub retained_memory_bytes: f64, /// Index-aligned with the input RQE slice; display IDs need not be unique. pub query_latency_secs: Vec, } @@ -37,6 +42,24 @@ pub fn score( }) .sum(); + let mut retained_instances: BTreeMap = BTreeMap::new(); + for (r, &di) in rqes.iter().zip(mapping) { + let count = deployments[di] + .retained_instance_count(r.lookback_secs) + .expect("mapping only contains eligible pairs"); + let held = retained_instances.entry(di).or_default(); + *held = (*held).max(count); + } + let retained_memory_bytes = retained_instances + .iter() + .map(|(&di, &count)| { + let d = &deployments[di]; + label_sets[&d.labels].cardinality as f64 + * d.config.mem_bytes_per_instance + * count as f64 + }) + .sum(); + let mut peak_query_memory_bytes = 0.0_f64; let mut merge_cpu_secs_per_sec = 0.0; let mut query_cpu_secs_per_sec = 0.0; @@ -66,10 +89,61 @@ pub fn score( merge_cpu_secs_per_sec, query_cpu_secs_per_sec, tco_cpu_secs_per_sec, + retained_memory_bytes, query_latency_secs, } } +/// One instance type standing for an EC2 machine family (§4 of the +/// evaluation plan, ProjectASAP/ASAPQuery#777). +#[derive(Debug, Clone, PartialEq)] +pub struct MachineFamily { + pub family: String, + pub vcpu: f64, + pub memory_gib: f64, + pub usd_per_hour: f64, +} + +impl MachineFamily { + /// Reads the `families` of a `scripts/fetch_ec2_pricing.py` snapshot. + pub fn from_pricing_json(json: &str) -> Result, String> { + let snapshot: serde_json::Value = serde_json::from_str(json).map_err(|e| e.to_string())?; + let families = snapshot["families"] + .as_array() + .ok_or("pricing snapshot has no families array")?; + families + .iter() + .map(|f| { + let number = |key: &str| { + f[key] + .as_f64() + .ok_or_else(|| format!("family is missing numeric {key}")) + }; + Ok(Self { + family: f["family"] + .as_str() + .ok_or("family is missing its name")? + .to_string(), + vcpu: number("vcpu")?, + memory_gib: number("memory_gib")?, + usd_per_hour: number("usd_per_hour")?, + }) + }) + .collect() + } + + /// Fractional instances needed for the plan: whichever of CPU or memory + /// runs out first. + pub fn instances(&self, objectives: &Objectives) -> f64 { + (objectives.tco_cpu_secs_per_sec / self.vcpu) + .max(objectives.retained_memory_bytes / BYTES_PER_GIB / self.memory_gib) + } + + pub fn usd_per_hour(&self, objectives: &Objectives) -> f64 { + self.usd_per_hour * self.instances(objectives) + } +} + #[cfg(test)] mod tests { use super::*; @@ -118,5 +192,81 @@ mod tests { assert_eq!(result.query_latency_secs, vec![44.0]); // 4 × (5 query + 2 merges × 3) assert_eq!(result.query_cpu_secs_per_sec, 20.0 / 30.0); assert_eq!(result.merge_cpu_secs_per_sec, 24.0 / 30.0); + assert_eq!(result.retained_memory_bytes, 320.0); // 4 × 10 × (20 + 60) / 10 + } + + #[test] + fn shared_deployment_retains_for_its_longest_lookback() { + let labels = LabelSet::new(); + let deployment = Deployment { + capability: Capability::Freq, + labels: labels.clone(), + config: AtomicCostEntry { + sketch: "cms-fastpath-vector2d".into(), + sketch_config: serde_json::json!(null), + mem_bytes_per_instance: 10.0, + insert_cpu_secs: 1.0, + merge_cpu_secs: 1.0, + query_cpu_secs: 1.0, + query_accuracy: BTreeMap::from([("err".into(), 0.0)]), + }, + window_secs: 10, + slide_secs: 10, + }; + let rqe = |lookback_secs| Rqe { + id: "r".into(), + capability: Capability::Freq, + lookback_secs, + interval_secs: 10, + labels: labels.clone(), + accuracy_metric: "err".into(), + accuracy_tolerance: 1.0, + accuracy_direction: AccuracyDirection::LowerIsBetter, + }; + let tables = BTreeMap::from([( + labels.clone(), + LabelSetInfo { + cardinality: 2, + arrival_rate_per_sec: 1.0, + }, + )]); + let result = score(&[rqe(60), rqe(600)], &[deployment], &vec![0, 0], &tables); + // One copy of state, sized for the 600 s lookback: 2 × 10 × (10 + 600) / 10. + assert_eq!(result.retained_memory_bytes, 1220.0); + } + + #[test] + fn machine_family_is_sized_by_its_binding_resource() { + let family = MachineFamily { + family: "f".into(), + vcpu: 4.0, + memory_gib: 8.0, + usd_per_hour: 2.0, + }; + let mut objectives = Objectives { + peak_query_memory_bytes: 0.0, + ingest_cpu_secs_per_sec: 0.0, + merge_cpu_secs_per_sec: 0.0, + query_cpu_secs_per_sec: 0.0, + tco_cpu_secs_per_sec: 2.0, + retained_memory_bytes: BYTES_PER_GIB, + query_latency_secs: Vec::new(), + }; + assert_eq!(family.instances(&objectives), 0.5); // CPU binds + objectives.retained_memory_bytes = 16.0 * BYTES_PER_GIB; + assert_eq!(family.usd_per_hour(&objectives), 4.0); // memory binds: 2 instances + } + + #[test] + fn reads_committed_pricing_snapshot() { + let families = + MachineFamily::from_pricing_json(include_str!("../data/ec2-pricing-2026-10-04.json")) + .expect("committed snapshot parses"); + let names: Vec<_> = families.iter().map(|f| f.family.as_str()).collect(); + assert_eq!( + names, + ["compute_optimized", "general_purpose", "memory_optimized"] + ); + assert!(families.iter().all(|f| f.usd_per_hour > 0.0)); } } diff --git a/rqe-optimizer/src/pareto.rs b/rqe-optimizer/src/pareto.rs index e4d71257..a48b6181 100644 --- a/rqe-optimizer/src/pareto.rs +++ b/rqe-optimizer/src/pareto.rs @@ -82,6 +82,7 @@ mod tests { fn objectives(memory: f64, cpu: f64, latency: f64) -> Objectives { Objectives { peak_query_memory_bytes: memory, + retained_memory_bytes: 0.0, ingest_cpu_secs_per_sec: 0.0, merge_cpu_secs_per_sec: 0.0, query_cpu_secs_per_sec: 0.0, diff --git a/scripts/fetch_ec2_pricing.py b/scripts/fetch_ec2_pricing.py new file mode 100755 index 00000000..f33cba19 --- /dev/null +++ b/scripts/fetch_ec2_pricing.py @@ -0,0 +1,70 @@ +#!/usr/bin/env python3 +"""Snapshot on-demand EC2 prices for the rqe-optimizer machine families. + +Reads the public price file behind aws.amazon.com/ec2/pricing/on-demand +(us-east-1, Linux, shared tenancy) and writes the three instance types the +optimizer prices plans with. Never edit the output by hand; rerun this. + +Usage: + scripts/fetch_ec2_pricing.py rqe-optimizer/data/ec2-pricing-$(date +%F).json +""" +import datetime +import gzip +import json +import sys +import urllib.request + +SOURCE = ( + "https://b0.p.awsstatic.com/pricing/2.0/meteredUnitMaps/ec2/USD/current/" + "ec2-ondemand-without-sec-sel/US%20East%20(N.%20Virginia)/Linux/index.json" +) +REGION = "US East (N. Virginia)" +# One size per family: plans are priced in fractional instances, so the size +# within a family does not change the result. +INSTANCE_TYPES = { + "compute_optimized": "c7i.xlarge", + "general_purpose": "m7i.xlarge", + "memory_optimized": "r7i.xlarge", +} + + +def main(out_path): + with urllib.request.urlopen(SOURCE) as response: + body = response.read() + if body[:2] == b"\x1f\x8b": + body = gzip.decompress(body) + data = json.loads(body) + rows = {row["Instance Type"]: row for row in data["regions"][REGION].values()} + + families = [] + for family, instance_type in INSTANCE_TYPES.items(): + row = rows[instance_type] + memory, unit = row["Memory"].split() + assert unit == "GiB", row["Memory"] + families.append( + { + "family": family, + "instance_type": instance_type, + "vcpu": int(row["vCPU"]), + "memory_gib": float(memory), + "usd_per_hour": float(row["price"]), + "rate_code": row["rateCode"], + } + ) + + snapshot = { + "source": SOURCE, + "region": REGION, + "operating_system": "Linux", + "pricing": "on-demand", + "publication_date": data["manifest"]["hawkFilePublicationDate"], + "fetched": datetime.date.today().isoformat(), + "families": families, + } + with open(out_path, "w") as out: + json.dump(snapshot, out, indent=2) + out.write("\n") + + +if __name__ == "__main__": + main(sys.argv[1])