From c3f2661c87127d755c92bac1a60cf2a3bfe2a7c3 Mon Sep 17 00:00:00 2001 From: James Dumay Date: Tue, 11 Aug 2026 16:53:20 +1000 Subject: [PATCH 1/2] feat(bench): run SWE-Gym through Harbor --- crates/skippy-bench/README.md | 30 ++- crates/skippy-bench/src/cli.rs | 23 ++ crates/skippy-bench/src/evals.rs | 218 ++++++++++++++++-- .../skippy-bench/src/evals/adapters/harbor.rs | 144 ++++++++++++ crates/skippy-bench/src/evals/adapters/mod.rs | 8 +- crates/skippy-bench/src/evals/registry.rs | 31 ++- crates/skippy-bench/src/evals/run.rs | 87 ++++++- crates/skippy-bench/src/evals/sync.rs | 38 +-- .../local_generation/linear_decode.rs | 82 +++++++ 9 files changed, 602 insertions(+), 59 deletions(-) create mode 100644 crates/skippy-bench/src/evals/adapters/harbor.rs diff --git a/crates/skippy-bench/README.md b/crates/skippy-bench/README.md index c223eed0f..3817b999d 100644 --- a/crates/skippy-bench/README.md +++ b/crates/skippy-bench/README.md @@ -112,7 +112,8 @@ The current core pack is: | Eval id | External harness | Default run | |---|---|---| | `speed-bench` | llama.cpp `tools/server/bench/speed-bench` | Native SPEED-Bench qualitative run across all categories, no sample limit, `--osl 1024` | -| `terminal-bench` | Terminal-Bench CLI (`tb`) | `terminal-bench-core==0.1.1`, Terminus agent, no task-id filter | +| `terminal-bench` | Pinned Harbor (`ff69e554`) | Harbor `terminal-bench@2.0` dataset with Terminus2 | +| `swe-gym` | Pinned Harbor (`ff69e554`) | SWE-Gym Lite via Harbor `swegym-lite`; use `--task-id` for one task | | `swe-bench-pro` | Scale SWE-Bench Pro OS repo | Upstream SWE-agent patch generation, patch gathering, and `swe_bench_pro_eval.py` | | `mcp-atlas` | Scale MCP-Atlas repo | Native MCP-Atlas completion script with upstream `--no-filter`, plus scoring through auto-started MCP services | @@ -128,6 +129,26 @@ skippy-bench eval run terminal-bench \ --metrics-run-id run-local-qwen ``` +SWE-Gym Lite one-task smoke: + +```bash +skippy-bench eval sync swe-gym +skippy-bench eval run swe-gym \ + --task-id getmoto__moto-5752 \ + --dataset lite \ + --base-url http://127.0.0.1:9337/v1 \ + --session-id swegym-smoke \ + --model org/repo:Q4_K_M \ + --metrics-http http://127.0.0.1:18080 \ + --metrics-run-id swegym-smoke +``` + +Start `metrics-server` before the run for request correlation. Full SWE-Gym +runs use Harbor's official `-d swegym-lite` dataset path. Single-task +preparation uses `uv run --with swebench adapters/swegym/run_adapter.py`. +If the task container cannot reach Mesh, provide `--harbor-endpoint-url` with +a container-reachable URL. + `--timeout-secs` is forwarded to native harnesses as their request/task timeout where supported. It is not a SkippyBench dataset limit and does not cap full canonical runs. Use `--harness-timeout-secs` only when an operator wants a hard @@ -150,9 +171,10 @@ cache behind upstream. Each run records that resolved commit in Before launching native harness traffic, `eval run` enforces the same required tool checks as `eval doctor`, including Docker container-start readiness for Docker-backed evals. -Terminal-Bench is installed through `uv tool install --python 3.12` because the -current `tb` CLI is not compatible with Python 3.14. `eval doctor` checks that -Docker's daemon is reachable and can start a tiny container, not just that the +Harbor is synced once at pinned commit `ff69e554` and reused by both +Terminal-Bench and SWE-Gym; runs do not clone it per invocation. The legacy +direct `tb` runner is unsupported. `eval doctor` checks that Docker's daemon is +reachable and can start a tiny container, not just that the `docker` CLI exists or that `docker info` returns. MCP-Atlas starts its Docker agent environment and Python completion service when ports `1984` and `3000` are not already reachable, waits for readiness, diff --git a/crates/skippy-bench/src/cli.rs b/crates/skippy-bench/src/cli.rs index a1ec2c126..8a8597c4a 100644 --- a/crates/skippy-bench/src/cli.rs +++ b/crates/skippy-bench/src/cli.rs @@ -66,6 +66,7 @@ pub enum EvalCommandKind { pub enum EvalId { SpeedBench, TerminalBench, + SweGym, SweBenchPro, McpAtlas, } @@ -75,6 +76,7 @@ impl EvalId { match self { Self::SpeedBench => "speed-bench", Self::TerminalBench => "terminal-bench", + Self::SweGym => "swe-gym", Self::SweBenchPro => "swe-bench-pro", Self::McpAtlas => "mcp-atlas", } @@ -145,6 +147,27 @@ pub struct EvalRunArgs { pub model: String, #[arg(long, default_value = "skippy-bench")] pub api_key: String, + #[arg(long, help = "Run one Harbor task instead of the selected dataset")] + pub task_id: Option, + #[arg(long, default_value = "lite", help = "Harbor dataset or split")] + pub dataset: String, + #[arg(long, default_value = "terminus-2", help = "Harbor agent name")] + pub agent: String, + #[arg( + long, + help = "Endpoint URL reachable from the Harbor task container; required for Harbor runs when --base-url points at localhost" + )] + pub harbor_endpoint_url: Option, + #[arg( + long, + help = "Stable session ID sent to Mesh as X-Session-ID and session_id" + )] + pub session_id: Option, + #[arg( + long, + help = "Cacheline smoke state directory containing source/causal/ingestion assertion files" + )] + pub cacheline_state: Option, #[arg(long)] pub cache_root: Option, #[arg(long)] diff --git a/crates/skippy-bench/src/evals.rs b/crates/skippy-bench/src/evals.rs index 999aef31a..645df128a 100644 --- a/crates/skippy-bench/src/evals.rs +++ b/crates/skippy-bench/src/evals.rs @@ -20,9 +20,10 @@ use crate::cli::{ }; use crate::telemetry_report::{self, BenchTelemetry}; -const CORE_EVALS: [EvalId; 4] = [ +const CORE_EVALS: [EvalId; 5] = [ EvalId::SpeedBench, EvalId::TerminalBench, + EvalId::SweGym, EvalId::SweBenchPro, EvalId::McpAtlas, ]; @@ -134,7 +135,7 @@ struct RunArtifact { path: String, } -#[derive(Clone)] +#[derive(Clone, Debug)] struct CommandSpec { program: String, args: Vec, @@ -466,18 +467,16 @@ fn shell_quote(raw: &str) -> String { mod tests { use super::*; use super::{ - adapters::{ - mcp_atlas_command, speed_bench_command, swe_bench_pro_command, terminal_bench_command, - }, + adapters::{harbor_command, mcp_atlas_command, speed_bench_command, swe_bench_pro_command}, doctor::preflight_eval_run, registry::{definition, selected_evals}, run::{ - fill_client_rates, resolved_harness_commit, run_artifacts, speed_bench_metrics, - speed_bench_output_path, speed_bench_response_timings_path, swe_bench_pro_metrics, - swe_bench_pro_output_path, telemetry_or_unavailable, terminal_bench_metrics, - terminal_bench_output_path, + fill_client_rates, harbor_jobs_output_path, harbor_metrics, resolved_harness_commit, + run_artifacts, speed_bench_metrics, speed_bench_output_path, + speed_bench_response_timings_path, swe_bench_pro_metrics, swe_bench_pro_output_path, + telemetry_or_unavailable, terminal_bench_metrics, terminal_bench_output_path, }, - sync::existing_repo_sync_steps, + sync::{existing_repo_sync_steps, new_repo_sync_steps}, }; #[test] @@ -505,6 +504,12 @@ mod tests { base_url: "http://127.0.0.1:9337/v1".to_string(), model: "tiny-local".to_string(), api_key: "test".to_string(), + task_id: None, + dataset: "lite".to_string(), + agent: "terminus-2".to_string(), + harbor_endpoint_url: None, + session_id: None, + cacheline_state: None, cache_root: None, output_dir: None, timeout_secs: 30, @@ -595,6 +600,12 @@ mod tests { base_url: "http://127.0.0.1:9337/v1".to_string(), model: "tiny-local".to_string(), api_key: "terminal-secret-value".to_string(), + task_id: None, + dataset: "lite".to_string(), + agent: "terminus-2".to_string(), + harbor_endpoint_url: Some("http://host.docker.internal:9337/v1".to_string()), + session_id: None, + cacheline_state: None, cache_root: None, output_dir: None, timeout_secs: 30, @@ -606,24 +617,107 @@ mod tests { metrics_finalize_only: false, dry_run: true, }; - let command = terminal_bench_command(&args, Path::new("/tmp/skippy-run")); - assert!( - command - .args - .contains(&"terminal-bench-core==0.1.1".to_string()) - ); + let run_dir = temp_run_dir("terminal-harbor-command"); + fs::create_dir_all(run_dir.join("raw")).unwrap(); + let command = harbor_command( + definition(EvalId::TerminalBench), + &args, + Path::new("/tmp/skippy-cache"), + &run_dir, + ) + .unwrap(); assert!(!command.args.contains(&"--task-id".to_string())); - assert!(command.args.contains(&"--output-path".to_string())); - assert!( - command - .envs - .contains(&("OPENAI_BASE_URL".to_string(), args.base_url)) - ); + let script = fs::read_to_string(run_dir.join("raw/run-harbor.sh")).unwrap(); + assert!(script.contains("terminal-bench@2.0")); + assert!(script.contains("harbor jobs start")); + assert!(script.contains("session_id")); + assert!(command.envs.contains(&( + "SKIPPY_BENCH_TASK_ENDPOINT_URL".to_string(), + "http://host.docker.internal:9337/v1".to_string() + ))); assert!(command.secret_envs.contains(&( "OPENAI_API_KEY".to_string(), "terminal-secret-value".to_string() ))); + assert!(script.contains("--n-concurrent 1")); + assert!(!script.contains("--n-concurrent-trials")); + assert!(!script.contains("--ae OPENAI_API_KEY")); assert!(!command.display().contains("terminal-secret-value")); + let _ = fs::remove_dir_all(run_dir); + } + + #[test] + fn swe_gym_command_selects_one_instance_and_redacts_endpoint_credentials() { + let run_dir = temp_run_dir("swegym-harbor-command"); + fs::create_dir_all(run_dir.join("raw")).unwrap(); + let mut args = eval_run_args(EvalId::SweGym, "swegym-secret"); + args.task_id = Some("getmoto__moto-5752".to_string()); + args.run_id = Some("smoke/run".to_string()); + args.harbor_endpoint_url = Some("http://host.docker.internal:9337/v1".to_string()); + + let command = harbor_command( + definition(EvalId::SweGym), + &args, + Path::new("/tmp/skippy-cache"), + &run_dir, + ) + .unwrap(); + let script = fs::read_to_string(run_dir.join("raw/run-harbor.sh")).unwrap(); + assert!(script.contains("--dataset lite")); + assert!(script.contains("uv run --with swebench adapters/swegym/run_adapter.py")); + assert!(script.contains("--instance-id getmoto__moto-5752")); + assert!(script.contains("harbor jobs start \"$@\"")); + assert!(script.contains("--ak session_id=\"$SKIPPY_BENCH_SESSION_ID\"")); + assert!(script.contains("--job-name skippy-smoke-run")); + assert!(!script.contains("--ae OPENAI_API_KEY")); + assert!( + command + .secret_envs + .contains(&("OPENAI_API_KEY".to_string(), "swegym-secret".to_string())) + ); + assert!(!command.display().contains("swegym-secret")); + let _ = fs::remove_dir_all(run_dir); + } + + #[test] + fn swe_gym_without_task_id_uses_official_harbor_dataset() { + let run_dir = temp_run_dir("harbor-local-endpoint"); + fs::create_dir_all(run_dir.join("raw")).unwrap(); + let args = eval_run_args(EvalId::SweGym, "secret"); + let command = harbor_command( + definition(EvalId::SweGym), + &args, + Path::new("/tmp/skippy-cache"), + &run_dir, + ) + .unwrap(); + let script = fs::read_to_string(run_dir.join("raw/run-harbor.sh")).unwrap(); + assert!(script.contains("set -- -d swegym-lite")); + assert!(!script.contains("run_adapter.py")); + assert!(command.envs.iter().any(|(key, value)| { + key == "SKIPPY_BENCH_BASE_URL" && value == "http://127.0.0.1:9337/v1" + })); + let _ = fs::remove_dir_all(run_dir); + } + + #[test] + fn swe_gym_rejects_unsafe_task_ids() { + for task_id in ["/absolute", "nested/task", "nested\\task", "..", "foo..bar"] { + let run_dir = temp_run_dir("unsafe-task-id"); + fs::create_dir_all(run_dir.join("raw")).unwrap(); + let mut args = eval_run_args(EvalId::SweGym, "secret"); + args.task_id = Some(task_id.to_string()); + let error = harbor_command( + definition(EvalId::SweGym), + &args, + Path::new("/tmp/skippy-cache"), + &run_dir, + ) + .unwrap_err() + .to_string(); + assert!(error.contains("safe task name"), "{task_id}: {error}"); + let _ = fs::remove_dir_all(run_dir); + } } #[test] @@ -703,6 +797,41 @@ mod tests { ); } + #[test] + fn new_repo_sync_fetches_and_checks_out_pinned_sha_after_branchless_clone() { + let steps = new_repo_sync_steps( + Path::new("/tmp/harbor"), + "https://github.com/harbor-framework/harbor.git", + "ff69e554fac1c751aa608e03de027db9043a2eac", + ); + + assert_eq!( + steps[0].args, + [ + "clone", + "--recurse-submodules", + "https://github.com/harbor-framework/harbor.git", + "/tmp/harbor" + ] + ); + assert_eq!( + steps[1].args, + [ + "-C", + "/tmp/harbor", + "fetch", + "--prune", + "origin", + "ff69e554fac1c751aa608e03de027db9043a2eac" + ] + ); + assert_eq!( + steps[2].args, + ["-C", "/tmp/harbor", "checkout", "--detach", "FETCH_HEAD"] + ); + assert!(!steps[0].args.contains(&"--branch".to_string())); + } + #[test] fn resolved_harness_commit_reads_checked_out_revision() { let root = temp_run_dir("harness-revision"); @@ -940,6 +1069,45 @@ mod tests { let _ = fs::remove_dir_all(run_dir); } + #[test] + fn harbor_metrics_normalize_trial_rewards() { + let run_dir = temp_run_dir("harbor-metrics"); + let job_dir = harbor_jobs_output_path(&run_dir).join("job"); + fs::create_dir_all(&job_dir).unwrap(); + fs::write( + job_dir.join("result.json"), + r#"{"n_total_trials": 5, "stats": {"n_completed_trials": 5}}"#, + ) + .unwrap(); + fs::create_dir_all(job_dir.join("trial-1")).unwrap(); + fs::write( + job_dir.join("trial-1/result.json"), + r#"{"verifier_result": {"rewards": {"reward": 1}}}"#, + ) + .unwrap(); + fs::create_dir_all(job_dir.join("trial-2")).unwrap(); + fs::write( + job_dir.join("trial-2/result.json"), + r#"{"verifier_result": {"rewards": {"reward": 0}}}"#, + ) + .unwrap(); + fs::create_dir_all(job_dir.join("trial-3")).unwrap(); + fs::write( + job_dir.join("trial-3/result.json"), + r#"{"verifier_result": {}}"#, + ) + .unwrap(); + fs::create_dir_all(job_dir.join("trial-4")).unwrap(); + fs::write(job_dir.join("trial-4/result.json"), "not-json").unwrap(); + fs::create_dir_all(job_dir.join("trial-5")).unwrap(); + + let metrics = harbor_metrics(&run_dir).unwrap(); + assert_eq!(metrics.request_count, Some(5)); + assert_eq!(metrics.failed_count, Some(4)); + assert_eq!(metrics.pass_rate, Some(0.2)); + let _ = fs::remove_dir_all(run_dir); + } + #[test] fn fill_client_rates_computes_missing_rates() { let mut metrics = EvalMetrics { @@ -968,6 +1136,12 @@ mod tests { base_url: "http://127.0.0.1:9337/v1".to_string(), model: "tiny-local".to_string(), api_key: api_key.to_string(), + task_id: None, + dataset: "lite".to_string(), + agent: "terminus-2".to_string(), + harbor_endpoint_url: None, + session_id: None, + cacheline_state: None, cache_root: None, output_dir: None, timeout_secs: 30, diff --git a/crates/skippy-bench/src/evals/adapters/harbor.rs b/crates/skippy-bench/src/evals/adapters/harbor.rs new file mode 100644 index 000000000..745653169 --- /dev/null +++ b/crates/skippy-bench/src/evals/adapters/harbor.rs @@ -0,0 +1,144 @@ +use super::super::{run::harbor_jobs_output_path, *}; + +pub(in crate::evals) fn harbor_command( + definition: EvalDefinition, + args: &EvalRunArgs, + root: &Path, + run_dir: &Path, +) -> Result { + let harness = harness_dir(root, definition); + let session_id = args + .session_id + .clone() + .or_else(|| args.run_id.clone().map(|run_id| format!("skippy-{run_id}"))) + .unwrap_or_else(|| "skippy-bench-session".to_string()); + let jobs_dir = harbor_jobs_output_path(run_dir); + let task_root = run_dir.join("raw/harbor-tasks"); + let script_path = run_dir.join("raw/run-harbor.sh"); + let job_name = format!( + "skippy-{}", + harbor_job_slug(args.run_id.as_deref().unwrap_or("eval")) + ); + let model = litellm_model_name(&args.model); + let dataset = match definition.id { + EvalId::TerminalBench => { + if args.dataset != "lite" { + args.dataset.clone() + } else { + "terminal-bench@2.0".to_string() + } + } + EvalId::SweGym => match args.dataset.as_str() { + "lite" | "full" => args.dataset.clone(), + other => bail!("SWE-Gym dataset must be `lite` or `full`, got {other:?}"), + }, + _ => bail!("{} is not a Harbor-backed eval", definition.id.as_str()), + }; + + let task_selection = match definition.id { + EvalId::SweGym => { + if let Some(task_id) = args.task_id.as_deref() { + validate_task_id(task_id)?; + } + if args.task_id.is_none() { + let harbor_dataset = if dataset == "lite" { + "swegym-lite" + } else { + "swegym" + }; + format!("set -- -d {}\n", shell_quote(harbor_dataset)) + } else { + let instance = args + .task_id + .as_deref() + .map(|task_id| format!(" --instance-id {}", shell_quote(task_id))) + .unwrap_or_default(); + let task_path = args + .task_id + .as_deref() + .map(|task_id| { + format!( + "{}", + shell_quote(&task_root.join(task_id).display().to_string()) + ) + }) + .unwrap_or_else(|| shell_quote(&task_root.display().to_string())); + format!( + "set -- -p {}\nuv run --with swebench adapters/swegym/run_adapter.py --dataset {}{} --task-dir {}\n", + task_path, + shell_quote(&dataset), + instance, + shell_quote(&task_root.display().to_string()) + ) + } + } + EvalId::TerminalBench => { + if args.task_id.is_some() { + bail!("--task-id is currently supported only for swe-gym"); + } + format!("set -- -d {}\n", shell_quote(&dataset)) + } + _ => unreachable!(), + }; + + let task_container_endpoint = args + .harbor_endpoint_url + .as_ref() + .map(|_| " --ae OPENAI_BASE_URL=\"$SKIPPY_BENCH_TASK_ENDPOINT_URL\"") + .unwrap_or(""); + let script = format!( + "#!/bin/sh\nset -eu\n{task_selection}exec uv run harbor jobs start \"$@\" -a {agent} -m {model} -o {jobs} --job-name {job} --n-concurrent {concurrency} --ak api_base=\"$SKIPPY_BENCH_BASE_URL\" --ak session_id=\"$SKIPPY_BENCH_SESSION_ID\"{task_container_endpoint}\n", + task_selection = task_selection, + agent = shell_quote(&args.agent), + model = shell_quote(&model), + jobs = shell_quote(&jobs_dir.display().to_string()), + job = shell_quote(&job_name), + concurrency = args.endpoint_concurrency, + task_container_endpoint = task_container_endpoint, + ); + fs::write(&script_path, script) + .with_context(|| format!("write Harbor launcher {}", script_path.display()))?; + + let mut command = CommandSpec::new("sh") + .args([script_path.display().to_string()]) + .cwd(harness) + .env("SKIPPY_BENCH_BASE_URL", args.base_url.clone()) + .env("SKIPPY_BENCH_SESSION_ID", session_id.clone()) + .secret_env("OPENAI_API_KEY", args.api_key.clone()); + if let Some(endpoint) = &args.harbor_endpoint_url { + command = command.env("SKIPPY_BENCH_TASK_ENDPOINT_URL", endpoint.clone()); + } + Ok(command) +} + +fn validate_task_id(task_id: &str) -> Result<()> { + if task_id.is_empty() + || Path::new(task_id).is_absolute() + || task_id.contains('/') + || task_id.contains('\\') + || task_id.contains("..") + { + bail!( + "SWE-Gym --task-id must be a single safe task name without separators, absolute paths, or '..': {task_id:?}" + ); + } + Ok(()) +} + +fn harbor_job_slug(run_id: &str) -> String { + let mut slug = String::new(); + for character in run_id.chars() { + if character.is_ascii_alphanumeric() || matches!(character, '-' | '_' | '.') { + slug.push(character.to_ascii_lowercase()); + } else if !slug.ends_with('-') { + slug.push('-'); + } + } + let slug = slug.trim_matches('-').to_string(); + let slug = if slug.is_empty() { + "eval".to_string() + } else { + slug + }; + slug.chars().take(48).collect() +} diff --git a/crates/skippy-bench/src/evals/adapters/mod.rs b/crates/skippy-bench/src/evals/adapters/mod.rs index 56f6dc40e..694b239f8 100644 --- a/crates/skippy-bench/src/evals/adapters/mod.rs +++ b/crates/skippy-bench/src/evals/adapters/mod.rs @@ -1,13 +1,13 @@ use super::*; +mod harbor; mod mcp_atlas; mod speed_bench; mod swe_bench_pro; -mod terminal_bench; pub(super) use self::{ - mcp_atlas::mcp_atlas_command, speed_bench::speed_bench_command, - swe_bench_pro::swe_bench_pro_command, terminal_bench::terminal_bench_command, + harbor::harbor_command, mcp_atlas::mcp_atlas_command, speed_bench::speed_bench_command, + swe_bench_pro::swe_bench_pro_command, }; pub(super) fn run_command( @@ -18,7 +18,7 @@ pub(super) fn run_command( ) -> Result { Ok(match definition.id { EvalId::SpeedBench => speed_bench_command(definition, args, root, run_dir)?, - EvalId::TerminalBench => terminal_bench_command(args, run_dir), + EvalId::TerminalBench | EvalId::SweGym => harbor_command(definition, args, root, run_dir)?, EvalId::SweBenchPro => swe_bench_pro_command(args, root, run_dir)?, EvalId::McpAtlas => mcp_atlas_command(args, root, run_dir)?, }) diff --git a/crates/skippy-bench/src/evals/registry.rs b/crates/skippy-bench/src/evals/registry.rs index aae1fb199..f4ca8fc92 100644 --- a/crates/skippy-bench/src/evals/registry.rs +++ b/crates/skippy-bench/src/evals/registry.rs @@ -1,5 +1,8 @@ use super::*; +pub(super) const HARBOR_REPO: &str = "https://github.com/harbor-framework/harbor.git"; +pub(super) const HARBOR_REF: &str = "ff69e554fac1c751aa608e03de027db9043a2eac"; + pub(super) fn list_evals(args: EvalListArgs) -> Result<()> { let root = cache_root(args.cache_root.clone())?; let views = selected_evals(&[], EvalPack::Core) @@ -70,14 +73,30 @@ pub(super) fn definition(id: EvalId) -> EvalDefinition { EvalId::TerminalBench => EvalDefinition { id, name: "Terminal-Bench", - repo_url: "https://github.com/harbor-framework/terminal-bench.git", - repo_ref: "main", - cache_name: "terminal-bench", + repo_url: HARBOR_REPO, + repo_ref: HARBOR_REF, + cache_name: "harbor", description: "Agent benchmark for real terminal tasks in Docker sandboxes.", disk_estimate: "5-20GB depending on Docker images", - required_tools: &["git", "uv", "python3.12", "tb", "docker"], - sync_notes: &["Clones task repo and installs the terminal-bench uv tool."], - run_notes: &["Runs the Terminal-Bench CLI against the full selected dataset."], + required_tools: &["git", "uv", "python3.12", "docker"], + sync_notes: &["Clones the pinned Harbor checkout and prepares its uv environment."], + run_notes: &[ + "Runs Terminal-Bench through Harbor jobs; the legacy tb runner is unsupported.", + ], + }, + EvalId::SweGym => EvalDefinition { + id, + name: "SWE-Gym Lite", + repo_url: HARBOR_REPO, + repo_ref: HARBOR_REF, + cache_name: "harbor", + description: "Software-engineering agent benchmark over the Harbor SWE-Gym Lite adapter.", + disk_estimate: "10-30GB including task Docker images", + required_tools: &["git", "uv", "python3.12", "docker"], + sync_notes: &["Clones the pinned Harbor checkout and prepares its uv environment."], + run_notes: &[ + "Generates SWE-Gym Lite tasks with Harbor's adapter, then runs Harbor jobs.", + ], }, EvalId::SweBenchPro => EvalDefinition { id, diff --git a/crates/skippy-bench/src/evals/run.rs b/crates/skippy-bench/src/evals/run.rs index 24334ac93..b2deacd7c 100644 --- a/crates/skippy-bench/src/evals/run.rs +++ b/crates/skippy-bench/src/evals/run.rs @@ -179,9 +179,9 @@ pub(super) fn run_artifacts(definition: EvalDefinition, run_dir: &Path) -> Vec artifacts.push(RunArtifact { - kind: "terminal-bench-results", - path: terminal_bench_output_path(run_dir).display().to_string(), + EvalId::TerminalBench | EvalId::SweGym => artifacts.push(RunArtifact { + kind: "harbor-jobs", + path: harbor_jobs_output_path(run_dir).display().to_string(), }), } artifacts @@ -192,7 +192,7 @@ fn collect_metrics(definition: EvalDefinition, run_dir: &Path, duration_ms: f64) EvalId::SpeedBench => speed_bench_metrics(run_dir).unwrap_or_default(), EvalId::SweBenchPro => swe_bench_pro_metrics(run_dir).unwrap_or_default(), EvalId::McpAtlas => mcp_atlas_metrics(run_dir).unwrap_or_default(), - EvalId::TerminalBench => terminal_bench_metrics(run_dir).unwrap_or_default(), + EvalId::TerminalBench | EvalId::SweGym => harbor_metrics(run_dir).unwrap_or_default(), }; metrics.duration_ms = Some(duration_ms); fill_client_rates(&mut metrics, duration_ms); @@ -344,6 +344,81 @@ pub(super) fn terminal_bench_metrics(run_dir: &Path) -> Result { Ok(metrics) } +pub(super) fn harbor_metrics(run_dir: &Path) -> Result { + let root = harbor_jobs_output_path(run_dir); + let mut rewards = Vec::new(); + collect_harbor_trial_rewards(&root, &mut rewards)?; + if rewards.is_empty() { + bail!("no Harbor trial results under {}", root.display()); + } + let resolved = rewards + .iter() + .filter(|reward| reward.is_some_and(|reward| reward >= 1.0)) + .count() as u64; + let total = rewards.len() as u64; + Ok(EvalMetrics { + request_count: Some(total), + failed_count: Some(total.saturating_sub(resolved)), + pass_rate: Some(resolved as f64 / total as f64), + ..EvalMetrics::default() + }) +} + +fn collect_harbor_trial_rewards(root: &Path, rewards: &mut Vec>) -> Result<()> { + if !root.exists() { + bail!("Harbor jobs directory does not exist: {}", root.display()); + } + for entry in fs::read_dir(root).with_context(|| format!("read {}", root.display()))? { + let path = entry?.path(); + if !path.is_dir() { + continue; + } + + let job_result = path.join("result.json"); + let expected_trials = if job_result.is_file() { + read_json(&job_result) + .ok() + .and_then(|value| value.get("n_total_trials").and_then(Value::as_u64)) + } else { + None + }; + let mut child_count = 0_u64; + for child in fs::read_dir(&path).with_context(|| format!("read {}", path.display()))? { + let child_path = child?.path(); + if !child_path.is_dir() { + continue; + } + child_count += 1; + let trial_result = child_path.join("result.json"); + let reward = fs::read(&trial_result) + .ok() + .and_then(|bytes| serde_json::from_slice::(&bytes).ok()) + .and_then(|value| harbor_trial_reward(&value)); + rewards.push(reward); + } + if let Some(expected_trials) = expected_trials { + rewards.extend( + std::iter::repeat(None).take(expected_trials.saturating_sub(child_count) as usize), + ); + } + } + Ok(()) +} + +fn harbor_trial_reward(value: &Value) -> Option { + let rewards = value + .get("verifier_result") + .and_then(|result| result.get("rewards"))?; + rewards + .get("reward") + .or_else(|| { + rewards + .as_object() + .and_then(|object| object.values().next()) + }) + .and_then(Value::as_f64) +} + fn terminal_bench_results_path(run_dir: &Path) -> Result { let root = terminal_bench_output_path(run_dir); for entry in fs::read_dir(&root).with_context(|| format!("read {}", root.display()))? { @@ -454,6 +529,10 @@ pub(super) fn terminal_bench_output_path(run_dir: &Path) -> PathBuf { run_dir.join("raw/terminal-bench") } +pub(super) fn harbor_jobs_output_path(run_dir: &Path) -> PathBuf { + run_dir.join("raw/harbor-jobs") +} + fn metrics_report_path(run_dir: &Path) -> PathBuf { run_dir.join("raw/metrics-report.json") } diff --git a/crates/skippy-bench/src/evals/sync.rs b/crates/skippy-bench/src/evals/sync.rs index ce4d776e7..23bdf7c0a 100644 --- a/crates/skippy-bench/src/evals/sync.rs +++ b/crates/skippy-bench/src/evals/sync.rs @@ -28,17 +28,23 @@ fn sync_repo(definition: EvalDefinition, root: &Path, dry_run: bool) -> Result<( return Ok(()); } - run_step( - &CommandSpec::new("git").args([ - "clone", - "--recurse-submodules", - "--branch", - definition.repo_ref, - definition.repo_url, - &target.display().to_string(), - ]), - dry_run, - ) + for step in new_repo_sync_steps(&target, definition.repo_url, definition.repo_ref) { + run_step(&step, dry_run)?; + } + Ok(()) +} + +pub(super) fn new_repo_sync_steps( + target: &Path, + repo_url: &str, + repo_ref: &str, +) -> [CommandSpec; 3] { + let target = target.display().to_string(); + [ + CommandSpec::new("git").args(["clone", "--recurse-submodules", repo_url, &target]), + CommandSpec::new("git").args(["-C", &target, "fetch", "--prune", "origin", repo_ref]), + CommandSpec::new("git").args(["-C", &target, "checkout", "--detach", "FETCH_HEAD"]), + ] } pub(super) fn existing_repo_sync_steps(target: &Path, repo_ref: &str) -> [CommandSpec; 2] { @@ -53,14 +59,8 @@ fn sync_steps(definition: EvalDefinition, root: &Path) -> Vec { let harness = harness_dir(root, definition); match definition.id { EvalId::SpeedBench => Vec::new(), - EvalId::TerminalBench => { - vec![CommandSpec::new("uv").args([ - "tool", - "install", - "--python", - "3.12", - "terminal-bench", - ])] + EvalId::TerminalBench | EvalId::SweGym => { + vec![CommandSpec::new("uv").args(["sync"]).cwd(harness)] } EvalId::SweBenchPro => vec![ CommandSpec::new("git") diff --git a/crates/skippy-server/src/frontend/local_generation/linear_decode.rs b/crates/skippy-server/src/frontend/local_generation/linear_decode.rs index 9d14bdd11..1ba514378 100644 --- a/crates/skippy-server/src/frontend/local_generation/linear_decode.rs +++ b/crates/skippy-server/src/frontend/local_generation/linear_decode.rs @@ -206,6 +206,10 @@ impl StageOpenAiBackend { .committed_tokens .last() .ok_or_else(|| OpenAiError::backend("linear proposal receipt committed no tokens"))?; + append_pending_linear_proposal_tokens( + &mut state.pending_linear_proposal_tokens, + &receipt.committed_tokens, + ); committed_token_ids.extend_from_slice(&receipt.committed_tokens); let stopped_by_proposal = receipt.disposition == LinearProposalDisposition::Stopped; if state.emit_token_debug { @@ -247,3 +251,81 @@ impl StageOpenAiBackend { } } } + +fn append_pending_linear_proposal_tokens(pending: &mut Vec, committed_tokens: &[i32]) { + if !committed_tokens.is_empty() && !pending.ends_with(committed_tokens) { + pending.extend_from_slice(committed_tokens); + } +} + +#[cfg(test)] +mod tests { + use std::sync::{Arc, Mutex}; + + use anyhow::Result; + + use super::*; + use crate::frontend::linear_proposal::{ + LinearProposalDiscardReason, LinearProposalIngress, LinearProposalIngressConfig, + LinearProposalQuery, LinearProposalReceipt, LinearProposalSourceResponse, + query_linear_proposal, + }; + + #[derive(Default)] + struct PendingTokenIngress { + pending: Mutex>>, + } + + impl LinearProposalIngress for PendingTokenIngress { + fn propose(&self, query: LinearProposalQuery) -> Result { + self.pending.lock().unwrap().push(query.pending_token_ids); + Ok(LinearProposalSourceResponse::new(None)) + } + + fn report(&self, _receipt: &LinearProposalReceipt) -> Result<()> { + Ok(()) + } + + fn discard( + &self, + _decision_id: &crate::frontend::linear_proposal::OpaqueProposalDecisionId, + _reason: LinearProposalDiscardReason, + ) -> Result<()> { + Ok(()) + } + } + + #[test] + fn accepted_proposal_tokens_are_pending_on_the_next_query() { + let source = Arc::new(PendingTokenIngress::default()); + let config = + LinearProposalIngressConfig::new(source.clone(), Duration::from_secs(1), 4).unwrap(); + let committed_tokens = [41, 42]; + let mut pending = Vec::new(); + + append_pending_linear_proposal_tokens(&mut pending, &committed_tokens); + append_pending_linear_proposal_tokens(&mut pending, &committed_tokens); + + assert!(matches!( + query_linear_proposal( + &config, + LinearProposalQueryParams { + request_id: 7, + session_id: 8, + prompt_token_count: 2, + decode_step: 1, + committed_token_count: 3, + remaining_new_tokens: 4, + runtime_max_proposal_tokens: 4, + pending_token_ids: pending.into_boxed_slice(), + }, + ) + .unwrap(), + LinearProposalQueryOutcome::NoProposal { .. } + )); + assert_eq!( + source.pending.lock().unwrap().as_slice(), + &[committed_tokens.to_vec().into_boxed_slice()] + ); + } +} From acb95873014547e1c9f29d6295f3eaa30e5efe66 Mon Sep 17 00:00:00 2001 From: James Dumay Date: Wed, 12 Aug 2026 06:55:02 +1000 Subject: [PATCH 2/2] fix: address Harbor benchmark review feedback --- crates/skippy-bench/README.md | 7 +- crates/skippy-bench/src/cli.rs | 2 +- crates/skippy-bench/src/evals.rs | 37 +-- .../skippy-bench/src/evals/adapters/harbor.rs | 7 +- .../src/evals/adapters/terminal_bench.rs | 27 -- crates/skippy-bench/src/evals/run.rs | 45 +-- .../src/frontend/linear_proposal/execution.rs | 4 +- .../local_generation/linear_decode.rs | 259 ++++++++++++++++-- 8 files changed, 245 insertions(+), 143 deletions(-) delete mode 100644 crates/skippy-bench/src/evals/adapters/terminal_bench.rs diff --git a/crates/skippy-bench/README.md b/crates/skippy-bench/README.md index 3817b999d..6da18a756 100644 --- a/crates/skippy-bench/README.md +++ b/crates/skippy-bench/README.md @@ -229,10 +229,9 @@ SkippyBench launches it through a small adapter that adds the bearer token from `--api-key` without modifying the upstream harness. SWE-Bench Pro records OpenAI usage tokens and client-side tok/s when the upstream flow produces them. -Terminal-Bench records pass rate, resolved/unresolved task counts, token totals -when the agent reports them, and raw harness artifacts. The MCP-Atlas adapter -records wall time, raw completion CSV artifacts, the native scoring output -directory, and CSV task row count. +Terminal-Bench and SWE-Gym record Harbor trial counts, pass rates, and raw +Harbor job artifacts. The MCP-Atlas adapter records wall time, raw completion +CSV artifacts, the native scoring output directory, and CSV task row count. `eval run` requires metrics-server for every external benchmark. `--metrics-http` defaults to `http://127.0.0.1:18080`; the command creates a metrics-server run before the diff --git a/crates/skippy-bench/src/cli.rs b/crates/skippy-bench/src/cli.rs index 8a8597c4a..d76517e57 100644 --- a/crates/skippy-bench/src/cli.rs +++ b/crates/skippy-bench/src/cli.rs @@ -58,7 +58,7 @@ pub enum EvalCommandKind { Sync(EvalSyncArgs), Install(EvalSyncArgs), Doctor(EvalDoctorArgs), - Run(EvalRunArgs), + Run(Box), } #[derive(Clone, Copy, Debug, Eq, PartialEq, ValueEnum)] diff --git a/crates/skippy-bench/src/evals.rs b/crates/skippy-bench/src/evals.rs index 645df128a..3cd0ff5c8 100644 --- a/crates/skippy-bench/src/evals.rs +++ b/crates/skippy-bench/src/evals.rs @@ -39,7 +39,7 @@ pub fn eval_command(args: EvalArgs) -> Result<()> { EvalCommandKind::Info(args) => registry::info_eval(args), EvalCommandKind::Sync(args) | EvalCommandKind::Install(args) => sync::sync_evals(args), EvalCommandKind::Doctor(args) => doctor::doctor_evals(args), - EvalCommandKind::Run(args) => run::run_eval(args), + EvalCommandKind::Run(args) => run::run_eval(*args), } } @@ -474,7 +474,7 @@ mod tests { fill_client_rates, harbor_jobs_output_path, harbor_metrics, resolved_harness_commit, run_artifacts, speed_bench_metrics, speed_bench_output_path, speed_bench_response_timings_path, swe_bench_pro_metrics, swe_bench_pro_output_path, - telemetry_or_unavailable, terminal_bench_metrics, terminal_bench_output_path, + telemetry_or_unavailable, }, sync::{existing_repo_sync_steps, new_repo_sync_steps}, }; @@ -1036,39 +1036,6 @@ mod tests { let _ = fs::remove_dir_all(run_dir); } - #[test] - fn terminal_bench_metrics_extract_accuracy_and_tokens() { - let run_dir = temp_run_dir("terminal-metrics"); - let result_dir = terminal_bench_output_path(&run_dir).join("run-id"); - fs::create_dir_all(&result_dir).unwrap(); - fs::write( - result_dir.join("results.json"), - r#"{ - "results": [ - { - "task_id": "hello-world", - "is_resolved": true, - "total_input_tokens": 100, - "total_output_tokens": 25 - } - ], - "n_resolved": 1, - "n_unresolved": 0, - "accuracy": 1.0 - }"#, - ) - .unwrap(); - - let metrics = terminal_bench_metrics(&run_dir).unwrap(); - assert_eq!(metrics.request_count, Some(1)); - assert_eq!(metrics.failed_count, Some(0)); - assert_eq!(metrics.pass_rate, Some(1.0)); - assert_eq!(metrics.prompt_tokens, Some(100)); - assert_eq!(metrics.completion_tokens, Some(25)); - assert_eq!(metrics.total_tokens, Some(125)); - let _ = fs::remove_dir_all(run_dir); - } - #[test] fn harbor_metrics_normalize_trial_rewards() { let run_dir = temp_run_dir("harbor-metrics"); diff --git a/crates/skippy-bench/src/evals/adapters/harbor.rs b/crates/skippy-bench/src/evals/adapters/harbor.rs index 745653169..9c6a3cde3 100644 --- a/crates/skippy-bench/src/evals/adapters/harbor.rs +++ b/crates/skippy-bench/src/evals/adapters/harbor.rs @@ -56,12 +56,7 @@ pub(in crate::evals) fn harbor_command( let task_path = args .task_id .as_deref() - .map(|task_id| { - format!( - "{}", - shell_quote(&task_root.join(task_id).display().to_string()) - ) - }) + .map(|task_id| shell_quote(&task_root.join(task_id).display().to_string())) .unwrap_or_else(|| shell_quote(&task_root.display().to_string())); format!( "set -- -p {}\nuv run --with swebench adapters/swegym/run_adapter.py --dataset {}{} --task-dir {}\n", diff --git a/crates/skippy-bench/src/evals/adapters/terminal_bench.rs b/crates/skippy-bench/src/evals/adapters/terminal_bench.rs deleted file mode 100644 index 4b88d209d..000000000 --- a/crates/skippy-bench/src/evals/adapters/terminal_bench.rs +++ /dev/null @@ -1,27 +0,0 @@ -use super::super::{run::terminal_bench_output_path, *}; - -pub(in crate::evals) fn terminal_bench_command(args: &EvalRunArgs, run_dir: &Path) -> CommandSpec { - let model = litellm_model_name(&args.model); - CommandSpec::new("tb") - .args([ - "run".to_string(), - "--dataset".to_string(), - "terminal-bench-core==0.1.1".to_string(), - "--agent".to_string(), - "terminus".to_string(), - "--model".to_string(), - model, - "--n-concurrent".to_string(), - args.endpoint_concurrency.to_string(), - "--output-path".to_string(), - terminal_bench_output_path(run_dir).display().to_string(), - "--global-agent-timeout-sec".to_string(), - args.timeout_secs.to_string(), - "--global-test-timeout-sec".to_string(), - "60".to_string(), - "--no-upload-results".to_string(), - "--no-livestream".to_string(), - ]) - .env("OPENAI_BASE_URL", args.base_url.clone()) - .secret_env("OPENAI_API_KEY", args.api_key.clone()) -} diff --git a/crates/skippy-bench/src/evals/run.rs b/crates/skippy-bench/src/evals/run.rs index b2deacd7c..8fac424fc 100644 --- a/crates/skippy-bench/src/evals/run.rs +++ b/crates/skippy-bench/src/evals/run.rs @@ -322,28 +322,6 @@ fn swe_bench_pro_pass_rate(value: &Value) -> Option { (total > 0).then_some(resolved as f64 / total as f64) } -pub(super) fn terminal_bench_metrics(run_dir: &Path) -> Result { - let value = read_json(&terminal_bench_results_path(run_dir)?)?; - let results = value - .get("results") - .and_then(Value::as_array) - .map(Vec::as_slice) - .unwrap_or(&[]); - let mut metrics = EvalMetrics { - request_count: Some(results.len() as u64), - failed_count: value.get("n_unresolved").and_then(Value::as_u64), - pass_rate: value.get("accuracy").and_then(Value::as_f64), - ..EvalMetrics::default() - }; - metrics.prompt_tokens = sum_u64_field(results, "total_input_tokens"); - metrics.completion_tokens = sum_u64_field(results, "total_output_tokens"); - metrics.total_tokens = match (metrics.prompt_tokens, metrics.completion_tokens) { - (Some(prompt), Some(completion)) => Some(prompt + completion), - _ => None, - }; - Ok(metrics) -} - pub(super) fn harbor_metrics(run_dir: &Path) -> Result { let root = harbor_jobs_output_path(run_dir); let mut rewards = Vec::new(); @@ -397,9 +375,10 @@ fn collect_harbor_trial_rewards(root: &Path, rewards: &mut Vec>) -> rewards.push(reward); } if let Some(expected_trials) = expected_trials { - rewards.extend( - std::iter::repeat(None).take(expected_trials.saturating_sub(child_count) as usize), - ); + rewards.extend(std::iter::repeat_n( + None, + expected_trials.saturating_sub(child_count) as usize, + )); } } Ok(()) @@ -419,18 +398,6 @@ fn harbor_trial_reward(value: &Value) -> Option { .and_then(Value::as_f64) } -fn terminal_bench_results_path(run_dir: &Path) -> Result { - let root = terminal_bench_output_path(run_dir); - for entry in fs::read_dir(&root).with_context(|| format!("read {}", root.display()))? { - let entry = entry?; - let path = entry.path().join("results.json"); - if path.is_file() { - return Ok(path); - } - } - bail!("no Terminal-Bench results.json under {}", root.display()) -} - fn mcp_atlas_metrics(run_dir: &Path) -> Result { let mut reader = csv::Reader::from_path(mcp_atlas_output_path(run_dir))?; let data_rows = reader.records().filter(|record| record.is_ok()).count(); @@ -525,10 +492,6 @@ pub(super) fn mcp_atlas_score_dir(run_dir: &Path) -> PathBuf { run_dir.join("raw/mcp-atlas-evaluation-results") } -pub(super) fn terminal_bench_output_path(run_dir: &Path) -> PathBuf { - run_dir.join("raw/terminal-bench") -} - pub(super) fn harbor_jobs_output_path(run_dir: &Path) -> PathBuf { run_dir.join("raw/harbor-jobs") } diff --git a/crates/skippy-server/src/frontend/linear_proposal/execution.rs b/crates/skippy-server/src/frontend/linear_proposal/execution.rs index 64bc28ae6..b0576b86f 100644 --- a/crates/skippy-server/src/frontend/linear_proposal/execution.rs +++ b/crates/skippy-server/src/frontend/linear_proposal/execution.rs @@ -63,7 +63,7 @@ impl StageOpenAiBackend { params: LinearProposalExecutionParams<'_>, queried: QueriedLinearProposal, cancellation: Option<&openai_frontend::CancellationToken>, - on_token: &mut impl FnMut(i32) -> OpenAiResult, + on_token: &mut (impl FnMut(i32) -> OpenAiResult + ?Sized), ) -> OpenAiResult> { ensure_request_active(cancellation)?; let proposal_token_count = queried.proposal.token_ids.len(); @@ -142,7 +142,7 @@ impl StageOpenAiBackend { proposal_tokens: &[i32], verify_inputs: &[i32], cancellation: Option<&openai_frontend::CancellationToken>, - on_token: &mut impl FnMut(i32) -> OpenAiResult, + on_token: &mut (impl FnMut(i32) -> OpenAiResult + ?Sized), ) -> OpenAiResult> { ensure_request_active(cancellation)?; let verify_timer = Instant::now(); diff --git a/crates/skippy-server/src/frontend/local_generation/linear_decode.rs b/crates/skippy-server/src/frontend/local_generation/linear_decode.rs index 1ba514378..3beef077a 100644 --- a/crates/skippy-server/src/frontend/local_generation/linear_decode.rs +++ b/crates/skippy-server/src/frontend/local_generation/linear_decode.rs @@ -22,6 +22,36 @@ impl StageOpenAiBackend { state: &mut DecodeState, emit_token: &mut impl FnMut(i32) -> OpenAiResult, ) -> OpenAiResult { + self.try_execute_linear_proposal_with_executor( + request, + session_id, + state, + emit_token, + |backend, params, queried, cancellation, emit_token| { + backend.execute_local_linear_proposal(params, queried, cancellation, emit_token) + }, + ) + } + + fn try_execute_linear_proposal_with_executor( + &self, + request: &LocalGeneration<'_>, + session_id: &str, + state: &mut DecodeState, + emit_token: &mut impl FnMut(i32) -> OpenAiResult, + execute: F, + ) -> OpenAiResult + where + F: for<'a> FnOnce( + &StageOpenAiBackend, + LinearProposalExecutionParams<'a>, + crate::frontend::linear_proposal::QueriedLinearProposal, + Option<&'a openai_frontend::CancellationToken>, + &'a mut dyn FnMut(i32) -> OpenAiResult, + ) -> OpenAiResult< + Option, + >, + { let Some(config) = self.linear_proposal_ingress.as_ref() else { return Ok(LinearProposalProgress::NotUsed); }; @@ -137,7 +167,8 @@ impl StageOpenAiBackend { }; let decision_id = queried.proposal.decision_id.clone(); let receipt = execute_linear_proposal_with_terminal_discard(config, &decision_id, || { - self.execute_local_linear_proposal( + execute( + self, LinearProposalExecutionParams { request_id: request.ids.request_id, request_session_id: request.ids.session_id, @@ -260,29 +291,56 @@ fn append_pending_linear_proposal_tokens(pending: &mut Vec, committed_token #[cfg(test)] mod tests { - use std::sync::{Arc, Mutex}; + use std::sync::{Arc, Mutex, atomic::AtomicUsize}; use anyhow::Result; + use skippy_protocol::{LoadMode, StageConfig}; + use skippy_runtime::SamplingConfig; + use tokio::sync::Semaphore; use super::*; + use crate::binary_transport::DecodeFrameBatcher; + use crate::frontend::admission::GenerationTokenBudget; + use crate::frontend::decode_batcher::DecodeBatcher; + use crate::frontend::generation::{OpenAiBackendMode, OpenAiCacheHints, OpenAiGenerationIds}; use crate::frontend::linear_proposal::{ - LinearProposalDiscardReason, LinearProposalIngress, LinearProposalIngressConfig, - LinearProposalQuery, LinearProposalReceipt, LinearProposalSourceResponse, - query_linear_proposal, + LinearProposal, LinearProposalDiscardReason, LinearProposalIngress, + LinearProposalIngressConfig, LinearProposalQuery, LinearProposalReceipt, + LinearProposalSourceResponse, OpaqueProposalDecisionId, }; + use crate::frontend::native_mtp::{NativeMtpDecodeOptions, NativeMtpVerifier}; + use crate::frontend::{EmbeddedOpenAiRequestDefaults, SpeculativeDecodeConfig}; + use crate::runtime_state::RuntimeState; + use crate::telemetry::{Telemetry, TelemetryLevel}; #[derive(Default)] struct PendingTokenIngress { pending: Mutex>>, + reports: Mutex>>, + proposals: AtomicUsize, } impl LinearProposalIngress for PendingTokenIngress { fn propose(&self, query: LinearProposalQuery) -> Result { self.pending.lock().unwrap().push(query.pending_token_ids); - Ok(LinearProposalSourceResponse::new(None)) + if self + .proposals + .fetch_add(1, std::sync::atomic::Ordering::Relaxed) + == 0 + { + Ok(LinearProposalSourceResponse::new(Some( + LinearProposal::new(OpaqueProposalDecisionId::new("test-decision")?, [41, 42]), + ))) + } else { + Ok(LinearProposalSourceResponse::new(None)) + } } - fn report(&self, _receipt: &LinearProposalReceipt) -> Result<()> { + fn report(&self, receipt: &LinearProposalReceipt) -> Result<()> { + self.reports + .lock() + .unwrap() + .push(receipt.committed_tokens.clone()); Ok(()) } @@ -300,32 +358,179 @@ mod tests { let source = Arc::new(PendingTokenIngress::default()); let config = LinearProposalIngressConfig::new(source.clone(), Duration::from_secs(1), 4).unwrap(); - let committed_tokens = [41, 42]; - let mut pending = Vec::new(); + let stage_config = StageConfig { + run_id: "linear-proposal-test".to_string(), + topology_id: "linear-proposal-test".to_string(), + model_id: "linear-proposal-test".to_string(), + package_ref: None, + manifest_sha256: None, + source_model_path: None, + source_model_sha256: None, + source_model_bytes: None, + materialized_path: None, + materialized_pinned: false, + model_path: None, + projector_path: None, + stage_id: "stage-0".to_string(), + stage_index: 0, + layer_start: 0, + layer_end: 1, + ctx_size: 128, + lane_count: 1, + n_batch: Some(4), + n_ubatch: Some(4), + n_gpu_layers: 0, + mmap: Some(true), + mlock: false, + cache_type_k: "f16".to_string(), + cache_type_v: "f16".to_string(), + flash_attn_type: Default::default(), + filter_tensors_on_load: false, + selected_device: None, + kv_cache: None, + native_mtp_enabled: false, + load_mode: LoadMode::RuntimeSlice, + bind_addr: "127.0.0.1:0".to_string(), + upstream: None, + downstream: None, + }; + let runtime = Arc::new(Mutex::new(RuntimeState::new_modelless_for_test(1))); + let speculative = SpeculativeDecodeConfig::default(); + let backend = StageOpenAiBackend { + runtime: runtime.clone(), + config: stage_config.clone(), + telemetry: Telemetry::new(None, 1, stage_config, TelemetryLevel::Off), + model_id: "linear-proposal-test".to_string(), + default_max_tokens: 4, + request_defaults: EmbeddedOpenAiRequestDefaults::default(), + ctx_size: 128, + mode: OpenAiBackendMode::LocalRuntime, + draft: None, + speculative_window: 0, + adaptive_speculative_window: false, + ngram_max: 0, + speculative: speculative.clone(), + generation_limit: Arc::new(Semaphore::new(1)), + generation_queue_depth: Arc::new(AtomicUsize::new(0)), + generation_queue_limit: 1, + generation_token_budget: Arc::new(GenerationTokenBudget::new(128)), + hook_policy: None, + generation_receipt: None, + linear_proposal_ingress: Some(config), + kv: None, + decode_batcher: DecodeBatcher::new(runtime.clone(), 1), + decode_frame_batcher: DecodeFrameBatcher::new(runtime, 1), + }; + let sampling = SamplingConfig::default(); + let ids = OpenAiGenerationIds::new(OpenAiCacheHints::default(), None); + let prompt_token_ids = [1, 2]; + let request = LocalGeneration { + prompt_token_ids: &prompt_token_ids, + max_tokens: 5, + sampling: &sampling, + chat_sampling_metadata: None, + speculative: &speculative, + native_mtp_enabled: false, + hook_request: None, + hook_runtime: None, + cancellation: None, + ids: &ids, + }; + let mut state = DecodeState { + decoded_tokens: 0, + current: 2, + stopped: false, + runtime_lock_wait_ms: 0.0, + runtime_lock_wait_max_ms: 0.0, + runtime_lock_hold_ms: 0.0, + runtime_lock_hold_max_ms: 0.0, + runtime_lock_acquires: 0, + runtime_sessions_before: None, + runtime_sessions_after: None, + hook_request: None, + hook_runtime: None, + generation_hooks_active: false, + linear_proposal_max_tokens: 4, + linear_context_tokens: Some(prompt_token_ids.to_vec()), + pending_linear_proposal_tokens: Vec::new(), + emit_token_debug: false, + native_mtp_options: NativeMtpDecodeOptions::from_config(&speculative), + native_mtp: NativeMtpVerifier::default(), + post_prefill_hook_checked: false, + last_mid_generation_hook_at: None, + }; + let mut emitted = Vec::new(); + + assert!(matches!( + backend + .try_execute_linear_proposal_with_executor( + &request, + "linear-proposal-session", + &mut state, + &mut |token| { + emitted.push(token); + Ok(TokenControl::Continue) + }, + |_backend, _params, queried, _cancellation, _emit_token| { + Ok(Some(LinearProposalReceipt { + request_id: 7, + session_id: 8, + decision_id: queried.proposal.decision_id, + disposition: LinearProposalDisposition::FullAccept, + proposal_token_count: 2, + verification_rows: 3, + accepted_proposal_tokens: 2, + committed_tokens: [41, 42].into(), + verification_row_predictions: [41, 42, 43].into(), + canonical_prediction_count: 2, + correction_or_boundary_token: Some(43), + base_position: 1, + position_after_verification: 4, + canonical_position: 3, + trimmed_rows: 1, + proposal_elapsed_us: 0, + verification_elapsed_us: 0, + repair_elapsed_us: 0, + total_elapsed_us: 0, + runtime_lock_wait_us: 0, + runtime_lock_hold_us: 0, + runtime_lock_acquires: 0, + })) + }, + ) + .unwrap(), + LinearProposalProgress::Continue + )); - append_pending_linear_proposal_tokens(&mut pending, &committed_tokens); - append_pending_linear_proposal_tokens(&mut pending, &committed_tokens); + assert_eq!(state.pending_linear_proposal_tokens, vec![41, 42]); + assert_eq!( + state.linear_context_tokens.as_deref(), + Some(&[1, 2, 41, 42][..]) + ); assert!(matches!( - query_linear_proposal( - &config, - LinearProposalQueryParams { - request_id: 7, - session_id: 8, - prompt_token_count: 2, - decode_step: 1, - committed_token_count: 3, - remaining_new_tokens: 4, - runtime_max_proposal_tokens: 4, - pending_token_ids: pending.into_boxed_slice(), - }, - ) - .unwrap(), - LinearProposalQueryOutcome::NoProposal { .. } + backend + .try_execute_linear_proposal_with_executor( + &request, + "linear-proposal-session", + &mut state, + &mut |_| Ok(TokenControl::Continue), + |_backend, _params, _queried, _cancellation, _emit_token| Ok(None), + ) + .unwrap(), + LinearProposalProgress::NotUsed )); assert_eq!( source.pending.lock().unwrap().as_slice(), - &[committed_tokens.to_vec().into_boxed_slice()] + &[ + Vec::::new().into_boxed_slice(), + vec![41, 42].into_boxed_slice(), + ] + ); + assert_eq!( + source.reports.lock().unwrap().as_slice(), + &[vec![41, 42].into_boxed_slice()] ); + assert!(emitted.is_empty()); } }