Skip to content

Commit a095704

Browse files
fix(planner, tools): use milliseconds for repetition_delay (#449)
* fix(planner, tools): use milliseconds for repetition_delay * fix(planner): change more variables to ms precision * more fixes * fixes to python infra scripts * fixed more bugs * removed stale codel * small fixes from CR * removed unnecessary deny_unknown_fields --------- Co-authored-by: Claude Sonnet 4.6 <noreply@anthropic.com>
1 parent d92c4e8 commit a095704

22 files changed

Lines changed: 279 additions & 195 deletions

File tree

asap-planner-rs/docker-compose.yml.j2

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -12,8 +12,8 @@ services:
1212
"--streaming_engine", "{{ streaming_engine }}",
1313
"--prometheus-url", "{{ prometheus_url }}",
1414
"--query-language", "{{ query_language }}"{% if punting %},
15-
"--enable-punting"{% endif %}{% if data_ingestion_interval is not none %},
16-
"--data-ingestion-interval", "{{ data_ingestion_interval }}"{% endif %}
15+
"--enable-punting"{% endif %}{% if data_ingestion_interval_ms is not none %},
16+
"--data-ingestion-interval-ms", "{{ data_ingestion_interval_ms }}"{% endif %}
1717
]
1818
network_mode: "host"
1919
restart: no

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

Lines changed: 19 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ use asap_types::streaming_config::StreamingConfig;
44
use asap_types::PromQLSchema;
55
use promql_utilities::data_model::KeyByLabelNames;
66
use serde::Deserialize;
7+
use tracing::warn;
78

89
#[derive(Debug, Clone, Deserialize)]
910
#[serde(deny_unknown_fields)]
@@ -30,6 +31,22 @@ pub struct ControllerConfig {
3031
}
3132

3233
impl ControllerConfig {
34+
/// Warn if any query group has both SLAs at 0.0 (the serde Default),
35+
/// which indicates `controller_options` was omitted from the config.
36+
pub fn warn_default_slas(&self) {
37+
for qg in &self.query_groups {
38+
let opts = &qg.controller_options;
39+
if opts.accuracy_sla == 0.0 && opts.latency_sla == 0.0 {
40+
warn!(
41+
query_group_id = ?qg.id,
42+
"controller_options not set in query group; \
43+
accuracy_sla=0.0 and latency_sla=0.0 will be used — \
44+
add controller_options to your config"
45+
);
46+
}
47+
}
48+
}
49+
3350
/// Build a `PromQLSchema` from the `metrics` hints in this config.
3451
/// Returns an empty schema if no hints are present.
3552
pub fn schema_from_hints(&self) -> PromQLSchema {
@@ -133,7 +150,7 @@ pub struct SQLControllerConfig {
133150
pub struct SQLQueryGroup {
134151
pub id: Option<u32>,
135152
pub queries: Vec<String>,
136-
pub repetition_delay: u64,
153+
pub repetition_delay_ms: u64,
137154
pub controller_options: ControllerOptions,
138155
}
139156

@@ -157,7 +174,7 @@ pub struct ElasticDSLControllerConfig {
157174
pub struct ElasticDSLQueryGroup {
158175
pub id: Option<u32>,
159176
pub queries: Vec<String>,
160-
pub repetition_delay: u64,
177+
pub repetition_delay_ms: u64,
161178
pub index: String,
162179
pub time_field: String,
163180
pub controller_options: ControllerOptions,

asap-planner-rs/src/elastic_dsl/generator.rs

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -47,7 +47,7 @@ impl ElasticIndexSchemaBuilder {
4747

4848
pub struct ElasticRuntimeOptions {
4949
pub streaming_engine: StreamingEngine,
50-
pub data_ingestion_interval: u64,
50+
pub data_ingestion_interval_ms: u64,
5151
}
5252

5353
pub fn generate_elastic_plan(
@@ -60,12 +60,12 @@ pub fn generate_elastic_plan(
6060
.and_then(|c| c.policy)
6161
.unwrap_or(CleanupPolicy::ReadBased);
6262

63-
// Validate T % data_ingestion_interval == 0
63+
// Validate T % data_ingestion_interval_ms == 0
6464
for qg in &config.query_groups {
65-
if qg.repetition_delay % opts.data_ingestion_interval != 0 {
65+
if qg.repetition_delay_ms % opts.data_ingestion_interval_ms != 0 {
6666
return Err(ControllerError::PlannerError(format!(
67-
"repetition_delay {} is not a multiple of data_ingestion_interval {}",
68-
qg.repetition_delay, opts.data_ingestion_interval
67+
"repetition_delay_ms {} is not a multiple of data_ingestion_interval_ms {}",
68+
qg.repetition_delay_ms, opts.data_ingestion_interval_ms
6969
)));
7070
}
7171
}
@@ -111,8 +111,8 @@ pub fn generate_elastic_plan(
111111
for query_string in &qg.queries {
112112
let processor = ElasticSingleQueryProcessor::new(
113113
query_string.clone(),
114-
qg.repetition_delay,
115-
opts.data_ingestion_interval,
114+
qg.repetition_delay_ms,
115+
opts.data_ingestion_interval_ms,
116116
index_schema_builders[&qg.index].clone(),
117117
opts.streaming_engine,
118118
config.sketch_parameters.clone(),

asap-planner-rs/src/main.rs

Lines changed: 10 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -45,8 +45,8 @@ struct Args {
4545
#[arg(long = "query-language", value_enum, default_value = "promql")]
4646
query_language: QueryLanguage,
4747

48-
#[arg(long = "data-ingestion-interval", required = false)]
49-
data_ingestion_interval: Option<u64>,
48+
#[arg(long = "data-ingestion-interval-ms", required = false)]
49+
data_ingestion_interval_ms: Option<u64>,
5050

5151
/// ClickHouse base URL for auto-inferring metadata_columns when not listed
5252
/// in the config file. Example: http://localhost:8123
@@ -123,16 +123,16 @@ fn main() -> anyhow::Result<()> {
123123
controller.generate_to_dir(&args.output_dir)?;
124124
}
125125
QueryLanguage::sql | QueryLanguage::elastic_sql => {
126-
let interval = args.data_ingestion_interval.ok_or_else(|| {
127-
anyhow::anyhow!("--data-ingestion-interval is required for SQL mode")
126+
let interval = args.data_ingestion_interval_ms.ok_or_else(|| {
127+
anyhow::anyhow!("--data-ingestion-interval-ms is required for SQL mode")
128128
})?;
129129
let config_path = args
130130
.input_config
131131
.ok_or_else(|| anyhow::anyhow!("--input_config is required for SQL mode"))?;
132132
let opts = SQLRuntimeOptions {
133133
streaming_engine: engine,
134134
query_evaluation_time: None,
135-
data_ingestion_interval: interval,
135+
data_ingestion_interval_ms: interval,
136136
};
137137
let controller = match args.clickhouse_url {
138138
Some(ref url) => SQLController::from_file_with_discovery(
@@ -146,15 +146,17 @@ fn main() -> anyhow::Result<()> {
146146
controller.generate_to_dir(&args.output_dir)?;
147147
}
148148
QueryLanguage::elastic_querydsl => {
149-
let interval = args.data_ingestion_interval.ok_or_else(|| {
150-
anyhow::anyhow!("--data-ingestion-interval is required for Elasticsearch DSL mode")
149+
let interval = args.data_ingestion_interval_ms.ok_or_else(|| {
150+
anyhow::anyhow!(
151+
"--data-ingestion-interval-ms is required for Elasticsearch DSL mode"
152+
)
151153
})?;
152154
let config_path = args.input_config.ok_or_else(|| {
153155
anyhow::anyhow!("--input_config is required for Elasticsearch DSL mode")
154156
})?;
155157
let opts = ElasticRuntimeOptions {
156158
streaming_engine: engine,
157-
data_ingestion_interval: interval,
159+
data_ingestion_interval_ms: interval,
158160
};
159161
ElasticController::from_file(&config_path, opts)?.generate_to_dir(&args.output_dir)?;
160162
}

asap-planner-rs/src/optimizer/pipeline.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -78,7 +78,7 @@ pub fn run_greedy_pipeline(
7878
}
7979

8080
/// Convert a `ControllerConfig`'s query groups into a flat list of RQEs.
81-
/// Each (query, repetition_delay) pair becomes one RQE.
81+
/// Each (query, repetition_delay_ms) pair becomes one RQE.
8282
fn config_to_rqes(config: &ControllerConfig) -> Vec<RQE> {
8383
config
8484
.query_groups

asap-planner-rs/src/planner/cleanup.rs

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -75,19 +75,19 @@ pub fn get_cleanup_param(
7575
/// SQL cleanup param — SQL queries are always instant (no range_duration/step).
7676
pub fn get_sql_cleanup_param(
7777
cleanup_policy: CleanupPolicy,
78-
t_lookback: u64,
79-
t_repeat: u64,
78+
t_lookback_ms: u64,
79+
t_repeat_ms: u64,
8080
) -> Result<u64, String> {
8181
match cleanup_policy {
8282
CleanupPolicy::CircularBuffer | CleanupPolicy::ReadBased => {
83-
if t_repeat == 0 {
83+
if t_repeat_ms == 0 {
8484
return Err(
85-
"repetition_delay must be > 0 for cleanup param calculation; \
86-
set a non-zero repetition_delay in your query group config"
85+
"repetition_delay_ms must be > 0 for cleanup param calculation; \
86+
set a non-zero repetition_delay_ms in your query group config"
8787
.to_string(),
8888
);
8989
}
90-
Ok(t_lookback.div_ceil(t_repeat))
90+
Ok(t_lookback_ms.div_ceil(t_repeat_ms))
9191
}
9292
CleanupPolicy::NoCleanup => {
9393
Err("NoCleanup policy should not call get_sql_cleanup_param".to_string())

asap-planner-rs/src/planner/elastic_dsl.rs

Lines changed: 9 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -19,9 +19,9 @@ use indexmap::IndexSet;
1919

2020
pub struct ElasticSingleQueryProcessor {
2121
query_string: String,
22-
t_repeat: u64,
22+
t_repeat_ms: u64,
2323
#[allow(dead_code)]
24-
data_ingestion_interval: u64,
24+
data_ingestion_interval_ms: u64,
2525
index_schema: ElasticIndexSchemaBuilder,
2626
#[allow(dead_code)]
2727
streaming_engine: StreamingEngine,
@@ -33,17 +33,17 @@ impl ElasticSingleQueryProcessor {
3333
#[allow(clippy::too_many_arguments)]
3434
pub fn new(
3535
query_string: String,
36-
t_repeat: u64,
37-
data_ingestion_interval: u64,
36+
t_repeat_ms: u64,
37+
data_ingestion_interval_ms: u64,
3838
index_schema: ElasticIndexSchemaBuilder,
3939
streaming_engine: StreamingEngine,
4040
sketch_parameters: Option<SketchParameterOverrides>,
4141
cleanup_policy: CleanupPolicy,
4242
) -> Self {
4343
Self {
4444
query_string,
45-
t_repeat,
46-
data_ingestion_interval,
45+
t_repeat_ms,
46+
data_ingestion_interval_ms,
4747
index_schema,
4848
streaming_engine,
4949
sketch_parameters,
@@ -65,11 +65,7 @@ impl ElasticSingleQueryProcessor {
6565
// Get aggregation type and statistics
6666
let (treatment_type, statistics) = get_elastic_statistics(&query_info.aggregation)?;
6767

68-
// Build window config (always tumbling for Elasticsearch queries).
69-
// self.t_repeat is seconds (Elastic DSL is out of scope for the ms-precision
70-
// rename, see issue #427) — IntermediateWindowConfig is now ms-typed throughout,
71-
// so convert at this boundary, same as planner::sql::compute_sql_window.
72-
let t_repeat_ms = self.t_repeat * 1000;
68+
let t_repeat_ms = self.t_repeat_ms;
7369
let window_cfg = IntermediateWindowConfig {
7470
window_size_ms: t_repeat_ms,
7571
slide_interval_ms: t_repeat_ms,
@@ -134,7 +130,7 @@ impl ElasticSingleQueryProcessor {
134130
_ => false,
135131
})
136132
.and_then(|p| range_query_to_time_range(p, 0));
137-
let t_lookback = match time_range {
133+
let t_lookback_ms = match time_range {
138134
Some(tr) => tr.duration_ms().unwrap_or(t_repeat_ms),
139135
None => t_repeat_ms,
140136
};
@@ -144,7 +140,7 @@ impl ElasticSingleQueryProcessor {
144140
None
145141
} else {
146142
Some(
147-
get_sql_cleanup_param(self.cleanup_policy, t_lookback, t_repeat_ms)
143+
get_sql_cleanup_param(self.cleanup_policy, t_lookback_ms, t_repeat_ms)
148144
.map_err(ControllerError::PlannerError)?,
149145
)
150146
};

asap-planner-rs/src/planner/sql.rs

Lines changed: 26 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -20,8 +20,8 @@ use crate::StreamingEngine;
2020

2121
pub struct SQLSingleQueryProcessor {
2222
query_string: String,
23-
t_repeat: u64,
24-
data_ingestion_interval: u64,
23+
t_repeat_ms: u64,
24+
data_ingestion_interval_ms: u64,
2525
table_definitions: Vec<TableDefinition>,
2626
#[allow(dead_code)]
2727
streaming_engine: StreamingEngine,
@@ -33,17 +33,17 @@ impl SQLSingleQueryProcessor {
3333
#[allow(clippy::too_many_arguments)]
3434
pub fn new(
3535
query_string: String,
36-
t_repeat: u64,
37-
data_ingestion_interval: u64,
36+
t_repeat_ms: u64,
37+
data_ingestion_interval_ms: u64,
3838
table_definitions: Vec<TableDefinition>,
3939
streaming_engine: StreamingEngine,
4040
sketch_parameters: Option<SketchParameterOverrides>,
4141
cleanup_policy: CleanupPolicy,
4242
) -> Self {
4343
Self {
4444
query_string,
45-
t_repeat,
46-
data_ingestion_interval,
45+
t_repeat_ms,
46+
data_ingestion_interval_ms,
4747
table_definitions,
4848
streaming_engine,
4949
sketch_parameters,
@@ -72,8 +72,12 @@ impl SQLSingleQueryProcessor {
7272
})?;
7373

7474
// Match query to pattern
75-
let sql_query = SQLPatternMatcher::new(schema.clone(), self.data_ingestion_interval as f64)
76-
.query_info_to_pattern(&qdata);
75+
// SQLPatternMatcher.scrape_interval is in seconds (SQL timestamps are seconds-based).
76+
let sql_query = SQLPatternMatcher::new(
77+
schema.clone(),
78+
self.data_ingestion_interval_ms as f64 / 1000.0,
79+
)
80+
.query_info_to_pattern(&qdata);
7781

7882
if !sql_query.is_valid() {
7983
return Err(ControllerError::SqlParse(sql_query.msg.unwrap_or_default()));
@@ -97,8 +101,11 @@ impl SQLSingleQueryProcessor {
97101
let value_column = agg_info.get_value_column_name().to_string();
98102

99103
// Compute window
100-
let window_cfg =
101-
compute_sql_window(query_type, self.data_ingestion_interval, self.t_repeat);
104+
let window_cfg = compute_sql_window(
105+
query_type,
106+
self.data_ingestion_interval_ms,
107+
self.t_repeat_ms,
108+
);
102109

103110
// Get all metadata columns for the table
104111
let all_metadata = get_all_metadata_columns(&self.table_definitions, table_name)?;
@@ -157,16 +164,17 @@ impl SQLSingleQueryProcessor {
157164
}
158165
}
159166

160-
let t_lookback = match query_type {
161-
QueryType::Spatial => self.data_ingestion_interval,
162-
_ => sql_query.query_data[0].time_info.get_duration() as u64,
167+
let t_lookback_ms = match query_type {
168+
QueryType::Spatial => self.data_ingestion_interval_ms,
169+
// SQLPatternParser always produces second-based durations; convert to ms.
170+
_ => (sql_query.query_data[0].time_info.get_duration() * 1000.0).round() as u64,
163171
};
164172

165173
let cleanup_param = if self.cleanup_policy == CleanupPolicy::NoCleanup {
166174
None
167175
} else {
168176
Some(
169-
get_sql_cleanup_param(self.cleanup_policy, t_lookback, self.t_repeat)
177+
get_sql_cleanup_param(self.cleanup_policy, t_lookback_ms, self.t_repeat_ms)
170178
.map_err(ControllerError::PlannerError)?,
171179
)
172180
};
@@ -215,16 +223,12 @@ fn get_sql_statistics(name: &str) -> Result<Vec<Statistic>, ControllerError> {
215223

216224
fn compute_sql_window(
217225
query_type: &QueryType,
218-
data_ingestion_interval: u64,
219-
t_repeat: u64,
226+
data_ingestion_interval_ms: u64,
227+
t_repeat_ms: u64,
220228
) -> IntermediateWindowConfig {
221-
// data_ingestion_interval and t_repeat are seconds (SQL mode is out of scope for the
222-
// ms-precision rename, see issue #427) — IntermediateWindowConfig is now ms-typed
223-
// throughout, so convert at this boundary, same as the pre-existing pattern in
224-
// asap-query-engine/src/engines/simple_engine/sql.rs.
225229
let window_size_ms = match query_type {
226-
QueryType::Spatial => data_ingestion_interval * 1000,
227-
_ => t_repeat * 1000,
230+
QueryType::Spatial => data_ingestion_interval_ms,
231+
_ => t_repeat_ms,
228232
};
229233
IntermediateWindowConfig {
230234
window_size_ms,

asap-planner-rs/src/promql/controller.rs

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,7 @@ impl Controller {
3737
) -> Result<Self, ControllerError> {
3838
let yaml_str = std::fs::read_to_string(path)?;
3939
let config: ControllerConfig = serde_yaml::from_str(&yaml_str)?;
40+
config.warn_default_slas();
4041
let all_queries: Vec<String> = config
4142
.query_groups
4243
.iter()
@@ -78,6 +79,7 @@ impl Controller {
7879
) -> Result<Self, ControllerError> {
7980
let yaml_str = std::fs::read_to_string(path)?;
8081
let config: ControllerConfig = serde_yaml::from_str(&yaml_str)?;
82+
config.warn_default_slas();
8183
Ok(Self {
8284
config,
8385
schema,
@@ -92,6 +94,7 @@ impl Controller {
9294
opts: RuntimeOptions,
9395
) -> Result<Self, ControllerError> {
9496
let config: ControllerConfig = serde_yaml::from_str(yaml)?;
97+
config.warn_default_slas();
9598
Ok(Self {
9699
config,
97100
schema,

0 commit comments

Comments
 (0)