Skip to content

Commit c66f024

Browse files
Added support for asap-planner-rs to infer query repetitions and config from Prometheus query log
1 parent 87383ed commit c66f024

18 files changed

Lines changed: 679 additions & 19 deletions

File tree

‎Cargo.lock‎

Lines changed: 1 addition & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

‎asap-planner-rs/Cargo.toml‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ tracing.workspace = true
2323
tracing-subscriber.workspace = true
2424
clap.workspace = true
2525
indexmap.workspace = true
26+
chrono.workspace = true
2627
promql-parser = "0.5.0"
2728

2829
[dev-dependencies]

‎asap-planner-rs/src/config/input.rs‎

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -13,10 +13,17 @@ pub struct QueryGroup {
1313
pub id: Option<u32>,
1414
pub queries: Vec<String>,
1515
pub repetition_delay: u64,
16+
#[serde(default)]
1617
pub controller_options: ControllerOptions,
18+
/// Per-group step override (seconds). Falls back to `RuntimeOptions::step` when None.
19+
#[serde(default)]
20+
pub step: Option<u64>,
21+
/// Per-group range_duration override (seconds). Falls back to `RuntimeOptions::range_duration` when None.
22+
#[serde(default)]
23+
pub range_duration: Option<u64>,
1724
}
1825

19-
#[derive(Debug, Clone, Deserialize)]
26+
#[derive(Debug, Clone, Deserialize, Default)]
2027
pub struct ControllerOptions {
2128
pub accuracy_sla: f64,
2229
pub latency_sla: f64,

‎asap-planner-rs/src/lib.rs‎

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ pub mod config;
22
pub mod error;
33
pub mod output;
44
pub mod planner;
5+
pub mod query_log;
56

67
use serde_yaml::Value as YamlValue;
78
use std::path::Path;
@@ -178,6 +179,30 @@ impl Controller {
178179
})
179180
}
180181

182+
/// Build a `Controller` from a Prometheus query log file and a metrics config YAML.
183+
///
184+
/// - `log_path`: newline-delimited JSON query log (Prometheus `--query.log-file` output)
185+
/// - `metrics_path`: YAML file with a `metrics:` section listing metric names and labels
186+
pub fn from_query_log(
187+
log_path: &Path,
188+
metrics_path: &Path,
189+
opts: RuntimeOptions,
190+
) -> Result<Self, ControllerError> {
191+
let entries = query_log::parse_log_file(log_path)?;
192+
let (instants, ranges) =
193+
query_log::infer_queries(&entries, opts.prometheus_scrape_interval);
194+
195+
let metrics_yaml = std::fs::read_to_string(metrics_path)?;
196+
let metrics_config: query_log::MetricsConfig = serde_yaml::from_str(&metrics_yaml)?;
197+
198+
let config = query_log::to_controller_config(instants, ranges, metrics_config.metrics);
199+
200+
Ok(Self {
201+
config,
202+
options: opts,
203+
})
204+
}
205+
181206
pub fn generate(&self) -> Result<PlannerOutput, ControllerError> {
182207
let output = output::generator::generate_plan(&self.config, &self.options)?;
183208
Ok(PlannerOutput {

‎asap-planner-rs/src/main.rs‎

Lines changed: 22 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -5,8 +5,17 @@ use std::path::PathBuf;
55
#[derive(Parser, Debug)]
66
#[command(name = "asap-planner", about = "ASAP Query Planner")]
77
struct Args {
8-
#[arg(long = "input_config")]
9-
input_config: PathBuf,
8+
/// Path to a hand-authored YAML workload config. Mutually exclusive with --query-log.
9+
#[arg(long = "input_config", conflicts_with = "query_log")]
10+
input_config: Option<PathBuf>,
11+
12+
/// Path to a Prometheus query log file (newline-delimited JSON). Mutually exclusive with --input_config.
13+
#[arg(long = "query-log", conflicts_with = "input_config")]
14+
query_log: Option<PathBuf>,
15+
16+
/// Path to a metrics config YAML (required when using --query-log).
17+
#[arg(long = "metrics-config", requires = "query_log")]
18+
metrics_config: Option<PathBuf>,
1019

1120
#[arg(long = "output_dir")]
1221
output_dir: PathBuf,
@@ -60,9 +69,18 @@ fn main() -> anyhow::Result<()> {
6069
step: args.step,
6170
};
6271

63-
let controller = Controller::from_file(&args.input_config, opts)?;
64-
controller.generate_to_dir(&args.output_dir)?;
72+
let controller = match (args.input_config, args.query_log) {
73+
(Some(config_path), None) => Controller::from_file(&config_path, opts)?,
74+
(None, Some(log_path)) => {
75+
let metrics_path = args
76+
.metrics_config
77+
.expect("--metrics-config is required when using --query-log");
78+
Controller::from_query_log(&log_path, &metrics_path, opts)?
79+
}
80+
_ => anyhow::bail!("exactly one of --input_config or --query-log must be provided"),
81+
};
6582

83+
controller.generate_to_dir(&args.output_dir)?;
6684
println!("Generated configs in {}", args.output_dir.display());
6785
Ok(())
6886
}

‎asap-planner-rs/src/output/generator.rs‎

Lines changed: 20 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -57,8 +57,8 @@ pub fn generate_plan(
5757
metric_schema.clone(),
5858
opts.streaming_engine,
5959
controller_config.sketch_parameters.clone(),
60-
opts.range_duration,
61-
opts.step,
60+
qg.range_duration.unwrap_or(opts.range_duration),
61+
qg.step.unwrap_or(opts.step),
6262
cleanup_policy,
6363
);
6464

@@ -73,14 +73,25 @@ pub fn generate_plan(
7373
}
7474

7575
if should_process {
76-
let (configs, cleanup_param) = processor.get_streaming_aggregation_configs()?;
77-
let mut keys_for_query = Vec::new();
78-
for config in configs {
79-
let key = config.identifying_key();
80-
keys_for_query.push((key.clone(), cleanup_param));
81-
dedup_map.entry(key).or_insert(config);
76+
match processor.get_streaming_aggregation_configs() {
77+
Ok((configs, cleanup_param)) => {
78+
let mut keys_for_query = Vec::new();
79+
for config in configs {
80+
let key = config.identifying_key();
81+
keys_for_query.push((key.clone(), cleanup_param));
82+
dedup_map.entry(key).or_insert(config);
83+
}
84+
query_keys_map.insert(query_string.clone(), keys_for_query);
85+
}
86+
Err(ControllerError::UnknownMetric(ref metric)) => {
87+
tracing::warn!(
88+
query = %query_string,
89+
metric = %metric,
90+
"skipping query referencing unknown metric"
91+
);
92+
}
93+
Err(e) => return Err(e),
8294
}
83-
query_keys_map.insert(query_string.clone(), keys_for_query);
8495
}
8596
}
8697
}
Lines changed: 55 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,55 @@
1+
use serde::Deserialize;
2+
3+
use crate::config::input::{
4+
AggregateCleanupConfig, ControllerConfig, MetricDefinition, QueryGroup,
5+
};
6+
7+
use super::frequency::{InstantQueryInfo, RangeQueryInfo};
8+
9+
/// Subset of ControllerConfig used when loading from a metrics-only YAML file.
10+
#[derive(Deserialize)]
11+
pub struct MetricsConfig {
12+
pub metrics: Vec<MetricDefinition>,
13+
}
14+
15+
/// Build a `ControllerConfig` from extracted instant and range queries plus a metrics definition.
16+
///
17+
/// Each query becomes its own `QueryGroup` (one query per group, no SLA fields needed).
18+
pub fn to_controller_config(
19+
instants: Vec<InstantQueryInfo>,
20+
ranges: Vec<RangeQueryInfo>,
21+
metrics: Vec<MetricDefinition>,
22+
) -> ControllerConfig {
23+
let mut query_groups: Vec<QueryGroup> = Vec::new();
24+
25+
for info in instants {
26+
query_groups.push(QueryGroup {
27+
id: None,
28+
queries: vec![info.query],
29+
repetition_delay: info.repetition_delay,
30+
controller_options: Default::default(),
31+
step: None,
32+
range_duration: None,
33+
});
34+
}
35+
36+
for info in ranges {
37+
query_groups.push(QueryGroup {
38+
id: None,
39+
queries: vec![info.query],
40+
repetition_delay: info.repetition_delay,
41+
controller_options: Default::default(),
42+
step: Some(info.step),
43+
range_duration: Some(info.range_duration),
44+
});
45+
}
46+
47+
ControllerConfig {
48+
query_groups,
49+
metrics,
50+
sketch_parameters: None,
51+
aggregate_cleanup: Some(AggregateCleanupConfig {
52+
policy: Some("read_based".to_string()),
53+
}),
54+
}
55+
}

0 commit comments

Comments
 (0)