Skip to content

Commit 754d23b

Browse files
Added support for asap-planner-rs to infer query repetitions and config from Prometheus query log (#254)
1 parent e6d8afd commit 754d23b

18 files changed

Lines changed: 816 additions & 120 deletions

File tree

‎Cargo.lock‎

Lines changed: 133 additions & 100 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
@@ -25,6 +25,7 @@ tracing.workspace = true
2525
tracing-subscriber.workspace = true
2626
clap.workspace = true
2727
indexmap.workspace = true
28+
chrono.workspace = true
2829
promql-parser = "0.5.0"
2930

3031
[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;
@@ -291,6 +292,30 @@ impl Controller {
291292
})
292293
}
293294

295+
/// Build a `Controller` from a Prometheus query log file and a metrics config YAML.
296+
///
297+
/// - `log_path`: newline-delimited JSON query log (Prometheus `--query.log-file` output)
298+
/// - `metrics_path`: YAML file with a `metrics:` section listing metric names and labels
299+
pub fn from_query_log(
300+
log_path: &Path,
301+
metrics_path: &Path,
302+
opts: RuntimeOptions,
303+
) -> Result<Self, ControllerError> {
304+
let entries = query_log::parse_log_file(log_path)?;
305+
let (instants, ranges) =
306+
query_log::infer_queries(&entries, opts.prometheus_scrape_interval);
307+
308+
let metrics_yaml = std::fs::read_to_string(metrics_path)?;
309+
let metrics_config: query_log::MetricsConfig = serde_yaml::from_str(&metrics_yaml)?;
310+
311+
let config = query_log::to_controller_config(instants, ranges, metrics_config.metrics);
312+
313+
Ok(Self {
314+
config,
315+
options: opts,
316+
})
317+
}
318+
294319
pub fn generate(&self) -> Result<PlannerOutput, ControllerError> {
295320
let output = output::generator::generate_plan(&self.config, &self.options)?;
296321
Ok(PlannerOutput {

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

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

1221
#[arg(long = "output_dir")]
1322
output_dir: PathBuf,
@@ -71,20 +80,33 @@ fn main() -> anyhow::Result<()> {
7180
range_duration: args.range_duration,
7281
step: args.step,
7382
};
74-
let controller = Controller::from_file(&args.input_config, opts)?;
83+
let controller = match (args.input_config, args.query_log) {
84+
(Some(config_path), None) => Controller::from_file(&config_path, opts)?,
85+
(None, Some(log_path)) => {
86+
let metrics_path = args
87+
.metrics_config
88+
.expect("--metrics-config is required when using --query-log");
89+
Controller::from_query_log(&log_path, &metrics_path, opts)?
90+
}
91+
_ => anyhow::bail!(
92+
"exactly one of --input_config or --query-log must be provided for PromQL mode"
93+
),
94+
};
7595
controller.generate_to_dir(&args.output_dir)?;
7696
}
7797
QueryLanguage::sql | QueryLanguage::elastic_sql => {
7898
let interval = args.data_ingestion_interval.ok_or_else(|| {
7999
anyhow::anyhow!("--data-ingestion-interval is required for SQL mode")
80100
})?;
101+
let config_path = args
102+
.input_config
103+
.ok_or_else(|| anyhow::anyhow!("--input_config is required for SQL mode"))?;
81104
let opts = SQLRuntimeOptions {
82105
streaming_engine: engine,
83106
query_evaluation_time: None,
84107
data_ingestion_interval: interval,
85108
};
86-
SQLController::from_file(&args.input_config, opts)?
87-
.generate_to_dir(&args.output_dir)?;
109+
SQLController::from_file(&config_path, opts)?.generate_to_dir(&args.output_dir)?;
88110
}
89111
QueryLanguage::elastic_querydsl => {
90112
anyhow::bail!("ElasticQueryDSL is not yet supported");

‎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)