Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
85 changes: 18 additions & 67 deletions asap-planner-rs/src/bin/optimizer_cli.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,15 +3,12 @@

use std::path::PathBuf;

use anyhow::Context;
use asap_planner::optimizer::{
build_milp_workload, load_flat_atomic_cost_table, load_workload_facts, plan_to_planner_output,
reject_avg_queries, solve_milp, MilpError,
parse_weight, plan_milp, plan_to_planner_output, reject_unwritable_queries, MilpInputs,
MilpPlan,
};
use asap_planner::ControllerConfig;
use clap::Parser;
use rqe_optimizer::milp::Objective;
use rqe_optimizer::saturation::SaturationCurves;

#[derive(Parser, Debug)]
#[command(
Expand Down Expand Up @@ -84,51 +81,26 @@ fn main() -> anyhow::Result<()> {
run_milp(&args, &config)
}

/// Objective weights must be finite and non-negative: a negative weight
/// rewards cost, and NaN poisons every coefficient.
fn parse_weight(s: &str) -> Result<f64, String> {
match s.parse::<f64>() {
Ok(w) if w.is_finite() && w >= 0.0 => Ok(w),
Ok(w) => Err(format!("must be finite and >= 0, got {w}")),
Err(e) => Err(e.to_string()),
}
}

fn run_milp(args: &Args, config: &ControllerConfig) -> anyhow::Result<()> {
config.warn_default_slas();
let Some(hints) = config.metrics.as_deref() else {
return Err(MilpError::MissingMetricHints.into());
};
let facts = load_workload_facts(&args.workload_facts, hints, args.data_ingestion_interval_ms)?;
let curves = SaturationCurves::load(&args.saturation_dir).with_context(|| {
format!(
"loading saturation curves from --saturation-dir {}",
args.saturation_dir.display()
)
})?;
let costs = load_flat_atomic_cost_table(&args.atomic_costs)?;
let Objective::AUCCost { w_cpu, w_mem } = Objective::default();
let (w_cpu, w_mem) = (args.w_cpu.unwrap_or(w_cpu), args.w_mem.unwrap_or(w_mem));
// All-zero weights make every plan cost 0, so the solver's pick is arbitrary.
anyhow::ensure!(
w_cpu > 0.0 || w_mem > 0.0,
"--w-cpu and --w-mem are both 0; at least one must be positive"
);
let objective = Objective::AUCCost { w_cpu, w_mem };
tracing::debug!(?objective, cost_rows = costs.len(), "milp: inputs loaded");

// Fail before solving when the plan would be written but can't be.
if args.output_dir.is_some() {
reject_avg_queries(config)?;
reject_unwritable_queries(config)?;
}
let workload = build_milp_workload(config, &facts, args.data_ingestion_interval_ms)?;
let solution = solve_milp(
&workload,
&facts,
&costs,
let MilpPlan {
workload,
solution,
objective,
args.allow_undeployable_families,
&|raqe, deployment| curves.accuracy(raqe, deployment, &facts),
} = plan_milp(
config,
&MilpInputs {
workload_facts: &args.workload_facts,
atomic_costs: &args.atomic_costs,
saturation_dir: &args.saturation_dir,
scrape_interval_ms: args.data_ingestion_interval_ms,
w_cpu: args.w_cpu,
w_mem: args.w_mem,
allow_undeployable_families: args.allow_undeployable_families,
},
)?;

println!("=== Deployments: {} ===", solution.deployments.len());
Expand Down Expand Up @@ -180,29 +152,8 @@ fn run_milp(args: &Args, config: &ControllerConfig) -> anyhow::Result<()> {
}

if let Some(dir) = &args.output_dir {
let output = plan_to_planner_output(config, &workload, &solution)?;
// Serialize both before writing either, so a failure can't leave a
// new streaming config next to a stale inference config.
let streaming = output.to_streaming_yaml_string()?;
let inference = output.to_inference_yaml_string()?;
std::fs::create_dir_all(dir)?;
std::fs::write(dir.join("streaming_config.yaml"), streaming)?;
std::fs::write(dir.join("inference_config.yaml"), inference)?;
plan_to_planner_output(config, &workload, &solution)?.write_to_dir(dir)?;
println!("\nwrote configs to {}", dir.display());
}
Ok(())
}

#[cfg(test)]
mod tests {
use super::parse_weight;

#[test]
fn weights_must_be_finite_and_non_negative() {
assert_eq!(parse_weight("0.5"), Ok(0.5));
assert_eq!(parse_weight("0"), Ok(0.0));
for bad in ["-1", "NaN", "inf", "x"] {
assert!(parse_weight(bad).is_err(), "{bad}");
}
}
}
6 changes: 1 addition & 5 deletions asap-planner-rs/src/elastic_dsl/controller.rs
Original file line number Diff line number Diff line change
Expand Up @@ -35,11 +35,7 @@ impl ElasticController {

pub fn generate_to_dir(&self, dir: &Path) -> Result<PlannerOutput, ControllerError> {
let output = self.generate()?;
std::fs::create_dir_all(dir)?;
let streaming_str = serde_yaml::to_string(output.streaming_yaml())?;
let inference_str = serde_yaml::to_string(output.inference_yaml())?;
std::fs::write(dir.join("streaming_config.yaml"), streaming_str)?;
std::fs::write(dir.join("inference_config.yaml"), inference_str)?;
output.write_to_dir(dir)?;
Ok(output)
}
}
120 changes: 118 additions & 2 deletions asap-planner-rs/src/main.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,9 @@
use asap_planner::optimizer::{
parse_weight, plan_milp, plan_to_planner_output, reject_unwritable_queries, MilpInputs,
};
use asap_planner::{
Controller, ElasticController, ElasticRuntimeOptions, RuntimeOptions, SQLController,
SQLRuntimeOptions, StreamingEngine,
Controller, ControllerConfig, ElasticController, ElasticRuntimeOptions, RuntimeOptions,
SQLController, SQLRuntimeOptions, StreamingEngine,
};
use asap_types::enums::QueryLanguage;
use clap::Parser;
Expand Down Expand Up @@ -53,6 +56,36 @@ struct Args {
#[arg(long = "clickhouse-database", required = false)]
clickhouse_database: Option<String>,

/// `milp` plans with sketch-bench's rqe-optimizer: PromQL with
/// --input_config only, labels from its `metrics:` hints.
#[arg(long, value_enum, default_value = "legacy")]
planner: PlannerArg,

/// MILP only. YAML workload facts: per metric, positive `value_range` and
/// `cardinality` per label set, including all labels (the series count),
/// plus the `shape` of each grouping sketches may serve.
#[arg(long = "workload-facts", required_if_eq("planner", "milp"))]
workload_facts: Option<PathBuf>,

/// MILP only. The flat cost table sketch-bench's
/// `study_saturation.py --phase optimizer-cost` writes
/// (`rqe_atomic_costs.json`).
#[arg(long = "atomic-costs", required_if_eq("planner", "milp"))]
atomic_costs: Option<PathBuf>,

/// MILP only. sketch-bench's saturation-study directory: sketch accuracy
/// is read off its error-vs-N curves at each grouping's `shape`.
#[arg(long = "saturation-dir", required_if_eq("planner", "milp"))]
saturation_dir: Option<PathBuf>,

/// MILP only. Objective weight on CPU-sec/sec. Default: rqe-optimizer's.
#[arg(long = "w-cpu", value_parser = parse_weight)]
w_cpu: Option<f64>,

/// MILP only. Objective weight on memory GiB. Default: rqe-optimizer's.
#[arg(long = "w-mem", value_parser = parse_weight)]
w_mem: Option<f64>,

#[arg(short, long, action = clap::ArgAction::Count)]
verbose: u8,
}
Expand All @@ -62,6 +95,12 @@ enum EngineArg {
Precompute,
}

#[derive(clap::ValueEnum, Debug, Clone, Copy, PartialEq, Eq)]
enum PlannerArg {
Legacy,
Milp,
}

fn main() -> anyhow::Result<()> {
let args = Args::parse();

Expand All @@ -77,6 +116,21 @@ fn main() -> anyhow::Result<()> {
EngineArg::Precompute => StreamingEngine::Precompute,
};

if args.planner == PlannerArg::Milp {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

With --planner milp, --clickhouse-url, --clickhouse-database and --streaming_engine are silently ignored, while --prometheus-url, --enable-punting, --step-ms etc. are explicitly rejected in run_milp. Should these be rejected too for consistency?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 0c4bcf5: --clickhouse-url and --clickhouse-database are now rejected under --planner milp. I left --streaming_engine as is: its only value is precompute, which is what MILP emits, so it isn't ignored.

run_milp(&args)?;
println!("Generated configs in {}", args.output_dir.display());
return Ok(());
}
anyhow::ensure!(
args.workload_facts.is_none()
&& args.atomic_costs.is_none()
&& args.saturation_dir.is_none()
&& args.w_cpu.is_none()
&& args.w_mem.is_none(),
"--workload-facts, --atomic-costs, --saturation-dir, --w-cpu and --w-mem require \
--planner milp"
);

match args.query_language {
QueryLanguage::promql => {
let scrape_interval_ms = args.data_ingestion_interval_ms.ok_or_else(|| {
Expand Down Expand Up @@ -158,3 +212,65 @@ fn main() -> anyhow::Result<()> {
println!("Generated configs in {}", args.output_dir.display());
Ok(())
}

fn run_milp(args: &Args) -> anyhow::Result<()> {
anyhow::ensure!(
matches!(args.query_language, QueryLanguage::promql),
"--planner milp supports only --query-language promql"
);
anyhow::ensure!(
args.query_log.is_none(),
"--planner milp needs --input_config: query logs have no `metrics:` hints"
);
anyhow::ensure!(
args.prometheus_url.is_none(),
"--planner milp takes labels from the config's `metrics:` hints, not --prometheus-url"
);
anyhow::ensure!(
!args.enable_punting && args.range_duration_ms == 0 && args.step_ms == 0,
"--enable-punting, --range-duration-ms and --step-ms don't apply to --planner milp"
);
anyhow::ensure!(
args.clickhouse_url.is_none() && args.clickhouse_database.is_none(),
"--clickhouse-url and --clickhouse-database don't apply to --planner milp"
);
let config_path = args
.input_config
.as_deref()
.ok_or_else(|| anyhow::anyhow!("--planner milp requires --input_config"))?;
let scrape_interval_ms = args
.data_ingestion_interval_ms
.ok_or_else(|| anyhow::anyhow!("--planner milp requires --data-ingestion-interval-ms"))?;
// Per-series sample rates divide by it.
anyhow::ensure!(
scrape_interval_ms > 0,
"--data-ingestion-interval-ms must be positive"
);
let config: ControllerConfig = serde_yaml::from_str(&std::fs::read_to_string(config_path)?)?;

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This parses the config with serde_yaml directly, skipping the windowing.validate() that every legacy Controller::from_* constructor runs. E.g. windowing: {type: sliding, window_size_ms: 1000} without slide_interval_ms, or window_size_ms: 0, is rejected by the legacy planner but planned and written under --planner milp.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Addressed in 0c4bcf5 by rejecting windowing under MILP (see the build_milp_workload thread). The MILP picks its own windows, so a validated windowing would still be ignored.

// Fail before solving: the plan is always written.
reject_unwritable_queries(&config)?;
let plan = plan_milp(
&config,
&MilpInputs {
workload_facts: args
.workload_facts
.as_deref()
.expect("clap requires --workload-facts with --planner milp"),
atomic_costs: args
.atomic_costs
.as_deref()
.expect("clap requires --atomic-costs with --planner milp"),
saturation_dir: args
.saturation_dir
.as_deref()
.expect("clap requires --saturation-dir with --planner milp"),
scrape_interval_ms,
w_cpu: args.w_cpu,
w_mem: args.w_mem,
allow_undeployable_families: false,
},
)?;
plan_to_planner_output(&config, &plan.workload, &plan.solution)?
.write_to_dir(&args.output_dir)?;
Ok(())
}
Loading
Loading