diff --git a/README.md b/README.md index 5f0631b..caa2b7f 100644 --- a/README.md +++ b/README.md @@ -139,6 +139,11 @@ show benchmark queue ## Supported Benchmarks +There is no allowlist: any name you pass is scheduled and resolved on the +runner. Names that match a Criterion `[[bench]]` target run via Criterion; +everything else runs through `bench.sh`. A name that matches neither fails on +the runner. The tables below list the benchmarks known to work today. + ### DataFusion (standard) | Benchmark | Description | @@ -190,7 +195,7 @@ controller/ Rust controller crate github_poller.rs Comment polling loop job_manager.rs K8s Job lifecycle reconciler db.rs SQLite queries (jobs, seen comments, scan state) - benchmarks.rs Trigger parsing, benchmark allowlists + benchmarks.rs Trigger parsing (no allowlist — any name is accepted) migrations/ SQLite schema runner/ Benchmark runner container (builds project, runs benchmarks, posts results) queries/ SQL query files for ClickBench @@ -293,7 +298,7 @@ One-time setup for GitHub Actions → GCP authentication: | `github_poller` | Poll GitHub API, detect triggers, insert jobs | | `job_manager` | Create K8s Jobs, monitor status, post results | | `db` | SQLite persistence (jobs, seen comments, scan state) | -| `benchmarks` | Trigger parsing, per-repo allowlists, classification | +| `benchmarks` | Trigger parsing (no allowlist; benchmarks resolved on the runner) | | `github` | GitHub REST API client (comments, reactions) | | `config` | Environment-based configuration | diff --git a/controller/src/benchmarks.rs b/controller/src/benchmarks.rs index 1cf443f..2323570 100644 --- a/controller/src/benchmarks.rs +++ b/controller/src/benchmarks.rs @@ -1,8 +1,7 @@ -//! Benchmark trigger detection and per-repo allowlists. +//! Benchmark trigger detection. //! -//! Parses PR comment bodies for "run benchmark …" trigger phrases, validates -//! benchmark names against repo-specific allowlists, and classifies them by -//! [`JobType`]. +//! Parses PR comment bodies for "run benchmark …" trigger phrases. There is no +//! allowlist: any requested name is accepted and resolved on the runner. use std::collections::{HashMap, HashSet}; @@ -10,8 +9,7 @@ use once_cell::sync::Lazy; use regex::Regex; use serde::Deserialize; -use crate::config::RepoEntry; -use crate::models::{BenchmarkRequest, JobType}; +use crate::models::BenchmarkRequest; /// Unified trigger regex: matches `run benchmark(s) [name1 name2 ...]`. static TRIGGER_RE: Lazy = Lazy::new(|| { @@ -44,23 +42,6 @@ pub enum DetectResult { None, } -impl RepoEntry { - /// Determine the [`JobType`] for a benchmark name, or `None` if not recognized. - pub fn classify_benchmark(&self, name: &str) -> Option { - if self.standard_set().contains(name) { - Some(JobType::Standard) - } else if self.criterion_allows_any() || self.criterion_set().contains(name) { - if self.criterion_type == "arrow" { - Some(JobType::ArrowCriterion) - } else { - Some(JobType::Criterion) - } - } else { - None - } - } -} - /// Parse the extra lines (after the trigger line) into structured env vars and refs. /// /// Supports an optional ` ```yaml ` / ` ``` ` fence around the YAML content. @@ -165,10 +146,11 @@ pub fn parse_trigger(trigger: &str) -> Option { /// /// Recognizes `run benchmarks` (default suite), `run benchmarks `, and /// `run benchmark `. `run benchmark` without names returns `None` (caller -/// should post a help message). +/// should post a help message). Any requested names are accepted; there is no +/// allowlist. /// /// Supports `baseline:`/`changed:` sections with `env:` and `ref:` sub-entries. -pub fn detect_benchmark(repo_entry: &RepoEntry, body: &str) -> DetectResult { +pub fn detect_benchmark(body: &str) -> DetectResult { let lines: Vec<&str> = body.trim().lines().collect(); if lines.is_empty() { return DetectResult::None; @@ -202,26 +184,17 @@ pub fn detect_benchmark(repo_entry: &RepoEntry, body: &str) -> DetectResult { return DetectResult::None; } - let standard = repo_entry.standard_set(); - let criterion = repo_entry.criterion_set(); - let criterion_any = repo_entry.criterion_allows_any(); - - let all_valid = names.iter().all(|n| { - standard.contains(n.as_str()) || criterion_any || criterion.contains(n.as_str()) - }); - - if all_valid { - DetectResult::Parsed(BenchmarkRequest { - benchmarks: names, - env_vars: shared_env, - baseline_env_vars: baseline_env, - changed_env_vars: changed_env, - baseline_ref, - changed_ref, - }) - } else { - DetectResult::None - } + // No allowlist: accept any requested names. Names that resolve to + // neither a Criterion bench target nor a `bench.sh` suite simply + // fail on the runner. + DetectResult::Parsed(BenchmarkRequest { + benchmarks: names, + env_vars: shared_env, + baseline_env_vars: baseline_env, + changed_env_vars: changed_env, + baseline_ref, + changed_ref, + }) } TriggerKind::SingularNoNames => DetectResult::None, } @@ -249,65 +222,23 @@ pub fn is_queue_request(body: &str) -> bool { body.trim().eq_ignore_ascii_case("show benchmark queue") } -/// Build a markdown message listing all valid benchmarks for a repo, highlighting any unsupported names. -pub fn supported_benchmarks_message(repo_entry: &RepoEntry, requested: &[String]) -> String { - let standard: Vec<&str> = { - let mut v: Vec<&str> = repo_entry.standard.iter().map(|s| s.as_str()).collect(); - v.sort(); - v - }; - let criterion: Vec<&str> = { - let mut v: Vec<&str> = repo_entry.criterion.iter().map(|s| s.as_str()).collect(); - v.sort(); - v - }; - - let standard_str = if standard.is_empty() { - "(none)".to_string() - } else { - standard.join(", ") - }; - let criterion_any = repo_entry.criterion_allows_any(); - let criterion_str = if criterion_any { - "(any)".to_string() - } else if criterion.is_empty() { - "(none)".to_string() - } else { - criterion.join(", ") - }; - - let standard_set = repo_entry.standard_set(); - let criterion_set = repo_entry.criterion_set(); - - let bad: Vec<&String> = requested - .iter() - .filter(|n| { - !standard_set.contains(n.as_str()) - && !criterion_any - && !criterion_set.contains(n.as_str()) - }) - .collect(); - - let unsupported = if bad.is_empty() { - String::new() - } else { - format!( - "\nUnsupported benchmarks: {}.", - bad.iter() - .map(|s| s.as_str()) - .collect::>() - .join(", ") - ) - }; - - format!( - "Supported benchmarks:\n- Standard: {standard_str}\n- Criterion: {criterion_str}\n\n\ - Usage:\n\ +/// Build the usage/help message shown for malformed triggers and config errors. +/// +/// There is no allowlist, so this no longer enumerates valid benchmarks — any +/// name is accepted and resolved on the runner (an unresolvable name fails +/// there). The benchmark `bench.sh` suites and Criterion benches available are +/// whatever the target repo defines. +pub fn usage_message() -> String { + "Usage:\n\ ```\n\ run benchmark # run specific benchmark(s)\n\ run benchmarks # run default suite\n\ run benchmarks # run specific benchmarks\n\ - ```\n\n\ + ```\n\ + Any benchmark name is accepted: `bench.sh` suite names (e.g. `tpch`, \ + `clickbench_partitioned`, `wide_schema`) and Criterion bench targets \ + (e.g. `sql_planner`) are resolved automatically. A name that matches \ + neither fails on the runner.\n\n\ Per-side configuration (`run benchmark tpch` followed by):\n\ ```yaml\n\ env:\n\ @@ -325,8 +256,8 @@ pub fn supported_benchmarks_message(repo_entry: &RepoEntry, requested: &[String] ref: v46.0.0\n\ env:\n\ DATAFUSION_RUNTIME_MEMORY_LIMIT: 2G\n\ - ```{unsupported}" - ) + ```" + .to_string() } /// Format the allowlist as a comma-separated list of GitHub profile links. @@ -343,30 +274,12 @@ pub fn allowed_users_markdown(allowed_users: &HashSet) -> String { #[cfg(test)] mod tests { use super::*; + use crate::config::RepoEntry; + use crate::models::JobType; fn df_entry() -> RepoEntry { RepoEntry { - standard: vec![ - "tpch".into(), - "tpch10".into(), - "tpch_mem".into(), - "tpch_mem10".into(), - "topk_tpch".into(), - "clickbench_partitioned".into(), - "clickbench_extended".into(), - "clickbench_1".into(), - "clickbench_pushdown".into(), - "external_aggr".into(), - "tpcds".into(), - "smj".into(), - "sort_pushdown".into(), - "sort_pushdown_sorted".into(), - "sort_pushdown_inexact".into(), - "sort_pushdown_inexact_unsorted".into(), - "sort_pushdown_inexact_overlap".into(), - ], - criterion: vec!["sql_planner".into(), "in_list".into(), "case_when".into()], - criterion_type: "datafusion".into(), + kind: "datafusion".into(), default_standard: vec![ "clickbench_partitioned".into(), "tpcds".into(), @@ -377,22 +290,11 @@ mod tests { fn arrow_entry() -> RepoEntry { RepoEntry { - standard: vec![], - criterion: vec!["arrow_reader".into(), "arrow_writer".into()], - criterion_type: "arrow".into(), + kind: "arrow".into(), default_standard: vec![], } } - fn wildcard_criterion_entry() -> RepoEntry { - RepoEntry { - standard: vec!["tpch".into()], - criterion: vec!["*".into()], - criterion_type: "datafusion".into(), - default_standard: vec!["tpch".into()], - } - } - // ── detect_benchmark ──────────────────────────────────────────── /// Helper to unwrap a DetectResult::Parsed or panic. @@ -414,7 +316,7 @@ mod tests { #[test] fn detect_default_suite() { - let req = unwrap_parsed(detect_benchmark(&df_entry(), "run benchmarks")); + let req = unwrap_parsed(detect_benchmark("run benchmarks")); assert!(req.benchmarks.is_empty()); assert!(req.env_vars.is_empty()); } @@ -422,7 +324,7 @@ mod tests { #[test] fn detect_default_suite_with_env_vars() { let body = "run benchmarks\nenv:\n DATAFUSION_RUNTIME_MEMORY_LIMIT: 1G"; - let req = unwrap_parsed(detect_benchmark(&df_entry(), body)); + let req = unwrap_parsed(detect_benchmark(body)); assert!(req.benchmarks.is_empty()); assert_eq!( req.env_vars.get("DATAFUSION_RUNTIME_MEMORY_LIMIT").unwrap(), @@ -432,77 +334,53 @@ mod tests { #[test] fn detect_single_named() { - let req = unwrap_parsed(detect_benchmark(&df_entry(), "run benchmark tpch_mem")); + let req = unwrap_parsed(detect_benchmark("run benchmark tpch_mem")); assert_eq!(req.benchmarks, vec!["tpch_mem"]); } #[test] fn detect_multiple_named() { - let req = unwrap_parsed(detect_benchmark( - &df_entry(), - "run benchmark tpch_mem tpch10", - )); + let req = unwrap_parsed(detect_benchmark("run benchmark tpch_mem tpch10")); assert_eq!(req.benchmarks, vec!["tpch_mem", "tpch10"]); } #[test] fn detect_criterion_benchmark() { - let req = unwrap_parsed(detect_benchmark(&df_entry(), "run benchmark sql_planner")); + let req = unwrap_parsed(detect_benchmark("run benchmark sql_planner")); assert_eq!(req.benchmarks, vec!["sql_planner"]); } #[test] - fn detect_bogus_name_returns_none() { - assert!(is_none(&detect_benchmark( - &df_entry(), - "run benchmark bogus_name" - ))); - } + fn detect_any_name_is_accepted() { + // No allowlist: previously-unknown names now parse and are scheduled. + let req = unwrap_parsed(detect_benchmark("run benchmark anything_goes")); + assert_eq!(req.benchmarks, vec!["anything_goes"]); - #[test] - fn detect_one_invalid_rejects_all() { - assert!(is_none(&detect_benchmark( - &df_entry(), - "run benchmark tpch_mem bogus" - ))); + let req = unwrap_parsed(detect_benchmark("run benchmark tpch_mem bogus")); + assert_eq!(req.benchmarks, vec!["tpch_mem", "bogus"]); } #[test] fn detect_not_a_trigger() { - assert!(is_none(&detect_benchmark(&df_entry(), "hello world"))); + assert!(is_none(&detect_benchmark("hello world"))); } #[test] fn detect_empty_string() { - assert!(is_none(&detect_benchmark(&df_entry(), ""))); + assert!(is_none(&detect_benchmark(""))); } #[test] fn detect_case_insensitive() { - assert!(is_parsed(&detect_benchmark(&df_entry(), "Run Benchmarks"))); - assert!(is_parsed(&detect_benchmark( - &df_entry(), - "RUN BENCHMARK tpch" - ))); - } - - #[test] - fn detect_arrow_criterion() { - let req = unwrap_parsed(detect_benchmark( - &arrow_entry(), - "run benchmark arrow_reader", - )); - assert_eq!(req.benchmarks, vec!["arrow_reader"]); + assert!(is_parsed(&detect_benchmark("Run Benchmarks"))); + assert!(is_parsed(&detect_benchmark("RUN BENCHMARK tpch"))); } // ── plural trigger with names (new) ───────────────────────────── #[test] fn detect_plural_with_names() { - let req = unwrap_parsed(detect_benchmark( - &df_entry(), - "run benchmarks tpch clickbench_1", - )); + let req = unwrap_parsed(detect_benchmark("run benchmarks tpch clickbench_1")); assert_eq!(req.benchmarks, vec!["tpch", "clickbench_1"]); } @@ -510,7 +388,7 @@ mod tests { #[test] fn detect_singular_no_names_returns_none() { - assert!(is_none(&detect_benchmark(&df_entry(), "run benchmark"))); + assert!(is_none(&detect_benchmark("run benchmark"))); } #[test] @@ -526,7 +404,7 @@ mod tests { #[test] fn parse_baseline_changed_env_vars() { let body = "run benchmark tpch\nbaseline:\n env:\n DATAFUSION_RUNTIME_MEMORY_LIMIT: 1G\nchanged:\n env:\n DATAFUSION_RUNTIME_MEMORY_LIMIT: 2G"; - let req = unwrap_parsed(detect_benchmark(&df_entry(), body)); + let req = unwrap_parsed(detect_benchmark(body)); assert_eq!( req.baseline_env_vars .get("DATAFUSION_RUNTIME_MEMORY_LIMIT") @@ -545,7 +423,7 @@ mod tests { #[test] fn parse_baseline_ref() { let body = "run benchmarks tpch clickbench_1\nbaseline:\n ref: abc1234def"; - let req = unwrap_parsed(detect_benchmark(&df_entry(), body)); + let req = unwrap_parsed(detect_benchmark(body)); assert_eq!(req.baseline_ref.as_deref(), Some("abc1234def")); assert!(req.changed_ref.is_none()); } @@ -553,7 +431,7 @@ mod tests { #[test] fn parse_both_refs_with_env() { let body = "run benchmark tpch\nbaseline:\n ref: v45.0.0\n env:\n FOO: old_value\nchanged:\n ref: v46.0.0\n env:\n FOO: new_value"; - let req = unwrap_parsed(detect_benchmark(&df_entry(), body)); + let req = unwrap_parsed(detect_benchmark(body)); assert_eq!(req.baseline_ref.as_deref(), Some("v45.0.0")); assert_eq!(req.changed_ref.as_deref(), Some("v46.0.0")); assert_eq!(req.baseline_env_vars.get("FOO").unwrap(), "old_value"); @@ -563,7 +441,7 @@ mod tests { #[test] fn parse_shared_plus_per_side() { let body = "run benchmark tpch\nenv:\n SHARED_SETTING: enabled\nbaseline:\n env:\n DATAFUSION_RUNTIME_MEMORY_LIMIT: 1G\nchanged:\n env:\n DATAFUSION_RUNTIME_MEMORY_LIMIT: 2G"; - let req = unwrap_parsed(detect_benchmark(&df_entry(), body)); + let req = unwrap_parsed(detect_benchmark(body)); assert_eq!(req.env_vars.get("SHARED_SETTING").unwrap(), "enabled"); assert_eq!( req.baseline_env_vars @@ -582,7 +460,7 @@ mod tests { #[test] fn parse_explicit_env_section() { let body = "run benchmark tpch\nenv:\n DATAFUSION_RUNTIME_MEMORY_LIMIT: 1G"; - let req = unwrap_parsed(detect_benchmark(&df_entry(), body)); + let req = unwrap_parsed(detect_benchmark(body)); assert_eq!( req.env_vars.get("DATAFUSION_RUNTIME_MEMORY_LIMIT").unwrap(), "1G" @@ -592,7 +470,7 @@ mod tests { #[test] fn parse_yaml_fenced_block() { let body = "run benchmark tpch\n```yaml\nbaseline:\n ref: v45.0.0\n env:\n FOO: bar\nchanged:\n ref: v46.0.0\n```"; - let req = unwrap_parsed(detect_benchmark(&df_entry(), body)); + let req = unwrap_parsed(detect_benchmark(body)); assert_eq!(req.baseline_ref.as_deref(), Some("v45.0.0")); assert_eq!(req.changed_ref.as_deref(), Some("v46.0.0")); assert_eq!(req.baseline_env_vars.get("FOO").unwrap(), "bar"); @@ -601,7 +479,7 @@ mod tests { #[test] fn parse_unknown_field_returns_config_error() { let body = "run benchmark tpch\ncurrent:\n ref: HEAD"; - match detect_benchmark(&df_entry(), body) { + match detect_benchmark(body) { DetectResult::ConfigError(e) => { assert!(e.contains("unknown field"), "error was: {e}"); } @@ -616,6 +494,27 @@ mod tests { } } + // ── RepoEntry::job_type ───────────────────────────────────────── + + #[test] + fn job_type_datafusion_repo() { + assert_eq!(df_entry().job_type(), JobType::Datafusion); + } + + #[test] + fn job_type_arrow_repo() { + assert_eq!(arrow_entry().job_type(), JobType::ArrowCriterion); + } + + // ── usage_message ─────────────────────────────────────────────── + + #[test] + fn usage_message_has_usage_and_no_allowlist() { + let msg = usage_message(); + assert!(msg.contains("run benchmark")); + assert!(msg.contains("Any benchmark name is accepted")); + } + // ── is_benchmark_trigger ──────────────────────────────────────── #[test] @@ -665,86 +564,6 @@ mod tests { assert!(!is_queue_request("run benchmarks")); } - // ── RepoEntry::classify_benchmark ────────────────────────────── - - #[test] - fn classify_df_standard() { - assert_eq!( - df_entry().classify_benchmark("tpch"), - Some(JobType::Standard) - ); - } - - #[test] - fn classify_df_criterion() { - assert_eq!( - df_entry().classify_benchmark("sql_planner"), - Some(JobType::Criterion) - ); - } - - #[test] - fn classify_df_bogus() { - assert_eq!(df_entry().classify_benchmark("bogus"), None); - } - - #[test] - fn classify_arrow_criterion() { - assert_eq!( - arrow_entry().classify_benchmark("arrow_reader"), - Some(JobType::ArrowCriterion) - ); - } - - // ── wildcard criterion ────────────────────────────────────────── - - #[test] - fn detect_wildcard_criterion_accepts_any_name() { - let req = unwrap_parsed(detect_benchmark( - &wildcard_criterion_entry(), - "run benchmark anything_goes", - )); - assert_eq!(req.benchmarks, vec!["anything_goes"]); - } - - #[test] - fn classify_wildcard_criterion() { - assert_eq!( - wildcard_criterion_entry().classify_benchmark("any_bench_name"), - Some(JobType::Criterion) - ); - } - - #[test] - fn classify_wildcard_standard_still_works() { - assert_eq!( - wildcard_criterion_entry().classify_benchmark("tpch"), - Some(JobType::Standard) - ); - } - - #[test] - fn supported_msg_wildcard_shows_any() { - let msg = - supported_benchmarks_message(&wildcard_criterion_entry(), &["unknown".to_string()]); - assert!(msg.contains("(any)")); - assert!(!msg.contains("Unsupported")); - } - - // ── supported_benchmarks_message ──────────────────────────────── - - #[test] - fn supported_msg_no_unsupported() { - let msg = supported_benchmarks_message(&df_entry(), &[]); - assert!(!msg.contains("Unsupported")); - } - - #[test] - fn supported_msg_with_unsupported() { - let msg = supported_benchmarks_message(&df_entry(), &["bogus".to_string()]); - assert!(msg.contains("Unsupported benchmarks: bogus")); - } - // ── allowed_users_markdown ────────────────────────────────────── #[test] diff --git a/controller/src/bin/runner.rs b/controller/src/bin/runner.rs index cf09fb2..3c47904 100644 --- a/controller/src/bin/runner.rs +++ b/controller/src/bin/runner.rs @@ -9,7 +9,7 @@ use tracing::{error, info}; use benchmark_controller::github; use benchmark_controller::runner::config::{BenchType, RunnerConfig}; use benchmark_controller::runner::poster::CommentPoster; -use benchmark_controller::runner::{bench_arrow, bench_criterion, bench_standard, shell}; +use benchmark_controller::runner::{bench_arrow, bench_datafusion, shell}; #[tokio::main] async fn main() { @@ -58,8 +58,9 @@ async fn main() { async fn run_benchmark(config: &RunnerConfig, poster: &CommentPoster) -> Result<()> { match config.bench_type { - BenchType::Standard | BenchType::MainTracking => bench_standard::run(config, poster).await, - BenchType::Criterion => bench_criterion::run(config, poster).await, + BenchType::Datafusion | BenchType::MainTracking => { + bench_datafusion::run(config, poster).await + } BenchType::ArrowCriterion => bench_arrow::run(config, poster).await, } } diff --git a/controller/src/config.rs b/controller/src/config.rs index acf86fd..c392bef 100644 --- a/controller/src/config.rs +++ b/controller/src/config.rs @@ -5,6 +5,8 @@ use std::collections::{HashMap, HashSet}; use anyhow::{Context, Result}; use serde::Deserialize; +use crate::models::JobType; + /// Maximum benchmark jobs a single user can have in the `running` state at once. /// Enforced at pickup in `db::get_pending_jobs`, so pending jobs sit in the /// queue until an earlier run finishes. @@ -15,38 +17,38 @@ pub const MAX_RUNNING_PER_USER: i64 = 5; /// to wait instead of silently dropping the request. pub const MAX_QUEUED_PER_USER: i64 = 15; -/// Per-repo benchmark allowlists loaded from JSON config. +/// Per-repo benchmark configuration loaded from JSON config. +/// +/// There is no benchmark allowlist: any name a user requests is scheduled and +/// resolved on the runner (a name that resolves to neither a Criterion bench +/// target nor a `bench.sh` suite simply fails there). `kind` selects how the +/// repo is benchmarked. #[derive(Debug, Clone, Deserialize)] pub struct RepoEntry { - #[serde(default)] - pub standard: Vec, - #[serde(default)] - pub criterion: Vec, - /// `"datafusion"` (default) or `"arrow"` — controls criterion `JobType`. - #[serde(default = "default_criterion_type")] - pub criterion_type: String, - /// Default standard benchmarks for `"run benchmarks"` (no specific names). + /// `"datafusion"` (default) — `bench.sh` suites plus Criterion benches, + /// resolved per-benchmark on the runner — or `"arrow"` — arrow-rs Criterion + /// benches only. + #[serde(default = "default_kind")] + pub kind: String, + /// Default benchmarks for `"run benchmarks"` (no specific names). #[serde(default)] pub default_standard: Vec, } -fn default_criterion_type() -> String { +fn default_kind() -> String { "datafusion".to_string() } impl RepoEntry { - pub fn standard_set(&self) -> HashSet<&str> { - self.standard.iter().map(|s| s.as_str()).collect() - } - - pub fn criterion_set(&self) -> HashSet<&str> { - self.criterion.iter().map(|s| s.as_str()).collect() - } - - /// Returns `true` if the criterion list contains `"*"`, meaning any - /// criterion benchmark name is accepted. - pub fn criterion_allows_any(&self) -> bool { - self.criterion.iter().any(|s| s == "*") + /// The [`JobType`] used to run every benchmark for this repo. The + /// datafusion runner resolves each name (Criterion vs `bench.sh`) itself, + /// so there is no per-benchmark classification. + pub fn job_type(&self) -> JobType { + if self.kind == "arrow" { + JobType::ArrowCriterion + } else { + JobType::Datafusion + } } } diff --git a/controller/src/github_poller.rs b/controller/src/github_poller.rs index 4289f9f..e69c5b6 100644 --- a/controller/src/github_poller.rs +++ b/controller/src/github_poller.rs @@ -10,7 +10,7 @@ use tracing::{info, warn}; use crate::benchmarks::{ allowed_users_markdown, detect_benchmark, is_benchmark_trigger, is_queue_request, - is_singular_no_names, supported_benchmarks_message, DetectResult, + is_singular_no_names, usage_message, DetectResult, }; use crate::config::{BenchmarkConfig, Config, RepoEntry, MAX_QUEUED_PER_USER}; use crate::db; @@ -208,7 +208,7 @@ async fn process_comment( } // Try to detect benchmark trigger - let request = match detect_benchmark(repo_entry, body) { + let request = match detect_benchmark(body) { DetectResult::Parsed(req) => req, DetectResult::ConfigError(err) => { // YAML config was present but invalid — post a helpful error @@ -216,7 +216,7 @@ async fn process_comment( let msg = format!( "Hi @{login}, your benchmark configuration could not be parsed ({comment_url}).\n\n\ **Error:** `{err}`\n\n{}{footer}", - supported_benchmarks_message(repo_entry, &[]) + usage_message() ); gh.post_comment(repo, pr_number, &msg).await?; } @@ -224,7 +224,9 @@ async fn process_comment( return Ok(()); } DetectResult::None => { - // Check if it looks like a failed trigger attempt + // With no allowlist, named/default triggers always parse, so the + // only `None` cases that look like a trigger are `run benchmark` + // (singular) with no names. if is_benchmark_trigger(body) { if !bench_cfg.allowed_users.contains(login) { let msg = not_allowed_message( @@ -235,16 +237,6 @@ async fn process_comment( ); gh.post_comment(repo, pr_number, &msg).await?; } else { - // Singular "run benchmark" with no names, or invalid names - let requested: Vec = body - .trim() - .lines() - .next() - .unwrap_or("") - .split_whitespace() - .skip(2) - .map(|s| s.to_string()) - .collect(); let prefix = if is_singular_no_names(body) { format!( "Hi @{login}, `run benchmark` requires benchmark names ({comment_url}).\n\n" @@ -252,10 +244,7 @@ async fn process_comment( } else { format!("Hi @{login}, thanks for the request ({comment_url}).\n\n") }; - let msg = format!( - "{prefix}{}{footer}", - supported_benchmarks_message(repo_entry, &requested) - ); + let msg = format!("{prefix}{}{footer}", usage_message()); gh.post_comment(repo, pr_number, &msg).await?; } } @@ -307,9 +296,12 @@ async fn process_comment( let baseline_env_json = serde_json::to_string(&request.baseline_env_vars)?; let changed_env_json = serde_json::to_string(&request.changed_env_vars)?; - // Determine job type(s) and insert jobs + // The job type is per-repo (the datafusion runner resolves each name + // itself), so there is no per-benchmark classification. + let job_type = repo_entry.job_type().as_str(); + if benchmarks.is_empty() { - // No defaults configured — insert a single standard job (runner uses its own default) + // No defaults configured — insert a single job (runner uses its own default) let benchmarks_json = serde_json::to_string(&benchmarks)?; db::insert_job( pool, @@ -325,18 +317,13 @@ async fn process_comment( changed_env_vars: &changed_env_json, baseline_ref: request.baseline_ref.as_deref(), changed_ref: request.changed_ref.as_deref(), - job_type: "standard", + job_type, }, ) .await?; } else { - // Group benchmarks by type + // One job per benchmark so each runs in its own pod. for bench in &benchmarks { - let job_type = repo_entry - .classify_benchmark(bench) - .map(|jt| jt.as_str()) - .unwrap_or("standard"); - let single_bench = serde_json::to_string(&[bench])?; db::insert_job( pool, diff --git a/controller/src/job_manager.rs b/controller/src/job_manager.rs index 60b2f95..836a668 100644 --- a/controller/src/job_manager.rs +++ b/controller/src/job_manager.rs @@ -426,9 +426,9 @@ async fn create_k8s_job( env.push(env_var("CHANGED_REF", changed_ref)); } - // For criterion benchmarks, set BENCH_NAME to the first benchmark - if (job.job_type == "criterion" || job.job_type == "arrow_criterion") && !benchmarks.is_empty() - { + // The arrow runner benches a single named target; the datafusion runner + // reads the BENCHMARKS list and resolves each name itself. + if job.job_type == "arrow_criterion" && !benchmarks.is_empty() { env.push(env_var("BENCH_NAME", benchmarks[0].clone())); } diff --git a/controller/src/models.rs b/controller/src/models.rs index 4e5d6af..7f7dbf9 100644 --- a/controller/src/models.rs +++ b/controller/src/models.rs @@ -74,13 +74,13 @@ pub struct BenchmarkRequest { /// Benchmark runner variant. /// -/// - `Standard` — shell-based DataFusion benchmarks (tpch, clickbench, etc.) -/// - `Criterion` — `cargo bench` criterion benchmarks in DataFusion -/// - `ArrowCriterion` — `cargo bench` criterion benchmarks in arrow-rs +/// - `Datafusion` — DataFusion benchmarks. The runner resolves each requested +/// name per-benchmark: real `cargo bench` Criterion targets run via Criterion, +/// everything else runs through `bench.sh`. +/// - `ArrowCriterion` — `cargo bench` criterion benchmarks in arrow-rs. #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum JobType { - Standard, - Criterion, + Datafusion, ArrowCriterion, } @@ -88,8 +88,7 @@ impl JobType { /// Returns the string stored in the `job_type` SQLite column. pub fn as_str(&self) -> &'static str { match self { - Self::Standard => "standard", - Self::Criterion => "criterion", + Self::Datafusion => "datafusion", Self::ArrowCriterion => "arrow_criterion", } } @@ -168,13 +167,8 @@ mod tests { // ── JobType::as_str ───────────────────────────────────────────── #[test] - fn job_type_standard() { - assert_eq!(JobType::Standard.as_str(), "standard"); - } - - #[test] - fn job_type_criterion() { - assert_eq!(JobType::Criterion.as_str(), "criterion"); + fn job_type_datafusion() { + assert_eq!(JobType::Datafusion.as_str(), "datafusion"); } #[test] diff --git a/controller/src/runner/bench_criterion.rs b/controller/src/runner/bench_criterion.rs deleted file mode 100644 index 6105047..0000000 --- a/controller/src/runner/bench_criterion.rs +++ /dev/null @@ -1,441 +0,0 @@ -//! Criterion benchmark runner — ports `run_criterion.sh`. - -use std::path::{Path, PathBuf}; - -use anyhow::{Context, Result}; -use tracing::{info, warn}; - -use crate::github; -use crate::runner::config::RunnerConfig; -use crate::runner::git; -use crate::runner::monitor; -use crate::runner::poster::CommentPoster; -use crate::runner::shell; - -/// Run a criterion benchmark comparing a PR branch to its merge-base. -pub async fn run(config: &RunnerConfig, poster: &CommentPoster) -> Result<()> { - let repo_url = config.repo_url(); - let bench_name = &config.bench_name; - let bench_filter = &config.bench_filter; - let bench_command_args = bench_command_args(bench_name); - - let branch_dir = PathBuf::from("/workspace/datafusion-branch"); - let base_dir = PathBuf::from("/workspace/datafusion-base"); - - // Clone and checkout PR branch - info!("=== Cloning PR branch ==="); - git::clone_shallow(&repo_url, &branch_dir, 200).await?; - let branch_name = git::checkout_pr(&config.pr_url, &branch_dir).await?; - let merge_base = git::merge_base(&branch_dir).await?; - let bench_branch_name = git::sanitize_branch_name(&branch_name); - - // If a custom changed ref is specified, checkout that instead of PR head - if let Some(ref changed_ref) = config.changed_ref { - info!(changed_ref, "=== Checking out custom changed ref ==="); - git::fetch_pr_ref(&config.pr_url, &branch_dir).await?; - git::fetch_origin(&branch_dir).await?; - git::checkout(&branch_dir, changed_ref).await?; - } - - // Determine baseline: custom ref or merge-base - let baseline_display: String; - info!("=== Cloning merge-base ==="); - git::clone_shallow(&repo_url, &base_dir, 200).await?; - if let Some(ref baseline_ref) = config.baseline_ref { - info!(baseline_ref, "=== Checking out custom baseline ref ==="); - git::fetch_pr_ref(&config.pr_url, &base_dir).await?; - git::fetch_origin(&base_dir).await?; - git::checkout(&base_dir, baseline_ref).await?; - baseline_display = baseline_ref.clone(); - } else { - git::checkout(&base_dir, &merge_base).await?; - baseline_display = merge_base.clone(); - } - - // Pre-install stable toolchain to avoid rustup race in parallel builds - git::rustup_stable().await?; - - // Set up required benchmark data (e.g. sql_planner needs clickbench_partitioned) - setup_benchmark_data(bench_name, &branch_dir, &base_dir).await; - - // Post "running" comment - let uname = shell::uname().await; - let instance_type = shell::node_instance_type().await; - let pod_resources = shell::pod_resources(); - let lscpu = shell::lscpu().await; - let bench_command_display = format!("cargo bench --features=parquet --bench {bench_name}"); - let changed_display = config.changed_ref.as_deref().unwrap_or(&branch_name); - let changed_sha = git::rev_parse_head(&branch_dir).await?; - let base_sha = git::rev_parse_head(&base_dir).await?; - let baseline_label = if config.baseline_ref.is_some() { - baseline_display.clone() - } else { - format!("{} (merge-base)", &base_sha[..7.min(base_sha.len())]) - }; - let footer = github::issues_footer(config.runner_repo_url.as_deref()); - let running_body = format!( - "\u{1f916} Criterion benchmark running (GKE) | [trigger]({})\n\ - **Instance:** `{instance_type}` ({pod_resources}) | `{uname}`\n\ -
CPU Details (lscpu)\n\n\ - ```\n\ - {lscpu}\n\ - ```\n\n\ -
\n\n\ - Comparing {changed_display} ({changed_sha}) to {baseline_label} \ - [diff](https://github.com/{repo}/compare/{base_sha}..{changed_sha})\n\ - BENCH_NAME={bench_name}\n\ - BENCH_COMMAND={bench_command_display}\n\ - BENCH_FILTER={bench_filter}\n\ - Results will be posted here when complete{footer}", - config.comment_url, - repo = config.repo, - ); - let pr_number = config.pr_number()?; - poster - .post_comment(&config.repo, pr_number, &running_body) - .await?; - - // Compile both in parallel - info!("=== Compiling PR branch and merge-base in parallel ==="); - let mut branch_args = bench_command_args.clone(); - branch_args.push("--no-run".to_string()); - let mut base_args = bench_command_args.clone(); - base_args.push("--no-run".to_string()); - - let branch_build = shell::spawn_command( - "cargo", - &str_slice(&branch_args), - &branch_dir, - "/tmp/branch_build.log", - ); - let base_build = shell::spawn_command( - "cargo", - &str_slice(&base_args), - &base_dir, - "/tmp/base_build.log", - ); - - branch_build - .await - .context("branch build task panicked")? - .context("branch build failed")?; - let baseline_available = match base_build.await { - Ok(Ok(())) => true, - Ok(Err(e)) => { - warn!("Baseline build failed (benchmark may be new): {e:#}"); - false - } - Err(e) => { - warn!("Baseline build task panicked (benchmark may be new): {e:#}"); - false - } - }; - info!("=== Compilation complete ==="); - - // Run benchmarks sequentially, applying per-side env vars via `env` wrapper - let base_stats = if baseline_available { - info!("=== Running benchmark on merge-base ==="); - let mut base_run_args = bench_command_args.clone(); - base_run_args.extend(["--", "--save-baseline", "main"].map(String::from)); - if !bench_filter.is_empty() { - base_run_args.push(bench_filter.clone()); - } - let baseline_extra_env = config.baseline_env_args(); - let (_, stats) = if baseline_extra_env.is_empty() { - shell::run_command_monitored("cargo", &str_slice(&base_run_args), &base_dir, None) - .await? - } else { - let mut env_args: Vec = baseline_extra_env; - env_args.push("cargo".to_string()); - env_args.extend(base_run_args); - shell::run_command_monitored("env", &str_slice(&env_args), &base_dir, None).await? - }; - Some(stats) - } else { - info!("=== Skipping merge-base benchmark (baseline build failed) ==="); - None - }; - - info!("=== Running benchmark on PR branch ==="); - let mut branch_run_args = bench_command_args.clone(); - branch_run_args.extend(["--", "--save-baseline"].map(String::from)); - branch_run_args.push(bench_branch_name.clone()); - if !bench_filter.is_empty() { - branch_run_args.push(bench_filter.clone()); - } - let changed_extra_env = config.changed_env_args(); - let (_, branch_stats) = if changed_extra_env.is_empty() { - shell::run_command_monitored("cargo", &str_slice(&branch_run_args), &branch_dir, None) - .await? - } else { - let mut env_args: Vec = changed_extra_env; - env_args.push("cargo".to_string()); - env_args.extend(branch_run_args); - shell::run_command_monitored("env", &str_slice(&env_args), &branch_dir, None).await? - }; - - // Compare and post results - let result_body = if baseline_available { - // Copy baselines into one target dir for critcmp - copy_criterion_baselines(&base_dir, &branch_dir).await; - - let report = shell::run_command("critcmp", &["main", &bench_branch_name], &branch_dir) - .await - .context("critcmp")?; - - let resource_section = format!( - "{}\n{}", - monitor::format_resource_comment("base (merge-base)", &base_stats.unwrap()), - monitor::format_resource_comment("branch", &branch_stats), - ); - format_result_comment( - &config.comment_url, - &report, - &resource_section, - &instance_type, - &pod_resources, - &lscpu, - &footer, - ) - } else { - let report = shell::run_command("critcmp", &[bench_branch_name.as_str()], &branch_dir) - .await - .context("critcmp")?; - - let resource_section = - monitor::format_resource_comment("branch", &branch_stats).to_string(); - format_branch_only_result_comment( - &config.comment_url, - &report, - &resource_section, - &instance_type, - &pod_resources, - &lscpu, - &footer, - ) - }; - poster - .post_comment(&config.repo, pr_number, &result_body) - .await?; - - Ok(()) -} - -/// Build cargo bench args for criterion. -fn bench_command_args(bench_name: &str) -> Vec { - vec![ - "bench".to_string(), - "--features=parquet".to_string(), - "--bench".to_string(), - bench_name.to_string(), - ] -} - -/// Copy criterion baselines from base to branch target directory. -pub(crate) async fn copy_criterion_baselines(base_dir: &Path, branch_dir: &Path) { - let src = base_dir.join("target/criterion"); - let dst = branch_dir.join("target/criterion"); - if src.exists() { - let _ = shell::run_command( - "cp", - &[ - "-r", - &format!("{}/.", src.to_string_lossy()), - &dst.to_string_lossy(), - ], - Path::new("/"), - ) - .await; - } -} - -/// Format the result comment body. -fn format_result_comment( - comment_url: &str, - report: &str, - resource_section: &str, - instance_type: &str, - pod_resources: &str, - lscpu: &str, - footer: &str, -) -> String { - format!( - "\u{1f916} Criterion benchmark completed (GKE) | [trigger]({comment_url})\n\n\ - **Instance:** `{instance_type}` ({pod_resources})\n\n\ -
CPU Details (lscpu)\n\n\ - ```\n\ - {lscpu}\n\ - ```\n\n\ -
\n\n\ -
Details\n\ -

\n\n\ - ```\n\ - {report}\ - ```\n\n\ -

\n\ -
\n\n\ -
Resource Usage\n\n\ - {resource_section}\ -
\n\ - {footer}" - ) -} - -/// Format the result comment body for branch-only runs (no baseline comparison). -fn format_branch_only_result_comment( - comment_url: &str, - report: &str, - resource_section: &str, - instance_type: &str, - pod_resources: &str, - lscpu: &str, - footer: &str, -) -> String { - format!( - "\u{1f916} Criterion benchmark completed (GKE) | [trigger]({comment_url})\n\n\ - **Instance:** `{instance_type}` ({pod_resources})\n\n\ -
CPU Details (lscpu)\n\n\ - ```\n\ - {lscpu}\n\ - ```\n\n\ -
\n\n\ - **New benchmark — branch-only results (no baseline comparison)**\n\n\ -
Details\n\ -

\n\n\ - ```\n\ - {report}\ - ```\n\n\ -

\n\ -
\n\n\ -
Resource Usage\n\n\ - {resource_section}\ -
\n\ - {footer}" - ) -} - -/// Return the dataset names required by a given criterion benchmark. -fn required_datasets(bench_name: &str) -> &'static [&'static str] { - match bench_name { - "sql_planner" => &["clickbench_partitioned"], - _ => &[], - } -} - -/// Set up benchmark data for criterion benchmarks that need it. -/// -/// Generates data once in the branch directory and symlinks it into the base -/// directory so both sides share the same input data. -async fn setup_benchmark_data(bench_name: &str, branch_dir: &Path, base_dir: &Path) { - let datasets = required_datasets(bench_name); - if datasets.is_empty() { - return; - } - - let branch_benchmarks = branch_dir.join("benchmarks"); - let bench_dir_str = branch_benchmarks.to_string_lossy().to_string(); - - for dataset in datasets { - info!("Setting up data for {dataset}"); - cache_data(dataset, &bench_dir_str).await; - } - - // Symlink the data directory into the base checkout so both sides can - // find it without duplicating storage. - let branch_data = branch_benchmarks.join("data"); - let base_data = base_dir.join("benchmarks/data"); - if branch_data.exists() && !base_data.exists() { - info!("Symlinking benchmark data into base directory"); - let _ = tokio::fs::symlink(&branch_data, &base_data).await; - } -} - -/// Run data generation with cache support via /scripts/cache_data.sh. -async fn cache_data(bench: &str, bench_dir: &str) { - let cache_script = Path::new("/scripts/cache_data.sh"); - if cache_script.exists() { - let _ = shell::run_command( - cache_script.to_str().unwrap(), - &[bench, bench_dir], - Path::new(bench_dir), - ) - .await; - } else { - let _ = shell::run_command("./bench.sh", &["data", bench], Path::new(bench_dir)).await; - } -} - -/// Convert Vec to a slice of &str for run_command. -fn str_slice(v: &[String]) -> Vec<&str> { - v.iter().map(|s| s.as_str()).collect() -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn required_datasets_sql_planner() { - assert_eq!( - required_datasets("sql_planner"), - &["clickbench_partitioned"] - ); - } - - #[test] - fn required_datasets_unknown_returns_empty() { - assert!(required_datasets("unknown_bench").is_empty()); - } - - #[test] - fn bench_args_construction() { - let args = bench_command_args("sql_planner"); - assert_eq!( - args, - vec!["bench", "--features=parquet", "--bench", "sql_planner"] - ); - } - - #[test] - fn result_comment_format() { - let comment = format_result_comment( - "https://example.com/comment", - "test report\n", - "resources\n", - "c4a-standard-48", - "12 vCPU / 65 GiB", - "lscpu output", - "", - ); - assert!(comment.contains("Criterion benchmark completed")); - assert!(comment.contains("[trigger](https://example.com/comment)")); - assert!(comment.contains("test report")); - assert!(comment.contains("
")); - assert!(comment.contains("Resource Usage")); - assert!(comment.contains("c4a-standard-48")); - assert!(comment.contains("12 vCPU / 65 GiB")); - assert!(comment.contains("lscpu output")); - } - - #[test] - fn branch_only_result_comment_format() { - let comment = format_branch_only_result_comment( - "https://example.com/comment", - "branch report\n", - "branch resources\n", - "c4a-standard-48", - "12 vCPU / 65 GiB", - "lscpu output", - "", - ); - assert!(comment.contains("Criterion benchmark completed")); - assert!(comment.contains("New benchmark — branch-only results")); - assert!(comment.contains("[trigger](https://example.com/comment)")); - assert!(comment.contains("branch report")); - assert!(comment.contains("Resource Usage")); - assert!(comment.contains("branch resources")); - assert!(comment.contains("c4a-standard-48")); - assert!(comment.contains("12 vCPU / 65 GiB")); - assert!(comment.contains("lscpu output")); - } -} diff --git a/controller/src/runner/bench_standard.rs b/controller/src/runner/bench_datafusion.rs similarity index 51% rename from controller/src/runner/bench_standard.rs rename to controller/src/runner/bench_datafusion.rs index 9571e71..549c253 100644 --- a/controller/src/runner/bench_standard.rs +++ b/controller/src/runner/bench_datafusion.rs @@ -1,28 +1,35 @@ -//! Standard bench.sh benchmark runner — ports `run_bench_sh.sh`. - +//! Unified DataFusion benchmark runner. +//! +//! There is no benchmark allowlist and no per-name classification baked into +//! the controller. For each requested benchmark the runner resolves how to run +//! it at runtime: +//! +//! * If the name is a real Criterion `[[bench]]` target in the `benchmarks` +//! crate (discovered via `cargo metadata`), it is run with +//! `cargo bench --bench -- --save-baseline `. +//! * Otherwise it is run through `bench.sh run ` (TPC-H variants take a +//! direct `dfbench` shortcut). `bench.sh` itself runs some suites through the +//! Criterion SQL harness (`wide_schema`, …) and others through `dfbench`. +//! +//! Comparison is then driven by the artifacts each run actually produced, not +//! by the benchmark name: `results//*.json` is diffed with +//! `bench.sh compare_detail`, and `target/criterion` baselines are diffed with +//! `critcmp`. A run mixing both families emits both sections. + +use std::collections::HashSet; use std::path::{Path, PathBuf}; use anyhow::{Context, Result}; -use tracing::info; +use tracing::{info, warn}; use crate::github; -use crate::runner::bench_criterion::copy_criterion_baselines; use crate::runner::config::RunnerConfig; use crate::runner::git; use crate::runner::monitor::{self, ResourceStats}; use crate::runner::poster::CommentPoster; use crate::runner::shell; -/// Benchmarks orchestrated by `bench.sh` but executed through the Criterion SQL -/// harness (`cargo bench --bench sql`). Unlike the dfbench-based suites, these -/// write timings to `target/criterion/` rather than `results/*.json`, so they -/// are compared with `critcmp` (against per-side `--save-baseline`s) instead of -/// `bench.sh compare_detail`. -fn is_criterion_harness(bench: &str) -> bool { - matches!(bench, "wide_schema" | "predicate_eval") -} - -/// Run standard bench.sh benchmarks comparing a PR branch to its merge-base. +/// Run DataFusion benchmarks comparing a PR branch to its merge-base. pub async fn run(config: &RunnerConfig, poster: &CommentPoster) -> Result<()> { let repo_url = config.repo_url(); let benchmarks = &config.benchmarks; @@ -63,23 +70,34 @@ pub async fn run(config: &RunnerConfig, poster: &CommentPoster) -> Result<()> { // Pre-install stable toolchain to avoid rustup race in parallel builds git::rustup_stable().await?; - // Compile both in parallel - info!("=== Compiling PR branch and merge-base in parallel ==="); - let branch_benchmarks = branch_dir.join("benchmarks"); - let branch_build = shell::spawn_command( - "cargo", - &["build", "--release", "--bin", "dfbench"], - &branch_benchmarks, - "/tmp/branch_build.log", - ); - - let base_benchmarks = base_dir.join("benchmarks"); - let base_build = shell::spawn_command( - "cargo", - &["build", "--release", "--bin", "dfbench"], - &base_benchmarks, - "/tmp/base_build.log", - ); + // Resolve which requested names are real Criterion bench targets. Anything + // else runs through bench.sh. Resolved from the branch checkout so a new + // bench target added in the PR is recognized. + let requested: Vec = benchmarks.split_whitespace().map(String::from).collect(); + let bench_targets = criterion_bench_targets(&branch_dir.join("benchmarks")).await; + let is_criterion = |bench: &str| bench_targets.contains(bench); + let any_shell = requested.iter().any(|b| !is_criterion(b)); + + // dfbench is only needed for the bench.sh / TPC-H path. Build it for both + // sides in parallel; skip entirely for Criterion-only runs. + let builds = if any_shell { + info!("=== Compiling dfbench for PR branch and merge-base in parallel ==="); + let branch_build = shell::spawn_command( + "cargo", + &["build", "--release", "--bin", "dfbench"], + &branch_dir.join("benchmarks"), + "/tmp/branch_build.log", + ); + let base_build = shell::spawn_command( + "cargo", + &["build", "--release", "--bin", "dfbench"], + &base_dir.join("benchmarks"), + "/tmp/base_build.log", + ); + Some((branch_build, base_build)) + } else { + None + }; // Post "running" comment let uname = shell::uname().await; @@ -119,55 +137,103 @@ pub async fn run(config: &RunnerConfig, poster: &CommentPoster) -> Result<()> { .await?; // Wait for builds - info!("=== Waiting for builds ==="); - branch_build - .await - .context("branch build task panicked")? - .context("branch build failed")?; - base_build - .await - .context("base build task panicked")? - .context("base build failed")?; - info!("=== Builds complete ==="); - - // Set up bench runner from a third checkout - info!("=== Setting up bench runner ==="); - git::clone_shallow(&repo_url, &bench_dir, 200).await?; - git::checkout(&bench_dir, "origin/main").await?; + if let Some((branch_build, base_build)) = builds { + info!("=== Waiting for builds ==="); + branch_build + .await + .context("branch build task panicked")? + .context("branch build failed")?; + base_build + .await + .context("base build task panicked")? + .context("base build failed")?; + info!("=== Builds complete ==="); + } + // Set up bench runner from a third checkout (only used by the bench.sh path) let bench_benchmarks = bench_dir.join("benchmarks"); + if any_shell { + info!("=== Setting up bench runner ==="); + git::clone_shallow(&repo_url, &bench_dir, 200).await?; + git::checkout(&bench_dir, "origin/main").await?; + + // Clean any prior results + let results_dir = bench_benchmarks.join("results"); + if results_dir.exists() { + let _ = tokio::fs::remove_dir_all(&results_dir).await; + } - // Clean any prior results - let results_dir = bench_benchmarks.join("results"); - if results_dir.exists() { - let _ = tokio::fs::remove_dir_all(&results_dir).await; + // Copy TPC-H expected answer files so bench.sh skips the docker-based copy + copy_tpch_answers(&bench_benchmarks).await; } - // Copy TPC-H expected answer files so bench.sh skips the docker-based copy - copy_tpch_answers(&bench_benchmarks).await; - // Run each benchmark - let mut base_stats_list: Vec<(&str, ResourceStats)> = Vec::new(); - let mut branch_stats_list: Vec<(&str, ResourceStats)> = Vec::new(); + let mut base_stats_list: Vec<(String, ResourceStats)> = Vec::new(); + let mut branch_stats_list: Vec<(String, ResourceStats)> = Vec::new(); let bench_dir_str = bench_benchmarks.to_string_lossy().to_string(); let baseline_extra_env = config.baseline_env_args(); let changed_extra_env = config.changed_env_args(); - // Explicit RESULTS_NAME ensures bench.sh saves to a predictable directory, - // regardless of whether DATAFUSION_DIR is on a branch or detached HEAD. + // Explicit RESULTS_NAME / baseline name ensures bench.sh and Criterion save + // to a predictable location per side, regardless of branch vs detached HEAD. let base_results_name = "HEAD".to_string(); let bench_branch_name = git::sanitize_branch_name(&branch_name); - for bench in benchmarks.split_whitespace() { + // Track whether any Criterion baselines were produced per side, so the + // comparison can pick `critcmp HEAD ` vs a branch-only `critcmp`. + let mut criterion_base_ok = false; + let mut criterion_branch_ok = false; + + for bench in &requested { + if is_criterion(bench) { + info!("** Setting up data for criterion bench {bench} **"); + setup_criterion_data(bench, &branch_dir, &base_dir).await; + + info!("** Running {bench} baseline (criterion) **"); + match run_criterion_side( + bench, + &base_dir, + &base_results_name, + &config.bench_filter, + &baseline_extra_env, + ) + .await + { + Ok(stats) => { + criterion_base_ok = true; + base_stats_list.push((bench.clone(), stats)); + } + Err(e) => { + // Most likely a new bench target absent on the base — fall + // back to a branch-only comparison for it. + warn!("Criterion baseline for {bench} failed (new bench?): {e:#}"); + } + } + + info!("** Running {bench} branch (criterion) **"); + let stats = run_criterion_side( + bench, + &branch_dir, + &bench_branch_name, + &config.bench_filter, + &changed_extra_env, + ) + .await + .with_context(|| format!("run {bench} (branch, criterion)"))?; + criterion_branch_ok = true; + branch_stats_list.push((bench.clone(), stats)); + continue; + } + info!("** Creating data if needed for {bench} **"); cache_data(bench, &bench_dir_str).await; info!("** Running {bench} baseline **"); let base_spill_dir = PathBuf::from(format!("/workspace/spill-base-{bench}")); let _ = tokio::fs::create_dir_all(&base_spill_dir).await; - let base_stats = run_one_side( + let base_stats = run_shell_side( bench, &base_dir, &bench_benchmarks, @@ -178,12 +244,12 @@ pub async fn run(config: &RunnerConfig, poster: &CommentPoster) -> Result<()> { .await .with_context(|| format!("run {bench} (base)"))?; let _ = tokio::fs::remove_dir_all(&base_spill_dir).await; - base_stats_list.push((bench, base_stats)); + base_stats_list.push((bench.clone(), base_stats)); info!("** Running {bench} branch **"); let branch_spill_dir = PathBuf::from(format!("/workspace/spill-branch-{bench}")); let _ = tokio::fs::create_dir_all(&branch_spill_dir).await; - let branch_stats = run_one_side( + let branch_stats = run_shell_side( bench, &branch_dir, &bench_benchmarks, @@ -194,19 +260,16 @@ pub async fn run(config: &RunnerConfig, poster: &CommentPoster) -> Result<()> { .await .with_context(|| format!("run {bench} (branch)"))?; let _ = tokio::fs::remove_dir_all(&branch_spill_dir).await; - branch_stats_list.push((bench, branch_stats)); + branch_stats_list.push((bench.clone(), branch_stats)); } - // Compare and post results. dfbench-based suites are compared from their - // results/*.json via `bench.sh compare_detail`; Criterion-harness suites - // (which only produce target/criterion) are compared with `critcmp`. - let ran: Vec<&str> = benchmarks.split_whitespace().collect(); - let has_standard = ran.iter().any(|b| !is_criterion_harness(b)); - let has_criterion = ran.iter().any(|b| is_criterion_harness(b)); - + // Compare and post results, driven by the artifacts that were produced: + // - results/*.json (dfbench suites) → `bench.sh compare_detail` + // - target/criterion baselines (Criterion targets + the bench.sh SQL + // harness) → `critcmp` let mut report = String::new(); - if has_standard { + if any_shell && results_have_json(&bench_benchmarks, &base_results_name).await { let detail = shell::run_command( "./bench.sh", &["compare_detail", &base_results_name, &bench_branch_name], @@ -217,17 +280,30 @@ pub async fn run(config: &RunnerConfig, poster: &CommentPoster) -> Result<()> { report.push_str(&detail); } - if has_criterion { - // Gather both sides' baselines into one target/criterion tree, then - // diff the base ("HEAD") and branch baselines. - copy_criterion_baselines(&base_dir, &branch_dir).await; - let critcmp = shell::run_command( - "critcmp", - &[base_results_name.as_str(), bench_branch_name.as_str()], - &branch_dir, - ) - .await - .context("critcmp")?; + // Criterion output (from either path) lands under /target/criterion. + criterion_branch_ok = criterion_branch_ok || criterion_dir_present(&branch_dir).await; + criterion_base_ok = criterion_base_ok || criterion_dir_present(&base_dir).await; + if criterion_branch_ok { + let critcmp = if criterion_base_ok { + // Gather both sides' baselines into one tree, then diff them. + copy_criterion_baselines(&base_dir, &branch_dir).await; + shell::run_command( + "critcmp", + &[base_results_name.as_str(), bench_branch_name.as_str()], + &branch_dir, + ) + .await + .context("critcmp")? + } else { + // No baseline available (e.g. a brand-new bench) — branch-only. + let mut out = String::from("New benchmark — branch-only Criterion results:\n\n"); + out.push_str( + &shell::run_command("critcmp", &[bench_branch_name.as_str()], &branch_dir) + .await + .context("critcmp")?, + ); + out + }; if !report.is_empty() { report.push('\n'); } @@ -251,7 +327,93 @@ pub async fn run(config: &RunnerConfig, poster: &CommentPoster) -> Result<()> { Ok(()) } -/// Run a single benchmark on one side (base or branch). +/// Discover Criterion `[[bench]]` target names in the `benchmarks` crate via +/// `cargo metadata`. Returns an empty set on any failure — the caller then +/// treats every name as a `bench.sh` suite (which itself reports unknown +/// names), so a metadata hiccup degrades gracefully rather than misrouting. +async fn criterion_bench_targets(crate_dir: &Path) -> HashSet { + let out = match shell::run_command( + "cargo", + &["metadata", "--no-deps", "--format-version", "1"], + crate_dir, + ) + .await + { + Ok(out) => out, + Err(e) => { + warn!("cargo metadata failed; treating all benches as bench.sh suites: {e:#}"); + return HashSet::new(); + } + }; + parse_bench_targets(&out) +} + +/// Extract names of `bench`-kind targets from `cargo metadata` JSON. +fn parse_bench_targets(metadata_json: &str) -> HashSet { + let mut targets = HashSet::new(); + let Ok(value) = serde_json::from_str::(metadata_json) else { + return targets; + }; + let Some(packages) = value.get("packages").and_then(|p| p.as_array()) else { + return targets; + }; + for pkg in packages { + let Some(pkg_targets) = pkg.get("targets").and_then(|t| t.as_array()) else { + continue; + }; + for target in pkg_targets { + let is_bench = target + .get("kind") + .and_then(|k| k.as_array()) + .map(|kinds| kinds.iter().any(|k| k.as_str() == Some("bench"))) + .unwrap_or(false); + if is_bench { + if let Some(name) = target.get("name").and_then(|n| n.as_str()) { + targets.insert(name.to_string()); + } + } + } + } + targets +} + +/// Run a Criterion bench target on one side, saving a named baseline. +/// +/// A non-empty `bench_filter` (from the `BENCH_FILTER` env var) is passed +/// through to Criterion as a test-name filter. +async fn run_criterion_side( + bench: &str, + side_dir: &Path, + baseline_name: &str, + bench_filter: &str, + extra_env: &[String], +) -> Result { + let mut bench_args: Vec = vec![ + "bench".into(), + "--features=parquet".into(), + "--bench".into(), + bench.into(), + "--".into(), + "--save-baseline".into(), + baseline_name.into(), + ]; + if !bench_filter.is_empty() { + bench_args.push(bench_filter.to_string()); + } + let (_, stats) = if extra_env.is_empty() { + let args_ref: Vec<&str> = bench_args.iter().map(|s| s.as_str()).collect(); + shell::run_command_monitored("cargo", &args_ref, side_dir, None).await? + } else { + let mut env_args: Vec = extra_env.to_vec(); + env_args.push("cargo".to_string()); + env_args.extend(bench_args); + let env_args_ref: Vec<&str> = env_args.iter().map(|s| s.as_str()).collect(); + shell::run_command_monitored("env", &env_args_ref, side_dir, None).await? + }; + Ok(stats) +} + +/// Run a single `bench.sh` benchmark on one side (base or branch). /// /// For TPC-H variants we bypass `bench.sh run` and invoke the prebuilt `dfbench` /// binary directly. Upstream PR apache/datafusion#21707 ported `bench.sh`'s @@ -262,7 +424,7 @@ pub async fn run(config: &RunnerConfig, poster: &CommentPoster) -> Result<()> { /// `bench_dir/benchmarks/results/`). The `dfbench tpch` subcommand still /// exists upstream, so we call it directly with the same args the old /// `run_tpch` used. Other benchmarks continue through `bench.sh`. -async fn run_one_side( +async fn run_shell_side( bench: &str, side_dir: &Path, bench_benchmarks: &Path, @@ -286,14 +448,12 @@ async fn run_one_side( format!("DATAFUSION_DIR={}", side_dir.display()), format!("RESULTS_NAME={results_name}"), format!("DATAFUSION_RUNTIME_TEMP_DIRECTORY={}", spill_dir.display()), + // Suites that bench.sh runs through the Criterion SQL harness read + // SQL_CARGO_COMMAND; setting it unconditionally is harmless for the + // dfbench-based suites (which never read it) and saves a named + // baseline per side for the ones that do, so we can critcmp them. + format!("SQL_CARGO_COMMAND=cargo bench --bench sql -- --save-baseline {results_name}"), ]; - // Criterion-harness suites don't emit results/*.json; have them save a - // named Criterion baseline per side so we can compare with critcmp. - if is_criterion_harness(bench) { - args.push(format!( - "SQL_CARGO_COMMAND=cargo bench --bench sql -- --save-baseline {results_name}" - )); - } args.extend(extra_env.iter().cloned()); args.extend([ "./bench.sh".to_string(), @@ -400,6 +560,78 @@ async fn run_tpch_direct( Ok(stats) } +/// Datasets a Criterion bench target needs generated before it can run. +fn required_datasets(bench_name: &str) -> &'static [&'static str] { + match bench_name { + "sql_planner" => &["clickbench_partitioned"], + _ => &[], + } +} + +/// Generate the data a Criterion bench needs, once in the branch checkout, and +/// symlink it into the base checkout so both sides share the same inputs. +async fn setup_criterion_data(bench: &str, branch_dir: &Path, base_dir: &Path) { + let datasets = required_datasets(bench); + if datasets.is_empty() { + return; + } + + let branch_benchmarks = branch_dir.join("benchmarks"); + let bench_dir_str = branch_benchmarks.to_string_lossy().to_string(); + for dataset in datasets { + info!("Setting up data for {dataset}"); + cache_data(dataset, &bench_dir_str).await; + } + + let branch_data = branch_benchmarks.join("data"); + let base_data = base_dir.join("benchmarks/data"); + if branch_data.exists() && !base_data.exists() { + info!("Symlinking benchmark data into base directory"); + let _ = tokio::fs::symlink(&branch_data, &base_data).await; + } +} + +/// Copy Criterion baselines from base into the branch target tree so a single +/// `critcmp` invocation can see both sides. +async fn copy_criterion_baselines(base_dir: &Path, branch_dir: &Path) { + let src = base_dir.join("target/criterion"); + let dst = branch_dir.join("target/criterion"); + if src.exists() { + let _ = shell::run_command( + "cp", + &[ + "-r", + &format!("{}/.", src.to_string_lossy()), + &dst.to_string_lossy(), + ], + Path::new("/"), + ) + .await; + } +} + +/// Whether `/target/criterion` exists (i.e. a Criterion run wrote there). +async fn criterion_dir_present(side_dir: &Path) -> bool { + tokio::fs::metadata(side_dir.join("target/criterion")) + .await + .map(|m| m.is_dir()) + .unwrap_or(false) +} + +/// Whether the base side produced any `results//*.json` (dfbench suites). +async fn results_have_json(bench_benchmarks: &Path, results_name: &str) -> bool { + let dir = bench_benchmarks.join("results").join(results_name); + let Ok(mut entries) = tokio::fs::read_dir(&dir).await else { + return false; + }; + while let Ok(Some(entry)) = entries.next_entry().await { + if entry.path().extension().and_then(|e| e.to_str()) == Some("json") { + return true; + } + } + false +} + /// Copy TPC-H answer files from the baked-in location into the benchmark data dirs. async fn copy_tpch_answers(bench_dir: &Path) { let answers_src = Path::new("/data/tpch-answers"); @@ -440,8 +672,8 @@ async fn cache_data(bench: &str, bench_dir: &str) { /// Build the resource usage section from collected stats. fn format_resource_section( - base_stats: &[(&str, ResourceStats)], - branch_stats: &[(&str, ResourceStats)], + base_stats: &[(String, ResourceStats)], + branch_stats: &[(String, ResourceStats)], ) -> String { let mut section = String::new(); for (bench, stats) in base_stats { @@ -520,15 +752,41 @@ mod tests { } #[test] - fn criterion_harness_classification() { - // Criterion SQL-harness suites → compared via critcmp. - assert!(is_criterion_harness("wide_schema")); - assert!(is_criterion_harness("predicate_eval")); - // dfbench/json suites → compared via bench.sh compare_detail. - assert!(!is_criterion_harness("tpch")); - assert!(!is_criterion_harness("clickbench_1")); - assert!(!is_criterion_harness("imdb")); - assert!(!is_criterion_harness("h2o")); + fn parse_bench_targets_picks_bench_kind() { + let json = r#"{ + "packages": [ + {"targets": [ + {"name": "dfbench", "kind": ["bin"]}, + {"name": "sql_planner", "kind": ["bench"]}, + {"name": "sql", "kind": ["bench"]}, + {"name": "datafusion-benchmarks", "kind": ["lib"]} + ]} + ] + }"#; + let targets = parse_bench_targets(json); + assert!(targets.contains("sql_planner")); + assert!(targets.contains("sql")); + assert!(!targets.contains("dfbench")); + assert!(!targets.contains("datafusion-benchmarks")); + } + + #[test] + fn parse_bench_targets_handles_garbage() { + assert!(parse_bench_targets("not json").is_empty()); + assert!(parse_bench_targets("{}").is_empty()); + } + + #[test] + fn required_datasets_sql_planner() { + assert_eq!( + required_datasets("sql_planner"), + &["clickbench_partitioned"] + ); + } + + #[test] + fn required_datasets_unknown_is_empty() { + assert!(required_datasets("wide_schema").is_empty()); } #[test] diff --git a/controller/src/runner/config.rs b/controller/src/runner/config.rs index e6e2e56..24dbf54 100644 --- a/controller/src/runner/config.rs +++ b/controller/src/runner/config.rs @@ -18,8 +18,7 @@ use crate::runner::poster::CommentPoster; /// Benchmark runner variant, parsed from `BENCH_TYPE`. #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum BenchType { - Standard, - Criterion, + Datafusion, ArrowCriterion, MainTracking, } @@ -27,8 +26,10 @@ pub enum BenchType { impl BenchType { fn from_str(s: &str) -> Result { match s { - "standard" => Ok(Self::Standard), - "criterion" => Ok(Self::Criterion), + // `standard`/`criterion` are legacy job-type strings that may still + // be present on in-flight DB rows; both map to the unified + // datafusion runner. + "datafusion" | "standard" | "criterion" => Ok(Self::Datafusion), "arrow_criterion" => Ok(Self::ArrowCriterion), "main_tracking" => Ok(Self::MainTracking), other => anyhow::bail!("unknown BENCH_TYPE: {other}"), @@ -212,18 +213,23 @@ mod tests { use super::*; #[test] - fn bench_type_standard() { + fn bench_type_datafusion() { assert_eq!( - BenchType::from_str("standard").unwrap(), - BenchType::Standard + BenchType::from_str("datafusion").unwrap(), + BenchType::Datafusion ); } #[test] - fn bench_type_criterion() { + fn bench_type_legacy_strings_map_to_datafusion() { + // In-flight DB rows may still carry the old job-type strings. + assert_eq!( + BenchType::from_str("standard").unwrap(), + BenchType::Datafusion + ); assert_eq!( BenchType::from_str("criterion").unwrap(), - BenchType::Criterion + BenchType::Datafusion ); } diff --git a/controller/src/runner/mod.rs b/controller/src/runner/mod.rs index bcbee28..7ce2b86 100644 --- a/controller/src/runner/mod.rs +++ b/controller/src/runner/mod.rs @@ -1,6 +1,5 @@ pub mod bench_arrow; -pub mod bench_criterion; -pub mod bench_standard; +pub mod bench_datafusion; pub mod config; pub mod controller_client; pub mod git; diff --git a/services/controller.ts b/services/controller.ts index 5dab562..36d3772 100644 --- a/services/controller.ts +++ b/services/controller.ts @@ -191,44 +191,27 @@ export const controllerStatefulSet = new k8s.apps.v1.StatefulSet("benchmark-cont "mzabaluev", "sdf-jkl", "liamzwbao", "cetra3", "brunal", "Fokko", "kunalsinghdadhwal", "grtlr", "codephage2020", "asubiotto", ], + // No benchmark allowlist: any requested name is scheduled and + // resolved on the runner. `kind` selects how the repo runs + // (datafusion = bench.sh + Criterion, resolved per-benchmark; + // arrow = arrow-rs Criterion). `default_standard` is the suite + // used by a bare `run benchmarks`. repos: { "adriangb/datafusion": { - standard: [ - "tpch", "tpch10", "tpch_mem", "tpch_mem10", - "topk_tpch", - "clickbench_partitioned", "clickbench_extended", - "clickbench_1", "clickbench_pushdown", - "external_aggr", "tpcds", "smj", "sort_pushdown", - "sort_pushdown_sorted", "sort_pushdown_inexact", - "sort_pushdown_inexact_unsorted", "sort_pushdown_inexact_overlap", - "wide_schema", "predicate_eval", - ], + kind: "datafusion", default_standard: [ "clickbench_partitioned", "tpcds", "tpch", ], - criterion: ["*"], }, "apache/datafusion": { - standard: [ - "tpch", "tpch10", "tpch_mem", "tpch_mem10", - "topk_tpch", - "clickbench_partitioned", "clickbench_extended", - "clickbench_1", "clickbench_pushdown", - "external_aggr", "tpcds", "smj", "sort_pushdown", - "sort_pushdown_sorted", "sort_pushdown_inexact", - "sort_pushdown_inexact_unsorted", "sort_pushdown_inexact_overlap", - "wide_schema", "predicate_eval", - ], + kind: "datafusion", default_standard: [ "clickbench_partitioned", "tpcds", "tpch", ], - criterion: ["*"], }, "apache/arrow-rs": { - standard: [], + kind: "arrow", default_standard: [], - criterion: ["*"], - criterion_type: "arrow", }, }, }) },