diff --git a/.github/workflows/session-benchmarks.yml b/.github/workflows/session-benchmarks.yml new file mode 100644 index 00000000..fc20d94c --- /dev/null +++ b/.github/workflows/session-benchmarks.yml @@ -0,0 +1,108 @@ +name: Session Benchmarks + +on: + pull_request: + branches: + - main + - release* + - release/* + - release-* + schedule: + - cron: "23 7 * * 2" + workflow_dispatch: + inputs: + mode: + description: Workload size + required: true + default: stress + type: choice + options: + - fast + - stress + +permissions: + contents: read + +jobs: + session: + name: ${{ github.event_name == 'pull_request' && 'Fast' || inputs.mode || 'Stress' }} (${{ matrix.os }}) + runs-on: ${{ matrix.os }} + timeout-minutes: 45 + strategy: + fail-fast: false + matrix: + os: + - ubuntu-latest + - windows-latest + - macos-latest + env: + CARGO_TARGET_DIR: ${{ github.workspace }}/.target-session-benchmarks + PET_SESSION_STRESS: ${{ (github.event_name == 'schedule' || inputs.mode == 'stress') && '1' || '' }} + RUST_BACKTRACE: 1 + RUST_LOG: warn + steps: + - name: Checkout + uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 + with: + persist-credentials: false + + - name: Set up Python fixture + uses: actions/setup-python@5fda3b95a4ea91299a34e894583c3862153e4b97 # v7.0.0 + with: + python-version: "3.12" + + - name: Set up Rust + uses: dtolnay/rust-toolchain@stable + with: + toolchain: stable + + - name: Test session metrics parser + run: python -B -m unittest scripts.tests.test_session_metrics -v + + - name: Run session benchmark + id: benchmark + timeout-minutes: 35 + run: | + mode=fast + if [ "${PET_SESSION_STRESS}" = "1" ]; then + mode=stress + fi + + set +e + cargo test --release --features ci-perf -p pet --test session_performance long_lived_session_benchmark -- --nocapture 2>&1 | tee session-output.txt + benchmark_status=${PIPESTATUS[0]} + python -B scripts/session_metrics.py extract \ + --input session-output.txt \ + --output session-metrics.json \ + --benchmark-status "$benchmark_status" \ + --mode "$mode" + metrics_status=$? + set -e + rm -f session-output.txt + if [ "$benchmark_status" -ne 0 ]; then + exit "$benchmark_status" + fi + exit "$metrics_status" + shell: bash + + - name: Ensure failure metrics + if: always() + run: | + mode=fast + if [ "${PET_SESSION_STRESS}" = "1" ]; then + mode=stress + fi + python -B scripts/session_metrics.py ensure-failure \ + --output session-metrics.json \ + --mode "$mode" \ + --reason benchmark-timeout-or-interruption + rm -f session-output.txt + shell: bash + + - name: Upload session metrics + if: always() + uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a # v7.0.1 + with: + name: session-metrics-${{ matrix.os }} + path: session-metrics.json + if-no-files-found: error diff --git a/crates/pet/tests/fixtures/session_sitecustomize.py b/crates/pet/tests/fixtures/session_sitecustomize.py new file mode 100644 index 00000000..30258709 --- /dev/null +++ b/crates/pet/tests/fixtures/session_sitecustomize.py @@ -0,0 +1,30 @@ +# Copyright (c) Microsoft Corporation. +# Licensed under the MIT License. + +import os +from pathlib import Path +import sys +import time + + +barrier = os.environ.get("PET_SESSION_RESOLVE_BARRIER") +resolve_root = os.environ.get("PET_SESSION_RESOLVE_ROOT") +if barrier and resolve_root: + prefix = os.path.normcase(os.path.realpath(sys.prefix)) + root_prefix = os.path.normcase(os.path.realpath(resolve_root)) + os.sep + if not prefix.startswith(root_prefix): + barrier = None +else: + barrier = None + +if barrier: + root = Path(barrier) + (root / f"entered-{os.getpid()}").write_text("entered", encoding="ascii") + deadline = time.monotonic() + 20 + release = root / "release" + while not release.exists(): + if time.monotonic() >= deadline: + (root / f"failed-{os.getpid()}").write_text("timeout", encoding="ascii") + os._exit(86) + time.sleep(0.005) + (root / f"released-{os.getpid()}").write_text("released", encoding="ascii") diff --git a/crates/pet/tests/jsonrpc_client.rs b/crates/pet/tests/jsonrpc_client.rs index 2fb12e41..57260516 100644 --- a/crates/pet/tests/jsonrpc_client.rs +++ b/crates/pet/tests/jsonrpc_client.rs @@ -57,6 +57,7 @@ pub struct EnvironmentNotification { pub kind: Option, pub name: Option, pub prefix: Option, + pub version: Option, pub error: Option, } @@ -70,10 +71,78 @@ pub struct ManagerNotification { pub struct JsonRpcNotification { pub method: String, pub params: Value, + pub received_at: Instant, +} + +pub struct PendingRequest { + id: u32, + method: String, + submitted_at: Instant, + receiver: mpsc::Receiver, + client: PetJsonRpcClient, +} + +struct TimedResponse { + result: Result, + received_at: Instant, +} + +pub struct TimedRefresh { + pub result: RefreshResult, + pub submitted_at: Instant, + pub round_trip: Duration, +} + +impl PendingRequest { + pub fn submitted_at(&self) -> Instant { + self.submitted_at + } + + pub fn wait(self, timeout: Duration) -> Result<(Value, Duration), String> { + let remaining = timeout.saturating_sub(self.submitted_at.elapsed()); + let response = match self.receiver.recv_timeout(remaining) { + Ok(response) => response, + Err(mpsc::RecvTimeoutError::Timeout) => { + return Err(format!( + "Timed out waiting for {} response after {timeout:?}", + self.method + )); + } + Err(mpsc::RecvTimeoutError::Disconnected) => { + return Err(format!( + "Response channel disconnected while waiting for {}; stderr: {}", + self.method, + self.client.stderr_output() + )); + } + }; + let round_trip = response + .received_at + .saturating_duration_since(self.submitted_at); + if round_trip > timeout { + return Err(format!( + "Timed out waiting for {} response after {timeout:?}", + self.method + )); + } + response.result.map(|result| (result, round_trip)) + } +} + +impl Drop for PendingRequest { + fn drop(&mut self) { + self.client + .inner + .state + .pending + .lock() + .unwrap() + .remove(&self.id); + } } struct ClientState { - pending: Mutex>>>, + pending: Mutex>>, notifications: Mutex>, stderr_lines: Mutex>, } @@ -189,6 +258,10 @@ impl PetJsonRpcClient { Ok(status) } + pub fn process_id(&self) -> u32 { + self.inner.child.lock().unwrap().id() + } + pub fn configure(&self, config: Value) -> Result<(), String> { self.send_request_value("configure", config, DEFAULT_REQUEST_TIMEOUT) .map(|_| ()) @@ -202,6 +275,19 @@ impl PetJsonRpcClient { ) } + pub fn refresh_with_timing(&self, params: Option) -> Result { + let pending = self.submit_request_value("refresh", params.unwrap_or_else(|| json!({})))?; + let submitted_at = pending.submitted_at(); + let (result, round_trip) = pending.wait(DEFAULT_REQUEST_TIMEOUT)?; + let result = serde_json::from_value(result) + .map_err(|e| format!("Failed to deserialize response for refresh: {e}"))?; + Ok(TimedRefresh { + result, + submitted_at, + round_trip, + }) + } + pub fn info(&self) -> Result { self.send_request("info", json!({}), DEFAULT_REQUEST_TIMEOUT) } @@ -215,6 +301,10 @@ impl PetJsonRpcClient { ) } + pub fn submit_resolve(&self, executable: &str) -> Result { + self.submit_request_value("resolve", json!({ "executable": executable })) + } + pub fn clear_notifications(&self) { self.inner.state.notifications.lock().unwrap().clear(); } @@ -327,6 +417,13 @@ impl PetJsonRpcClient { params: Value, timeout: Duration, ) -> Result { + self.submit_request_value(method, params)? + .wait(timeout) + .map(|(value, _)| value) + } + + fn submit_request_value(&self, method: &str, params: Value) -> Result { + let submitted_at = Instant::now(); let id = REQUEST_ID.fetch_add(1, Ordering::SeqCst); let request = json!({ "jsonrpc": "2.0", @@ -363,19 +460,13 @@ impl PetJsonRpcClient { return Err(err); } - match rx.recv_timeout(timeout) { - Ok(result) => result, - Err(mpsc::RecvTimeoutError::Timeout) => { - self.inner.state.pending.lock().unwrap().remove(&id); - Err(format!( - "Timed out waiting for {method} response after {timeout:?}" - )) - } - Err(mpsc::RecvTimeoutError::Disconnected) => Err(format!( - "Response channel disconnected while waiting for {method}; stderr: {}", - self.stderr_output() - )), - } + Ok(PendingRequest { + id, + method: method.to_string(), + submitted_at, + receiver: rx, + client: self.clone(), + }) } } @@ -385,6 +476,7 @@ fn spawn_stdout_reader(stdout: ChildStdout, state: Arc) -> JoinHand let read_result = loop { match read_message(&mut reader) { Ok(Some(message)) => { + let received_at = Instant::now(); if let Some(method) = message.get("method").and_then(|value| value.as_str()) { state .notifications @@ -393,20 +485,23 @@ fn spawn_stdout_reader(stdout: ChildStdout, state: Arc) -> JoinHand .push(JsonRpcNotification { method: method.to_string(), params: message.get("params").cloned().unwrap_or(Value::Null), + received_at, }); continue; } if let Some(id) = message.get("id").and_then(|value| value.as_u64()) { - if let Some(sender) = state.pending.lock().unwrap().remove(&(id as u32)) { - if let Some(error) = message.get("error") { - let _ = sender.send(Err(format!("JSONRPC error: {error:?}"))); + let sender = { state.pending.lock().unwrap().remove(&(id as u32)) }; + if let Some(sender) = sender { + let result = if let Some(error) = message.get("error") { + Err(format!("JSONRPC error: {error:?}")) } else { - let _ = sender.send(Ok(message - .get("result") - .cloned() - .unwrap_or(Value::Null))); - } + Ok(message.get("result").cloned().unwrap_or(Value::Null)) + }; + let _ = sender.send(TimedResponse { + result, + received_at, + }); } } } @@ -419,9 +514,13 @@ fn spawn_stdout_reader(stdout: ChildStdout, state: Arc) -> JoinHand Ok(()) => "PET stdout closed".to_string(), Err(err) => format!("Failed to read PET stdout: {err}"), }; + let received_at = Instant::now(); let pending = std::mem::take(&mut *state.pending.lock().unwrap()); for (_, sender) in pending { - let _ = sender.send(Err(failure.clone())); + let _ = sender.send(TimedResponse { + result: Err(failure.clone()), + received_at, + }); } }) } @@ -499,3 +598,113 @@ pub(crate) fn read_message(reader: &mut BufReader) -> io::Result (PendingRequest, impl FnOnce(Result, Instant)) { + let id = REQUEST_ID.fetch_add(1, Ordering::SeqCst); + let (sender, receiver) = mpsc::channel(); + client + .inner + .state + .pending + .lock() + .unwrap() + .insert(id, sender.clone()); + let request = PendingRequest { + id, + method: "controlled".to_string(), + submitted_at, + receiver, + client: client.clone(), + }; + let send = move |result, received_at| { + sender + .send(TimedResponse { + result, + received_at, + }) + .expect("controlled pending receiver should remain connected"); + }; + (request, send) + } + + fn pending_request_registered(client: &PetJsonRpcClient, request: &PendingRequest) -> bool { + client + .inner + .state + .pending + .lock() + .unwrap() + .contains_key(&request.id) + } + + fn pending_request_count(client: &PetJsonRpcClient) -> usize { + client.inner.state.pending.lock().unwrap().len() + } + + #[test] + fn pending_request_uses_receipt_time_and_operation_deadline() { + let client = PetJsonRpcClient::spawn().expect("failed to spawn idle PET client"); + + let submitted_at = Instant::now() - Duration::from_secs(10); + let received_at = submitted_at + Duration::from_millis(125); + let (received, send_received) = controlled_pending_request(&client, submitted_at); + send_received(Ok(json!({ "ok": true })), received_at); + let (result, latency) = received + .wait(Duration::from_secs(30)) + .expect("recorded response should be returned"); + assert_eq!(result, json!({ "ok": true })); + assert_eq!(latency, Duration::from_millis(125)); + assert_eq!(pending_request_count(&client), 0); + + let consumed_after_deadline_at = Instant::now() - Duration::from_secs(31); + let (consumed_after_deadline, send_timely) = + controlled_pending_request(&client, consumed_after_deadline_at); + send_timely( + Ok(json!({ "timely": true })), + consumed_after_deadline_at + Duration::from_secs(29), + ); + let (result, latency) = consumed_after_deadline + .wait(Duration::from_secs(30)) + .expect("a timely received response remains valid when consumed after the deadline"); + assert_eq!(result, json!({ "timely": true })); + assert_eq!(latency, Duration::from_secs(29)); + + let received_after_deadline_at = Instant::now() - Duration::from_secs(31); + let (received_after_deadline, send_late) = + controlled_pending_request(&client, received_after_deadline_at); + send_late( + Ok(json!({ "late": true })), + received_after_deadline_at + Duration::from_secs(30) + Duration::from_millis(1), + ); + let error = received_after_deadline + .wait(Duration::from_secs(30)) + .expect_err("a response received after the operation deadline must time out"); + assert!(error.contains("Timed out waiting for controlled response")); + assert_eq!(pending_request_count(&client), 0); + + let expired_at = Instant::now() - Duration::from_secs(31); + let (expired, _keep_sender_connected) = controlled_pending_request(&client, expired_at); + assert!(pending_request_registered(&client, &expired)); + let error = expired + .wait(Duration::from_secs(30)) + .expect_err("an expired operation must not receive a fresh wait budget"); + assert!(error.contains("Timed out waiting for controlled response")); + assert_eq!(pending_request_count(&client), 0); + + let (dropped, _keep_sender_connected) = controlled_pending_request(&client, Instant::now()); + assert!(pending_request_registered(&client, &dropped)); + drop(dropped); + assert_eq!( + pending_request_count(&client), + 0, + "dropping a pending request must remove its map entry" + ); + } +} diff --git a/crates/pet/tests/jsonrpc_server_test.rs b/crates/pet/tests/jsonrpc_server_test.rs index 49262fa1..054d29b4 100644 --- a/crates/pet/tests/jsonrpc_server_test.rs +++ b/crates/pet/tests/jsonrpc_server_test.rs @@ -13,7 +13,7 @@ use std::thread::JoinHandle; use std::time::{Duration, Instant}; use tempfile::TempDir; -mod jsonrpc_client; +pub mod jsonrpc_client; use jsonrpc_client::{EnvironmentNotification, PetJsonRpcClient}; diff --git a/crates/pet/tests/session_performance.rs b/crates/pet/tests/session_performance.rs new file mode 100644 index 00000000..5901765c --- /dev/null +++ b/crates/pet/tests/session_performance.rs @@ -0,0 +1,1379 @@ +// Copyright (c) Microsoft Corporation. +// Licensed under the MIT License. + +use pet_fs::path::norm_case; +use serde_json::json; +use std::collections::BTreeMap; +use std::ffi::OsString; +use std::fs; +use std::io; +use std::path::{Path, PathBuf}; +use std::process::Command; +use std::thread; +use std::time::{Duration, Instant}; +use tempfile::TempDir; + +#[allow(dead_code)] +mod jsonrpc_client; +use jsonrpc_client::{ + EnvironmentNotification, JsonRpcNotification, PendingRequest, PetJsonRpcClient, +}; + +const REQUEST_TIMEOUT: Duration = Duration::from_secs(30); +const FAST_SIZES: &[usize] = &[1, 10, 100]; +const STRESS_SIZES: &[usize] = &[1, 10, 100, 1000]; + +#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)] +struct EnvironmentIdentity { + executable: PathBuf, + prefix: PathBuf, + kind: String, + name: Option, + version: Option, +} + +#[derive(Debug, Clone)] +struct ResourceSample { + resident_bytes: u64, + threads: Option, + descriptors: Option, +} + +struct RefreshMeasurement { + inventory: Vec, + round_trip_us: u128, + ttfe_us: u128, + ambient_environment_count: usize, + ambient_manager_count: usize, +} + +#[derive(Default)] +struct AmbientCountPeak { + environment_count: usize, + manager_count: usize, +} + +impl AmbientCountPeak { + fn observe(&mut self, environment_count: usize, manager_count: usize) { + self.environment_count = self.environment_count.max(environment_count); + self.manager_count = self.manager_count.max(manager_count); + } +} + +struct Fixture { + _root: TempDir, + workspace: PathBuf, + cache: PathBuf, + barrier: PathBuf, + python_path: PathBuf, + resolve_version: String, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +struct ResolveFixtureIdentity { + executable: PathBuf, + prefix: PathBuf, + version: String, +} + +struct PendingFixtureResolve { + request: PendingRequest, + expected: ResolveFixtureIdentity, +} + +impl Fixture { + fn new() -> Self { + let resolve_python = session_python(); + let resolve_version = runtime_version(&resolve_python); + let root = tempfile::tempdir().expect("failed to create session fixture root"); + let workspace = root.path().join("workspace"); + let cache = root.path().join("cache"); + let barrier = root.path().join("resolve-barrier"); + let python_path = root.path().join("python-path"); + fs::create_dir_all(&workspace).expect("failed to create fixture workspace"); + fs::create_dir_all(&cache).expect("failed to create fixture cache"); + fs::create_dir_all(&barrier).expect("failed to create resolve barrier"); + fs::create_dir_all(&python_path).expect("failed to create fixture PYTHONPATH"); + fs::copy( + Path::new(env!("CARGO_MANIFEST_DIR")) + .join("tests") + .join("fixtures") + .join("session_sitecustomize.py"), + python_path.join("sitecustomize.py"), + ) + .expect("failed to install resolve barrier module"); + Self { + _root: root, + workspace, + cache, + barrier, + python_path, + resolve_version, + } + } + + fn reset_inventory(&self, size: usize) -> Vec { + self.reset_inventory_at(&self.workspace, size) + } + + fn reset_inventory_at(&self, workspace: &Path, size: usize) -> Vec { + if workspace.exists() { + fs::remove_dir_all(workspace).expect("failed to clear fixture workspace"); + } + fs::create_dir_all(workspace).expect("failed to recreate fixture workspace"); + (0..size) + .map(|index| { + self.create_fake_environment_at(workspace, &format!("env-{index:04}"), "3.11.0") + }) + .collect() + } + + fn create_fake_environment(&self, name: &str, version: &str) -> EnvironmentIdentity { + self.create_fake_environment_at(&self.workspace, name, version) + } + + fn create_fake_environment_at( + &self, + workspace: &Path, + name: &str, + version: &str, + ) -> EnvironmentIdentity { + let prefix = workspace.join(name); + let bin = bin_directory(&prefix); + fs::create_dir_all(&bin).expect("failed to create fake environment"); + fs::write( + prefix.join("pyvenv.cfg"), + format!("version = {version}\nprompt = {name}\n"), + ) + .expect("failed to write pyvenv.cfg"); + write_version_header(&prefix, version); + let executable = python_executable(&bin, false); + fs::write(&executable, b"fixture").expect("failed to create fake Python"); + EnvironmentIdentity { + executable: norm_case(executable), + prefix: norm_case(prefix), + kind: "Venv".to_string(), + name: Some(name.to_string()), + version: Some(version.to_string()), + } + } + + fn create_resolve_environments(&self, count: usize) -> Vec { + (0..count) + .map(|index| { + self.create_resolve_environment( + &self.resolve_root().join(format!("resolve-{index}")), + ) + }) + .collect() + } + + fn create_resolve_environment(&self, prefix: &Path) -> ResolveFixtureIdentity { + let python = session_python(); + let output = Command::new(&python) + .args(["-m", "venv", "--without-pip", "--copies"]) + .arg(prefix) + .output() + .expect("failed to start Python venv fixture setup"); + assert!( + output.status.success(), + "failed to create resolve environment at {}: {}", + prefix.display(), + String::from_utf8_lossy(&output.stderr) + ); + ResolveFixtureIdentity { + executable: resolve_fixture_path(&python_executable(&bin_directory(prefix), false)), + prefix: resolve_fixture_path(prefix), + version: self.resolve_version.clone(), + } + } + + fn resolve_root(&self) -> PathBuf { + self._root.path().join("resolve-environments") + } + + fn clear_barrier(&self) { + for entry in fs::read_dir(&self.barrier).expect("failed to read resolve barrier") { + let path = entry.expect("failed to read barrier entry").path(); + fs::remove_file(path).expect("failed to clear resolve barrier entry"); + } + } + + fn reset_cache(&self) { + if self.cache.exists() { + fs::remove_dir_all(&self.cache).expect("failed to clear fixture cache"); + } + fs::create_dir_all(&self.cache).expect("failed to recreate fixture cache"); + } + + fn entered_count(&self) -> usize { + fs::read_dir(&self.barrier) + .expect("failed to read resolve barrier") + .map(|entry| entry.expect("failed to read resolve barrier entry")) + .filter(|entry| entry.file_name().to_string_lossy().starts_with("entered-")) + .count() + } + + fn released_count(&self) -> usize { + fs::read_dir(&self.barrier) + .expect("failed to read resolve barrier") + .map(|entry| entry.expect("failed to read resolve barrier entry")) + .filter(|entry| entry.file_name().to_string_lossy().starts_with("released-")) + .count() + } + + fn failed_count(&self) -> usize { + fs::read_dir(&self.barrier) + .expect("failed to read resolve barrier") + .map(|entry| entry.expect("failed to read resolve barrier entry")) + .filter(|entry| entry.file_name().to_string_lossy().starts_with("failed-")) + .count() + } + + fn wait_for_entered(&self, expected: usize) { + let deadline = Instant::now() + Duration::from_secs(20); + while Instant::now() < deadline { + if self.entered_count() == expected { + return; + } + thread::sleep(Duration::from_millis(10)); + } + panic!( + "only {} of {expected} resolve subprocesses entered the barrier", + self.entered_count() + ); + } +} + +fn session_python() -> OsString { + std::env::var_os("PET_SESSION_PYTHON").unwrap_or_else(|| { + if cfg!(windows) { + "python".into() + } else { + "python3".into() + } + }) +} + +fn runtime_version(python: &OsString) -> String { + let output = Command::new(python) + .args([ + "-I", + "-c", + "import sys; print('.'.join(str(part) for part in sys.version_info))", + ]) + .output() + .expect("failed to query Python fixture runtime version"); + assert!( + output.status.success(), + "failed to query Python fixture runtime version: {}", + String::from_utf8_lossy(&output.stderr) + ); + String::from_utf8(output.stdout) + .expect("Python fixture runtime version was not UTF-8") + .trim() + .to_string() +} + +fn resolve_fixture_path(path: &Path) -> PathBuf { + #[cfg(target_os = "macos")] + let path = fs::canonicalize(path).expect("failed to resolve macOS fixture path"); + + norm_case(path) +} + +struct BarrierReleaseGuard { + release: PathBuf, + released: bool, +} + +impl BarrierReleaseGuard { + fn new(barrier: &Path) -> Self { + Self { + release: barrier.join("release"), + released: false, + } + } + + fn release(&mut self) { + fs::write(&self.release, b"release").expect("failed to release resolve fixture"); + self.released = true; + } +} + +impl Drop for BarrierReleaseGuard { + fn drop(&mut self) { + if !self.released { + let _ = fs::write(&self.release, b"release"); + } + } +} + +fn bin_directory(prefix: &Path) -> PathBuf { + if cfg!(windows) { + prefix.join("Scripts") + } else { + prefix.join("bin") + } +} + +fn write_version_header(prefix: &Path, version: &str) { + let include = prefix.join("include"); + fs::create_dir_all(&include).expect("failed to create fixture include directory"); + fs::write( + include.join("patchlevel.h"), + format!("#define PY_VERSION \"{version}\"\n"), + ) + .expect("failed to write fixture Python version header"); +} + +fn python_executable(bin: &Path, alias: bool) -> PathBuf { + match (cfg!(windows), alias) { + (true, false) => bin.join("python.exe"), + (true, true) => bin.join("python3.exe"), + (false, false) => bin.join("python"), + (false, true) => bin.join("python3"), + } +} + +fn fixture_inventory( + notifications: Vec, + workspace: &Path, +) -> (Vec, usize) { + let workspace = norm_case(workspace); + let (fixture, ambient): (Vec<_>, Vec<_>) = notifications.into_iter().partition(|environment| { + environment + .prefix + .as_ref() + .is_some_and(|prefix| norm_case(prefix).starts_with(&workspace)) + }); + let mut identities = fixture + .into_iter() + .map(|environment| { + assert_eq!( + environment.error, None, + "fixture environment reported an error" + ); + EnvironmentIdentity { + executable: norm_case( + environment + .executable + .expect("fixture environment had no executable"), + ), + prefix: norm_case( + environment + .prefix + .expect("fixture environment had no prefix"), + ), + kind: environment + .kind + .expect("fixture environment had no classification"), + name: environment.name, + version: environment.version, + } + }) + .collect::>(); + identities.sort_unstable(); + (identities, ambient.len()) +} + +fn manager_counts(client: &PetJsonRpcClient, workspace: &Path) -> (usize, usize) { + let workspace = norm_case(workspace); + let (fixture, ambient): (Vec<_>, Vec<_>) = client + .manager_notifications() + .into_iter() + .partition(|manager| { + manager + .executable + .as_ref() + .is_some_and(|executable| norm_case(executable).starts_with(&workspace)) + }); + (fixture.len(), ambient.len()) +} + +fn first_fixture_environment( + notifications: &[JsonRpcNotification], + workspace: &Path, +) -> Option { + let workspace = norm_case(workspace); + notifications + .iter() + .filter(|notification| { + notification.method == "environment" + && notification.params["prefix"] + .as_str() + .is_some_and(|prefix| norm_case(prefix).starts_with(&workspace)) + }) + .map(|notification| notification.received_at) + .min() +} + +#[test] +fn first_result_timing_ignores_ambient_environments_and_managers() { + let workspace = Path::new("workspace"); + let start = Instant::now(); + let notification = |method: &str, prefix: Option, millis| JsonRpcNotification { + method: method.to_string(), + params: json!({ "prefix": prefix }), + received_at: start + Duration::from_millis(millis), + }; + let mut notifications = vec![ + notification("environment", None, 1), + notification("environment", Some(PathBuf::from("workspace-other")), 2), + notification("manager", Some(workspace.join("manager")), 3), + ]; + assert_eq!(first_fixture_environment(¬ifications, workspace), None); + notifications.push(notification( + "environment", + Some(workspace.join("second")), + 20, + )); + notifications.push(notification( + "environment", + Some(workspace.join("first")), + 10, + )); + assert_eq!( + first_fixture_environment(¬ifications, workspace), + Some(start + Duration::from_millis(10)) + ); +} + +fn refresh_and_measure( + client: &PetJsonRpcClient, + workspace: &Path, + ambient_count_peak: &mut AmbientCountPeak, +) -> RefreshMeasurement { + client.clear_notifications(); + let timing = client + .refresh_with_timing(Some(json!({ "searchPaths": [workspace] }))) + .expect("session refresh failed"); + let first_environment = first_fixture_environment(&client.notifications(), workspace) + .expect("fixture refresh produced no environment"); + let ttfe = first_environment + .checked_duration_since(timing.submitted_at) + .expect("environment notification preceded refresh submission"); + assert!( + ttfe <= timing.round_trip, + "time-to-first environment exceeded refresh round trip" + ); + let _server_duration = timing.result.duration; + let (inventory, ambient_environment_count) = + fixture_inventory(client.environment_notifications(), workspace); + let (fixture_manager_count, ambient_manager_count) = manager_counts(client, workspace); + assert_eq!( + fixture_manager_count, 0, + "fixture workspace unexpectedly reported an environment manager" + ); + let measurement = RefreshMeasurement { + inventory, + round_trip_us: timing.round_trip.as_micros(), + ttfe_us: ttfe.as_micros(), + ambient_environment_count, + ambient_manager_count, + }; + ambient_count_peak.observe( + measurement.ambient_environment_count, + measurement.ambient_manager_count, + ); + measurement +} + +#[test] +fn ambient_count_peak_retains_middle_refresh_maximum() { + let mut peak = AmbientCountPeak::default(); + peak.observe(1, 2); + peak.observe(7, 6); + peak.observe(3, 4); + peak.observe(2, 1); + + assert_eq!(peak.environment_count, 7); + assert_eq!(peak.manager_count, 6); +} + +fn directory_usage(root: &Path) -> (usize, u64) { + let mut pending = vec![root.to_path_buf()]; + let mut files = 0; + let mut bytes = 0; + while let Some(directory) = pending.pop() { + for entry in fs::read_dir(&directory).expect("failed to read cache directory") { + let entry = entry.expect("failed to read cache entry"); + let metadata = entry.metadata().expect("failed to read cache metadata"); + if metadata.is_dir() { + pending.push(entry.path()); + } else if metadata.is_file() { + files += 1; + bytes += metadata.len(); + } + } + } + (files, bytes) +} + +fn cache_contents(root: &Path) -> io::Result>> { + let mut pending = vec![(PathBuf::new(), root.to_path_buf())]; + let mut contents = BTreeMap::new(); + while let Some((relative_directory, directory)) = pending.pop() { + for entry in fs::read_dir(directory)? { + let entry = entry?; + let file_type = entry.file_type()?; + let relative_path = relative_directory.join(entry.file_name()); + if file_type.is_dir() { + pending.push((relative_path, entry.path())); + } else if file_type.is_file() { + contents.insert(relative_path, fs::read(entry.path())?); + } + } + } + Ok(contents) +} + +fn original_cache_entries_unchanged( + original: &BTreeMap>, + current: &BTreeMap>, +) -> bool { + original + .iter() + .all(|(path, bytes)| current.get(path) == Some(bytes)) +} + +#[test] +fn cache_preservation_allows_additions_but_rejects_original_entry_changes() { + let original = BTreeMap::from([ + (PathBuf::from("flat.json"), b"original".to_vec()), + ( + PathBuf::from("nested").join("entry.json"), + b"nested".to_vec(), + ), + ]); + let mut with_addition = original.clone(); + with_addition.insert(PathBuf::from("ambient.json"), b"ambient".to_vec()); + assert!(original_cache_entries_unchanged(&original, &with_addition)); + + let mut changed = with_addition.clone(); + changed.insert(PathBuf::from("flat.json"), b"modified".to_vec()); + assert!(!original_cache_entries_unchanged(&original, &changed)); + + let mut deleted = with_addition; + deleted.remove(&PathBuf::from("nested").join("entry.json")); + assert!(!original_cache_entries_unchanged(&original, &deleted)); +} + +fn observe_resources(pid: u32) -> ResourceSample { + let samples = (0..3) + .map(|_| { + let sample = sample_process(pid); + thread::sleep(Duration::from_millis(20)); + sample + }) + .collect::>(); + ResourceSample { + resident_bytes: samples + .iter() + .map(|sample| sample.resident_bytes) + .max() + .unwrap(), + threads: samples.iter().filter_map(|sample| sample.threads).max(), + descriptors: samples.iter().filter_map(|sample| sample.descriptors).max(), + } +} + +fn observed_resource_peak(samples: &[ResourceSample]) -> ResourceSample { + ResourceSample { + resident_bytes: samples + .iter() + .map(|sample| sample.resident_bytes) + .max() + .expect("at least one resource sample is required"), + threads: samples.iter().filter_map(|sample| sample.threads).max(), + descriptors: samples.iter().filter_map(|sample| sample.descriptors).max(), + } +} + +fn resource_json(sample: &ResourceSample) -> serde_json::Value { + json!({ + "residentBytes": sample.resident_bytes, + "threads": sample.threads, + "handlesOrDescriptors": sample.descriptors, + }) +} + +fn cache_json((files, bytes): (usize, u64)) -> serde_json::Value { + json!({ + "files": files, + "bytes": bytes, + }) +} + +fn extract_resolved_fixture( + result: serde_json::Value, + expected: &ResolveFixtureIdentity, +) -> Result { + let environment: EnvironmentNotification = serde_json::from_value(result) + .map_err(|error| format!("resolve returned an invalid environment: {error}"))?; + if let Some(error) = environment.error { + return Err(format!("resolve reported an error: {error}")); + } + let executable = environment + .executable + .ok_or_else(|| "resolved environment had no executable".to_string())?; + let prefix = environment + .prefix + .ok_or_else(|| "resolved environment had no prefix".to_string())?; + let kind = environment + .kind + .ok_or_else(|| "resolved environment had no classification".to_string())?; + let version = environment + .version + .ok_or_else(|| "resolved environment had no version".to_string())?; + let identity = EnvironmentIdentity { + executable: resolve_fixture_path(Path::new(&executable)), + prefix: resolve_fixture_path(Path::new(&prefix)), + kind, + name: environment.name, + version: Some(version), + }; + if identity.executable != expected.executable { + return Err("resolve returned a different fixture executable".to_string()); + } + if identity.prefix != expected.prefix { + return Err("resolve returned a different fixture prefix".to_string()); + } + if identity.kind != "Venv" { + return Err(format!( + "resolve returned kind {}, expected Venv", + identity.kind + )); + } + if identity.version.as_deref() != Some(expected.version.as_str()) { + return Err("resolve returned an unexpected runtime version".to_string()); + } + Ok(identity) +} + +fn submit_fixture_resolve( + client: &PetJsonRpcClient, + expected: &ResolveFixtureIdentity, +) -> PendingFixtureResolve { + let executable = expected + .executable + .to_str() + .expect("resolve fixture path was not UTF-8"); + PendingFixtureResolve { + request: client + .submit_resolve(executable) + .expect("failed to submit fixture resolve"), + expected: expected.clone(), + } +} + +impl PendingFixtureResolve { + fn submitted_at(&self) -> Instant { + self.request.submitted_at() + } + + fn wait(self, context: &str) -> (EnvironmentIdentity, Duration) { + let (result, latency) = self + .request + .wait(REQUEST_TIMEOUT) + .unwrap_or_else(|error| panic!("{context} failed: {error}")); + let identity = extract_resolved_fixture(result, &self.expected) + .unwrap_or_else(|error| panic!("{context} returned the wrong environment: {error}")); + (identity, latency) + } +} + +fn resolve_and_measure( + client: &PetJsonRpcClient, + expected: &ResolveFixtureIdentity, +) -> (EnvironmentIdentity, u128) { + let (identity, latency) = + submit_fixture_resolve(client, expected).wait("cache-control resolve"); + (identity, latency.as_micros()) +} + +#[test] +fn resolved_fixture_validation_enforces_request_identity_and_response_shape() { + let root = tempfile::tempdir().expect("failed to create resolve validation fixture"); + let create_expected = |name: &str, version: &str| { + let prefix = root.path().join(name); + let executable = python_executable(&bin_directory(&prefix), false); + fs::create_dir_all(executable.parent().unwrap()) + .expect("failed to create resolve validation environment"); + fs::write(&executable, b"fixture").expect("failed to create resolve validation executable"); + ResolveFixtureIdentity { + executable: resolve_fixture_path(&executable), + prefix: resolve_fixture_path(&prefix), + version: version.to_string(), + } + }; + let first = create_expected("first", "3.12.10.final.0"); + let second = create_expected("second", "3.13.2.final.0"); + let response = |expected: &ResolveFixtureIdentity| { + json!({ + "executable": expected.executable, + "prefix": expected.prefix, + "kind": "Venv", + "name": "fixture", + "version": expected.version, + "error": null, + }) + }; + + let valid = extract_resolved_fixture(response(&first), &first) + .expect("valid resolve response was rejected"); + assert_eq!(valid.executable, first.executable); + assert_eq!(valid.prefix, first.prefix); + assert_eq!(valid.kind, "Venv"); + assert_eq!(valid.version.as_deref(), Some(first.version.as_str())); + + assert!( + extract_resolved_fixture(response(&second), &first).is_err(), + "a response assigned to the wrong requested fixture was accepted" + ); + assert!( + extract_resolved_fixture(json!(["not", "an", "environment"]), &first).is_err(), + "malformed resolve JSON was accepted" + ); + assert!( + extract_resolved_fixture(json!({ "error": null }), &first).is_err(), + "a resolve response with missing identity fields was accepted" + ); + + let mut errored = response(&first); + errored["error"] = json!("probe failed"); + assert!( + extract_resolved_fixture(errored, &first).is_err(), + "a resolve response containing an error was accepted" + ); + + let mut wrong_kind = response(&first); + wrong_kind["kind"] = json!("VirtualEnv"); + assert!( + extract_resolved_fixture(wrong_kind, &first).is_err(), + "a resolve response with the wrong kind was accepted" + ); + + let mut wrong_version = response(&first); + wrong_version["version"] = json!("3.12.10"); + assert!( + extract_resolved_fixture(wrong_version, &first).is_err(), + "a resolve response with the wrong runtime version was accepted" + ); +} + +fn shutdown_client(client: &PetJsonRpcClient, context: &str) { + let status = client + .shutdown(Duration::from_secs(10)) + .unwrap_or_else(|error| panic!("{context} did not shut down gracefully: {error}")); + assert!( + status.success(), + "{context} exited unsuccessfully with {status}" + ); +} + +#[cfg(target_os = "linux")] +fn sample_process(pid: u32) -> ResourceSample { + let status = fs::read_to_string(format!("/proc/{pid}/status")) + .expect("failed to read process-specific Linux status"); + let value = |name: &str| { + status + .lines() + .find_map(|line| line.strip_prefix(name)) + .and_then(|value| value.split_whitespace().next()) + .and_then(|value| value.parse::().ok()) + .unwrap_or_else(|| panic!("Linux process status omitted {name}")) + }; + ResourceSample { + resident_bytes: value("VmRSS:") * 1024, + threads: Some(value("Threads:") as usize), + descriptors: Some( + fs::read_dir(format!("/proc/{pid}/fd")) + .expect("failed to read process-specific descriptor directory") + .try_fold(0, |count, entry| entry.map(|_| count + 1)) + .expect("failed to read process descriptor entry"), + ), + } +} + +#[cfg(windows)] +fn sample_process(pid: u32) -> ResourceSample { + let script = format!( + "$p=Get-Process -Id {pid} -ErrorAction Stop; \ + Write-Output \"$($p.WorkingSet64),$($p.Threads.Count),$($p.HandleCount)\"" + ); + let output = Command::new("powershell") + .args(["-NoProfile", "-NonInteractive", "-Command", &script]) + .output() + .expect("failed to sample PET process"); + assert!( + output.status.success(), + "process-specific Windows sampling failed: {}", + String::from_utf8_lossy(&output.stderr) + ); + let fields = String::from_utf8(output.stdout) + .expect("Windows resource sample was not UTF-8") + .trim() + .split(',') + .map(|field| { + field + .parse::() + .expect("invalid Windows resource field") + }) + .collect::>(); + assert_eq!(fields.len(), 3); + ResourceSample { + resident_bytes: fields[0], + threads: Some(fields[1] as usize), + descriptors: Some(fields[2] as usize), + } +} + +#[cfg(target_os = "macos")] +fn sample_process(pid: u32) -> ResourceSample { + let output = Command::new("ps") + .args(["-o", "rss=", "-p", &pid.to_string()]) + .output() + .expect("failed to sample PET process"); + assert!( + output.status.success(), + "process-specific macOS sampling failed" + ); + let fields = String::from_utf8(output.stdout) + .expect("macOS resource sample was not UTF-8") + .split_whitespace() + .map(|field| field.parse::().expect("invalid macOS resource field")) + .collect::>(); + assert_eq!(fields.len(), 1); + ResourceSample { + resident_bytes: fields[0] * 1024, + threads: None, + descriptors: None, + } +} + +#[cfg_attr(feature = "ci-perf", test)] +#[allow(dead_code)] +fn long_lived_session_benchmark() { + let stress = std::env::var("PET_SESSION_STRESS").is_ok_and(|value| value == "1"); + let sizes = if stress { STRESS_SIZES } else { FAST_SIZES }; + let samples_per_size = if stress { 10 } else { 3 }; + let resolve_concurrency = if stress { 10 } else { 4 }; + let cold_resolve_batches = if stress { 5 } else { 2 }; + let fixture = Fixture::new(); + let mut refresh_scenario_samples = Vec::new(); + let mut cache_usage = Vec::new(); + let mut inventory_resource_samples: Vec<(usize, ResourceSample)> = Vec::new(); + let mut ambient_count_peak = AmbientCountPeak::default(); + + let barrier = fixture.barrier.as_os_str(); + let python_path = fixture.python_path.as_os_str(); + let resolve_root = fixture.resolve_root(); + let resolve_environment = + fixture.create_resolve_environment(&resolve_root.join("cache-control")); + let resolve_cache = fixture._root.path().join("resolve-cache"); + fs::create_dir_all(&resolve_cache).expect("failed to create resolve cache"); + let resolve_environment_variables = [ + ("PET_SESSION_RESOLVE_BARRIER", barrier), + ("PET_SESSION_RESOLVE_ROOT", resolve_root.as_os_str()), + ("PYTHONPATH", python_path), + ]; + + fixture.clear_barrier(); + let mut cold_release = BarrierReleaseGuard::new(&fixture.barrier); + cold_release.release(); + let cold_process = PetJsonRpcClient::spawn_with_environment(&resolve_environment_variables) + .expect("failed to spawn cold cache-control server"); + cold_process + .configure(json!({ + "workspaceDirectories": [&fixture.workspace], + "cacheDirectory": &resolve_cache, + })) + .expect("failed to configure cold cache-control server"); + let (cold_identity, cold_latency_us) = resolve_and_measure(&cold_process, &resolve_environment); + assert_eq!( + fixture.entered_count(), + 1, + "cold cache-control resolve must probe the interpreter once" + ); + assert_eq!(fixture.released_count(), 1); + assert_eq!(fixture.failed_count(), 0); + shutdown_client(&cold_process, "cold cache-control server"); + let cache_after_cold_resolve = directory_usage(&resolve_cache); + assert!( + cache_after_cold_resolve.0 > 0 && cache_after_cold_resolve.1 > 0, + "cold real resolve did not produce a persistent cache entry" + ); + + fixture.clear_barrier(); + let mut warm_release = BarrierReleaseGuard::new(&fixture.barrier); + warm_release.release(); + let warm_process = PetJsonRpcClient::spawn_with_environment(&resolve_environment_variables) + .expect("failed to spawn disk-warm cache-control server"); + warm_process + .configure(json!({ + "workspaceDirectories": [&fixture.workspace], + "cacheDirectory": &resolve_cache, + })) + .expect("failed to configure disk-warm cache-control server"); + let (disk_warm_identity, disk_warm_latency_us) = + resolve_and_measure(&warm_process, &resolve_environment); + assert_eq!( + fixture.entered_count(), + 0, + "disk-warm resolve spawned an interpreter instead of using the persistent cache" + ); + let (same_process_identity, same_process_warm_latency_us) = + resolve_and_measure(&warm_process, &resolve_environment); + assert_eq!( + fixture.entered_count(), + 0, + "same-process warm resolve unexpectedly spawned an interpreter" + ); + assert_eq!(cold_identity, disk_warm_identity); + assert_eq!(cold_identity, same_process_identity); + shutdown_client(&warm_process, "disk-warm cache-control server"); + let persistent_cache_resolve_samples = vec![ + json!({ + "scenario": "cold", + "sampleCount": 1, + "latencyUs": [cold_latency_us], + "interpreterProcessesStarted": 1, + }), + json!({ + "scenario": "diskWarm", + "sampleCount": 1, + "latencyUs": [disk_warm_latency_us], + "interpreterProcessesStarted": 0, + }), + json!({ + "scenario": "sameProcessWarm", + "sampleCount": 1, + "latencyUs": [same_process_warm_latency_us], + "interpreterProcessesStarted": 0, + }), + ]; + + for &size in sizes { + let mut expected = fixture.reset_inventory(size); + expected.sort_unstable(); + + fixture.reset_cache(); + let before_first_process = directory_usage(&fixture.cache); + assert_eq!( + before_first_process, + (0, 0), + "first-process scenario must start with an empty disk cache" + ); + let first_process = + PetJsonRpcClient::spawn().expect("failed to spawn first-process scenario server"); + first_process + .configure(json!({ + "workspaceDirectories": [&fixture.workspace], + "cacheDirectory": &fixture.cache, + })) + .expect("failed to configure first-process scenario server"); + let first = + refresh_and_measure(&first_process, &fixture.workspace, &mut ambient_count_peak); + assert_eq!( + first.inventory, expected, + "first-process refresh changed fixture identities at size {size}" + ); + refresh_scenario_samples.push(json!({ + "scenario": "firstProcessEmptyDiskCache", + "inventorySize": size, + "sampleCount": 1, + "roundTripUs": [first.round_trip_us], + "ttfeUs": [first.ttfe_us], + })); + shutdown_client(&first_process, "first-process scenario server"); + + let after_first_refresh = directory_usage(&fixture.cache); + let new_process = + PetJsonRpcClient::spawn().expect("failed to spawn reused-cache scenario server"); + new_process + .configure(json!({ + "workspaceDirectories": [&fixture.workspace], + "cacheDirectory": &fixture.cache, + })) + .expect("failed to configure reused-cache scenario server"); + let reused = refresh_and_measure(&new_process, &fixture.workspace, &mut ambient_count_peak); + assert_eq!( + reused.inventory, expected, + "new-process refresh changed fixture identities at size {size}" + ); + refresh_scenario_samples.push(json!({ + "scenario": "newProcessAfterFirstRefresh", + "inventorySize": size, + "sampleCount": 1, + "roundTripUs": [reused.round_trip_us], + "ttfeUs": [reused.ttfe_us], + })); + + let mut warm_round_trip_us = Vec::with_capacity(samples_per_size); + let mut warm_ttfe_us = Vec::with_capacity(samples_per_size); + for _ in 0..samples_per_size { + let warm = + refresh_and_measure(&new_process, &fixture.workspace, &mut ambient_count_peak); + assert_eq!( + warm.inventory, expected, + "same-process warm refresh changed fixture identities at size {size}" + ); + warm_round_trip_us.push(warm.round_trip_us); + warm_ttfe_us.push(warm.ttfe_us); + } + refresh_scenario_samples.push(json!({ + "scenario": "sameProcessWarm", + "inventorySize": size, + "sampleCount": samples_per_size, + "roundTripUs": warm_round_trip_us, + "ttfeUs": warm_ttfe_us, + })); + inventory_resource_samples.push((size, observe_resources(new_process.process_id()))); + shutdown_client(&new_process, "reused-cache scenario server"); + + let after_warm_refresh = directory_usage(&fixture.cache); + cache_usage.push(json!({ + "inventorySize": size, + "beforeFirstProcess": cache_json(before_first_process), + "afterFirstRefresh": cache_json(after_first_refresh), + "afterWarmRefresh": cache_json(after_warm_refresh), + "diskCacheAvailableForNewProcess": after_first_refresh.1 > 0, + })); + } + + let client = PetJsonRpcClient::spawn_with_environment(&[ + ("PET_SESSION_RESOLVE_BARRIER", barrier), + ("PET_SESSION_RESOLVE_ROOT", resolve_root.as_os_str()), + ("PYTHONPATH", python_path), + ]) + .expect("failed to spawn long-lived PET server"); + client + .configure(json!({ + "workspaceDirectories": [&fixture.workspace], + "cacheDirectory": &fixture.cache, + })) + .expect("failed to configure long-lived PET server"); + + let mut churn_expected = fixture.reset_inventory(10); + churn_expected.sort_unstable(); + let initial = refresh_and_measure(&client, &fixture.workspace, &mut ambient_count_peak); + assert_eq!(initial.inventory, churn_expected); + + let removed_prefix = fixture.workspace.join("env-0000"); + let removed_index = churn_expected + .iter() + .position(|identity| identity.prefix == norm_case(&removed_prefix)) + .expect("missing identity selected for replacement"); + fs::remove_dir_all(&removed_prefix).expect("failed to delete churn environment"); + churn_expected.remove(removed_index); + churn_expected.push(fixture.create_fake_environment("replacement", "3.12.1")); + churn_expected.sort_unstable(); + let replaced = refresh_and_measure(&client, &fixture.workspace, &mut ambient_count_peak); + assert_eq!(replaced.inventory.len(), initial.inventory.len()); + assert_ne!( + &replaced.inventory, &initial.inventory, + "same-count replacement was not detected" + ); + assert_eq!(replaced.inventory, churn_expected); + + let edited_prefix = fixture.workspace.join("env-0001"); + fs::write( + edited_prefix.join("pyvenv.cfg"), + "version = 3.13.2\nprompt = edited\n", + ) + .expect("failed to edit churn environment"); + write_version_header(&edited_prefix, "3.13.2"); + let previous_index = churn_expected + .iter() + .position(|identity| identity.prefix == norm_case(&edited_prefix)) + .expect("missing identity selected for edit"); + let previous = churn_expected.remove(previous_index); + churn_expected.push(EnvironmentIdentity { + executable: previous.executable, + prefix: previous.prefix, + kind: previous.kind, + name: Some("edited".to_string()), + version: Some("3.13.2".to_string()), + }); + churn_expected.sort_unstable(); + let edited = refresh_and_measure(&client, &fixture.workspace, &mut ambient_count_peak); + assert_eq!(edited.inventory, churn_expected); + + let alias_prefix = fixture.workspace.join("env-0002"); + let original_executable = python_executable(&bin_directory(&alias_prefix), false); + let alias_executable = python_executable(&bin_directory(&alias_prefix), true); + fs::hard_link(&original_executable, &alias_executable) + .expect("failed to create executable alias"); + fs::remove_file(&original_executable).expect("failed to remove original executable alias"); + let previous_index = churn_expected + .iter() + .position(|identity| identity.prefix == norm_case(&alias_prefix)) + .expect("missing identity selected for alias churn"); + let previous = churn_expected.remove(previous_index); + churn_expected.push(EnvironmentIdentity { + executable: norm_case(alias_executable), + #[cfg(windows)] + version: None, + ..previous + }); + churn_expected.sort_unstable(); + let aliased = refresh_and_measure(&client, &fixture.workspace, &mut ambient_count_peak); + assert_eq!(aliased.inventory, churn_expected); + + fixture.clear_barrier(); + let resolve_executables = + fixture.create_resolve_environments(1 + resolve_concurrency * (1 + cold_resolve_batches)); + let mut resolve_executables = resolve_executables.into_iter(); + let warmup_executable = resolve_executables + .next() + .expect("resolve fixture omitted warm-up interpreter"); + let mut warmup_release = BarrierReleaseGuard::new(&fixture.barrier); + let warmup = submit_fixture_resolve(&client, &warmup_executable); + fixture.wait_for_entered(1); + warmup_release.release(); + warmup.wait("warm-up resolve"); + assert_eq!(fixture.released_count(), 1); + assert_eq!(fixture.failed_count(), 0); + + fixture.clear_barrier(); + let pre_resolve_resources = observe_resources(client.process_id()); + let original_cache_before_overlap = + cache_contents(&fixture.cache).expect("failed to capture pre-overlap cache contents"); + assert!( + !original_cache_before_overlap.is_empty(), + "warm-up resolve did not populate the original process cache" + ); + let overlap_workspace = fixture._root.path().join("overlap-workspace"); + let mut overlap_expected = fixture.reset_inventory_at(&overlap_workspace, 3); + overlap_expected.sort_unstable(); + let overlap_executables = resolve_executables + .by_ref() + .take(resolve_concurrency) + .collect::>(); + assert_eq!(overlap_executables.len(), resolve_concurrency); + let mut overlap_release = BarrierReleaseGuard::new(&fixture.barrier); + let overlap_pending = overlap_executables + .iter() + .map(|expected| submit_fixture_resolve(&client, expected)) + .collect::>(); + fixture.wait_for_entered(resolve_concurrency); + assert_eq!( + fixture.entered_count(), + resolve_concurrency, + "every distinct resolve must start before any response is awaited" + ); + let barrier_observed_resources = observe_resources(client.process_id()); + let info = client + .info() + .expect("info request was not responsive while resolves were blocked"); + assert!( + !info.pet_version.is_empty(), + "info returned an empty PET version" + ); + client + .configure(json!({ + "workspaceDirectories": [&overlap_workspace], + "cacheDirectory": &fixture.cache, + })) + .expect("configure was not responsive while resolves were blocked"); + client.clear_notifications(); + if let Err(error) = client.refresh(None) { + panic!( + "refresh failed while cold resolves were held: {error}; entered={}, released={}, barrierReleased={}, stderr={}", + fixture.entered_count(), + fixture.released_count(), + fixture.barrier.join("release").exists(), + client.stderr_output() + ); + } + let (configured_overlap_inventory, overlap_ambient_environment_count) = + fixture_inventory(client.environment_notifications(), &overlap_workspace); + assert_eq!( + configured_overlap_inventory, overlap_expected, + "overlap refresh did not report the newly configured inventory" + ); + let (fixture_manager_count, overlap_ambient_manager_count) = + manager_counts(&client, &overlap_workspace); + assert_eq!( + fixture_manager_count, 0, + "overlap fixture unexpectedly reported an environment manager" + ); + ambient_count_peak.observe( + overlap_ambient_environment_count, + overlap_ambient_manager_count, + ); + let cache_after_overlap = + cache_contents(&fixture.cache).expect("failed to capture post-overlap cache contents"); + assert!( + original_cache_entries_unchanged(&original_cache_before_overlap, &cache_after_overlap), + "overlap reconfiguration modified or deleted an original cache entry" + ); + assert_eq!( + fixture.released_count(), + 0, + "a cold resolve crossed the fixture barrier before its release" + ); + assert!( + !fixture.barrier.join("release").exists(), + "the resolve barrier was released before responsiveness checks completed" + ); + let release_at = Instant::now(); + assert!( + overlap_pending + .iter() + .all(|request| request.submitted_at() < release_at), + "all resolve requests must be submitted before releasing the fixture barrier" + ); + overlap_release.release(); + for request in overlap_pending { + request.wait("barrier-proven concurrent resolve"); + } + assert_eq!( + fixture.released_count(), + resolve_concurrency, + "every overlap interpreter must confirm barrier release" + ); + assert_eq!( + fixture.failed_count(), + 0, + "an overlap interpreter timed out at the barrier" + ); + let post_overlap_resources = observe_resources(client.process_id()); + + let mut resolve_latency_us = Vec::new(); + let mut resolve_batch_resources: Vec<(usize, ResourceSample)> = Vec::new(); + for batch in 0..cold_resolve_batches { + fixture.clear_barrier(); + let mut latency_release = BarrierReleaseGuard::new(&fixture.barrier); + latency_release.release(); + let latency_executables = resolve_executables + .by_ref() + .take(resolve_concurrency) + .collect::>(); + assert_eq!(latency_executables.len(), resolve_concurrency); + let latency_pending = latency_executables + .iter() + .map(|expected| submit_fixture_resolve(&client, expected)) + .collect::>(); + resolve_latency_us.extend(latency_pending.into_iter().map(|request| { + let (_, latency) = request.wait("concurrent resolve"); + latency.as_micros() + })); + assert_eq!( + fixture.entered_count(), + resolve_concurrency, + "latency batch must cold-start every distinct interpreter" + ); + assert_eq!( + fixture.released_count(), + resolve_concurrency, + "latency batch interpreters must confirm barrier release" + ); + assert_eq!( + fixture.failed_count(), + 0, + "a latency batch interpreter timed out at the barrier" + ); + resolve_batch_resources.push((batch, observe_resources(client.process_id()))); + } + assert!( + resolve_executables.next().is_none(), + "resolve fixture count did not match the exercised cold batches" + ); + + let resource_after = observe_resources(client.process_id()); + if let (Some(baseline), Some(after)) = (post_overlap_resources.threads, resource_after.threads) + { + assert!( + after <= baseline, + "PET worker threads grew across repeated cold batches: post-overlap observation {baseline}, post-workload {after}" + ); + } + if let (Some(baseline), Some(after)) = ( + post_overlap_resources.descriptors, + resource_after.descriptors, + ) { + assert!( + after <= baseline, + "PET handles/descriptors grew across repeated cold batches: post-overlap observation {baseline}, post-workload {after}" + ); + } + + let mut all_resource_samples = vec![ + pre_resolve_resources.clone(), + barrier_observed_resources.clone(), + post_overlap_resources.clone(), + resource_after.clone(), + ]; + all_resource_samples.extend( + inventory_resource_samples + .iter() + .map(|(_, sample)| sample.clone()), + ); + all_resource_samples.extend( + resolve_batch_resources + .iter() + .map(|(_, sample)| sample.clone()), + ); + let observed_peak = observed_resource_peak(&all_resource_samples); + shutdown_client(&client, "long-lived benchmark server"); + + let first_process_sample_count = sizes.len(); + let new_process_sample_count = sizes.len(); + let same_process_warm_sample_count = sizes.len() * samples_per_size; + println!( + "SESSION_METRICS {}", + serde_json::to_string(&json!({ + "status": "passed", + "mode": if stress { "stress" } else { "fast" }, + "sizes": sizes, + "samplesPerSize": samples_per_size, + "measurementCounts": { + "firstProcessEmptyDiskCache": first_process_sample_count, + "newProcessAfterFirstRefresh": new_process_sample_count, + "sameProcessWarm": same_process_warm_sample_count, + "persistentCacheResolve": persistent_cache_resolve_samples.len(), + "cacheUsage": cache_usage.len(), + "inventoryResources": inventory_resource_samples.len(), + "resolveLatency": resolve_latency_us.len(), + "resolveBatchResources": resolve_batch_resources.len(), + }, + "refreshScenarioSamples": refresh_scenario_samples, + "cacheUsage": cache_usage, + "persistentCacheResolveSamples": persistent_cache_resolve_samples, + "persistentCacheAfterColdResolve": cache_json(cache_after_cold_resolve), + "resolveConcurrency": resolve_concurrency, + "coldResolveBatches": cold_resolve_batches, + "resolveLatencyUs": resolve_latency_us, + "overlapProcessesStarted": resolve_concurrency, + "overlapAmbientEnvironmentCount": overlap_ambient_environment_count, + "overlapAmbientManagerCount": overlap_ambient_manager_count, + "maxAmbientEnvironmentCount": ambient_count_peak.environment_count, + "maxAmbientManagerCount": ambient_count_peak.manager_count, + "latencyProcessesStarted": resolve_concurrency * cold_resolve_batches, + "inventoryResourceSamples": inventory_resource_samples + .iter() + .map(|(inventory_size, sample)| json!({ + "inventorySize": inventory_size, + "resources": resource_json(sample), + })) + .collect::>(), + "resolveBatchResourceSamples": resolve_batch_resources + .iter() + .map(|(batch, sample)| json!({ + "batch": batch, + "resources": resource_json(sample), + })) + .collect::>(), + "preResolveResources": resource_json(&pre_resolve_resources), + "barrierObservedResources": resource_json(&barrier_observed_resources), + "postOverlapResources": resource_json(&post_overlap_resources), + "observedResourcePeak": resource_json(&observed_peak), + "resourceAfter": resource_json(&resource_after), + "rssDeltaFromPreResolveBytes": i128::from(resource_after.resident_bytes) + - i128::from(pre_resolve_resources.resident_bytes), + })) + .expect("failed to serialize session metrics") + ); +} diff --git a/docs/SESSION_BENCHMARKS.md b/docs/SESSION_BENCHMARKS.md new file mode 100644 index 00000000..a88034f3 --- /dev/null +++ b/docs/SESSION_BENCHMARKS.md @@ -0,0 +1,98 @@ +# Long-lived session benchmarks + +`session_performance` runs deterministic virtual-environment fixtures against +real PET servers. For every inventory size it labels and measures three distinct +file-only refresh scenarios: a first server process with an explicitly empty +cache directory, a new process after that first refresh using the same cache +directory, and repeated warm refreshes in that second process. Process-cold does +not imply an OS-cold filesystem cache. These fake environments are identified +from files and normally produce no persistent resolve-cache entries, so their +new-process samples are not described as disk-cache hits. + +A separate, fixed-size control resolves one real venv. Its cold process must +start exactly one instrumented interpreter and write a nonempty persistent +cache entry before a normal bounded shutdown. A new PET process then resolves +the unchanged interpreter from the same cache without starting an interpreter, +and a same-process warm resolve must do the same. All three resolved identities +must match the submitted fixture's exact executable, prefix, `Venv` kind, and +full runtime `sys.version_info` string. The expected version is queried once +from the Python used to create the copied venvs, outside measured PET +operations. The artifact labels their latencies `cold`, `diskWarm`, and +`sameProcessWarm` and records one sample and the observed probe count for each. + +The benchmark also measures request-relative time to first fixture environment (excluding ambient host results), +concurrent resolve latency, and sampled process-specific resident memory, +threads, and handles or descriptors where the platform exposes them reliably. +Resource values are observed snapshot maxima, not lifetime peaks. Resident +memory is process RSS, not exact retained heap. macOS reports thread and +descriptor counts as unavailable rather than substituting zero. Cache and +resource samples are taken outside request timing. The JSON line prefixed with +`SESSION_METRICS` includes every sample, explicit pass status, and exact +measurement counts. + +The fast workload sweeps 1, 10, and 100 environments with one first-process +sample, one new-process-after-first-refresh sample, and three same-process warm +samples per size: + +```console +cargo test --release --features ci-perf -p pet --test session_performance long_lived_session_benchmark -- --nocapture +``` + +The stress workload sweeps 1, 10, 100, and 1000 environments with one +first-process sample, one new-process-after-first-refresh sample, ten +same-process warm samples per size, and ten overlapping resolves: + +```console +PET_SESSION_STRESS=1 cargo test --release --features ci-perf -p pet --test session_performance long_lived_session_benchmark -- --nocapture +``` + +In PowerShell, set `$env:PET_SESSION_STRESS = "1"` for the stress invocation. +Set `PET_SESSION_PYTHON` to an alternate Python executable when the default +`python`/`python3` cannot create a runnable copied venv (for example, +`PET_SESSION_PYTHON=/usr/bin/python3` in WSL). +`.github/workflows/session-benchmarks.yml` runs the fast workload for pull +requests and the stress workload weekly or through its manual `stress` mode. +The workflow keeps Cargo output in a target directory inside its checkout and +uploads only the privacy-safe `session-metrics.json` artifact, never fixture +paths or raw benchmark output. A dedicated parser rejects missing, malformed, +duplicate, mode-inconsistent, count-inconsistent, or shape-inconsistent +payloads. It also recomputes the observed resource peak from every reported +snapshot and the RSS delta from the pre-resolve and final snapshots, rejecting +either derived value unless it matches exactly. Artifact writes are atomic, and +the timeout fallback replaces corrupt artifacts rather than uploading invalid +JSON. Failed runs retain validated +measurements when available; failures without usable metrics have explicit failed +status and zero counts. The workflow preserves the original benchmark failure. + +Resolve overlap is established in an untimed proof pass by a fixture barrier. +Only environments below the fixture resolve root enter the barrier; resolved +real paths keep macOS temporary-directory aliases equivalent, while ambient +interpreters bypass it. During the proof pass, +every distinct interpreter process records entry before the client issues +`info`, reconfigures to a new workspace while retaining the process's original +cache directory, and refreshes that known inventory. Every cache file captured +before reconfiguration must still exist with identical bytes afterward; cache +files added for unrelated global discoveries are allowed. The barrier remains +held until those responsiveness checks complete. +Platform-global locators may also report host installations and managers. The +benchmark converts only configured workspace entries to strict fixture +identities, validates that fixture-scoped managers remain empty, and records +only counts of unrelated global discoveries and managers, never their paths. +The reported ambient maxima include every timed, churn, and overlap refresh; +artifact validation rejects overlap counts above those maxima. +Two fast or five stress pre-released batches then use fresh, distinct +interpreters for unobstructed client-latency and resource-cycling samples. +Every warm-up, overlap, and latency response is paired with its submitted +fixture and must return that same exact identity and runtime version; merely +returning non-null JSON is not sufficient. +Process creation is therefore intentional cold-resolve work; refresh timings +create no fixture process. Resource snapshots are taken after each inventory +size and resolve batch, outside request timing. The reported observed peak is +the maximum of those actual snapshots, while post-workload thread and +handle/descriptor counts must not exceed the measured post-overlap observation. +The pre-resolve snapshot remains in the metrics so lazy worker-pool growth from +the first concurrent batch is visible rather than treated as a leak or hidden. +An empty disk cache is reported as such and is not treated as evidence of a +retention bound. Every measured server must complete an explicitly successful +stdin-close shutdown before successful metrics are emitted; `Drop` cleanup is +only a failure fallback. diff --git a/scripts/session_metrics.py b/scripts/session_metrics.py new file mode 100644 index 00000000..61db2e41 --- /dev/null +++ b/scripts/session_metrics.py @@ -0,0 +1,522 @@ +#!/usr/bin/env python3 +# Copyright (c) Microsoft Corporation. +# Licensed under the MIT License. +"""Extract and validate privacy-safe long-lived session benchmark metrics.""" + +from __future__ import annotations + +import argparse +import json +import os +import sys +import tempfile +from pathlib import Path +from typing import Any, Sequence + + +METRICS_PREFIX = "SESSION_METRICS " +COUNT_KEYS = { + "firstProcessEmptyDiskCache", + "newProcessAfterFirstRefresh", + "sameProcessWarm", + "persistentCacheResolve", + "cacheUsage", + "inventoryResources", + "resolveLatency", + "resolveBatchResources", +} +SUCCESS_KEYS = { + "status", + "mode", + "sizes", + "samplesPerSize", + "measurementCounts", + "refreshScenarioSamples", + "cacheUsage", + "persistentCacheResolveSamples", + "persistentCacheAfterColdResolve", + "resolveConcurrency", + "coldResolveBatches", + "resolveLatencyUs", + "overlapProcessesStarted", + "overlapAmbientEnvironmentCount", + "overlapAmbientManagerCount", + "maxAmbientEnvironmentCount", + "maxAmbientManagerCount", + "latencyProcessesStarted", + "inventoryResourceSamples", + "resolveBatchResourceSamples", + "preResolveResources", + "barrierObservedResources", + "postOverlapResources", + "observedResourcePeak", + "resourceAfter", + "rssDeltaFromPreResolveBytes", +} +SCENARIOS = { + "firstProcessEmptyDiskCache", + "newProcessAfterFirstRefresh", + "sameProcessWarm", +} +PERSISTENT_CACHE_SCENARIOS = { + "cold": 1, + "diskWarm": 0, + "sameProcessWarm": 0, +} + + +class MetricsError(ValueError): + """Raised when benchmark metrics are missing or inconsistent.""" + + +def require_object(value: Any, name: str, keys: set[str]) -> dict[str, Any]: + if not isinstance(value, dict): + raise MetricsError(f"{name} must be an object") + if set(value) != keys: + missing = sorted(keys - set(value)) + extra = sorted(set(value) - keys) + raise MetricsError(f"{name} has invalid keys; missing={missing}, extra={extra}") + return value + + +def require_integer(value: Any, name: str, *, minimum: int = 0) -> int: + if isinstance(value, bool) or not isinstance(value, int) or value < minimum: + raise MetricsError(f"{name} must be an integer of at least {minimum}") + return value + + +def require_integer_list(value: Any, name: str, expected_length: int) -> list[int]: + if not isinstance(value, list) or len(value) != expected_length: + raise MetricsError(f"{name} must contain exactly {expected_length} samples") + return [require_integer(item, f"{name}[{index}]") for index, item in enumerate(value)] + + +def validate_resource(value: Any, name: str) -> dict[str, Any]: + resource = require_object( + value, name, {"residentBytes", "threads", "handlesOrDescriptors"} + ) + require_integer(resource["residentBytes"], f"{name}.residentBytes") + for key in ("threads", "handlesOrDescriptors"): + if resource[key] is not None: + require_integer(resource[key], f"{name}.{key}") + return resource + + +def observed_resource_peak(samples: Sequence[dict[str, Any]]) -> dict[str, Any]: + def optional_max(key: str) -> int | None: + known = [sample[key] for sample in samples if sample[key] is not None] + return max(known) if known else None + + return { + "residentBytes": max(sample["residentBytes"] for sample in samples), + "threads": optional_max("threads"), + "handlesOrDescriptors": optional_max("handlesOrDescriptors"), + } + + +def validate_cache_value(value: Any, name: str) -> dict[str, Any]: + cache = require_object(value, name, {"files", "bytes"}) + require_integer(cache["files"], f"{name}.files") + require_integer(cache["bytes"], f"{name}.bytes") + return cache + + +def validate_success_metrics(value: Any, expected_mode: str) -> dict[str, Any]: + metrics = require_object(value, "metrics", SUCCESS_KEYS) + if metrics["status"] != "passed": + raise MetricsError("metrics.status must be 'passed'") + if metrics["mode"] != expected_mode: + raise MetricsError( + f"metrics.mode must be {expected_mode!r}, got {metrics['mode']!r}" + ) + + expected_sizes = [1, 10, 100, 1000] if expected_mode == "stress" else [1, 10, 100] + expected_samples = 10 if expected_mode == "stress" else 3 + expected_concurrency = 10 if expected_mode == "stress" else 4 + expected_batches = 5 if expected_mode == "stress" else 2 + require_integer_list(metrics["sizes"], "metrics.sizes", len(expected_sizes)) + require_integer(metrics["samplesPerSize"], "metrics.samplesPerSize") + if metrics["sizes"] != expected_sizes: + raise MetricsError(f"metrics.sizes must equal {expected_sizes}") + if metrics["samplesPerSize"] != expected_samples: + raise MetricsError(f"metrics.samplesPerSize must equal {expected_samples}") + + counts = require_object(metrics["measurementCounts"], "measurementCounts", COUNT_KEYS) + for key in COUNT_KEYS: + require_integer(counts[key], f"measurementCounts.{key}") + expected_counts = { + "firstProcessEmptyDiskCache": len(expected_sizes), + "newProcessAfterFirstRefresh": len(expected_sizes), + "sameProcessWarm": len(expected_sizes) * expected_samples, + "persistentCacheResolve": len(PERSISTENT_CACHE_SCENARIOS), + "cacheUsage": len(expected_sizes), + "inventoryResources": len(expected_sizes), + "resolveLatency": expected_concurrency * expected_batches, + "resolveBatchResources": expected_batches, + } + if counts != expected_counts: + raise MetricsError( + f"measurementCounts must equal {expected_counts}, got {counts}" + ) + + scenarios = metrics["refreshScenarioSamples"] + if not isinstance(scenarios, list) or len(scenarios) != len(expected_sizes) * 3: + raise MetricsError("refreshScenarioSamples must contain three scenarios per size") + seen_scenarios: set[tuple[int, str]] = set() + for index, value in enumerate(scenarios): + sample = require_object( + value, + f"refreshScenarioSamples[{index}]", + {"scenario", "inventorySize", "sampleCount", "roundTripUs", "ttfeUs"}, + ) + scenario = sample["scenario"] + size = require_integer(sample["inventorySize"], "inventorySize", minimum=1) + if not isinstance(scenario, str) or scenario not in SCENARIOS or size not in expected_sizes: + raise MetricsError(f"refreshScenarioSamples[{index}] has an invalid scenario or size") + key = (size, scenario) + if key in seen_scenarios: + raise MetricsError(f"duplicate refresh scenario for size {size}: {scenario}") + seen_scenarios.add(key) + expected_count = expected_samples if scenario == "sameProcessWarm" else 1 + require_integer(sample["sampleCount"], "sampleCount", minimum=1) + if sample["sampleCount"] != expected_count: + raise MetricsError( + f"refreshScenarioSamples[{index}].sampleCount must equal {expected_count}" + ) + round_trips = require_integer_list( + sample["roundTripUs"], + f"refreshScenarioSamples[{index}].roundTripUs", + expected_count, + ) + ttfes = require_integer_list( + sample["ttfeUs"], + f"refreshScenarioSamples[{index}].ttfeUs", + expected_count, + ) + if any(ttfe > round_trip for ttfe, round_trip in zip(ttfes, round_trips)): + raise MetricsError( + f"refreshScenarioSamples[{index}] contains TTFE above round trip" + ) + + cache_usage = metrics["cacheUsage"] + if not isinstance(cache_usage, list) or len(cache_usage) != len(expected_sizes): + raise MetricsError("cacheUsage must contain one entry per size") + seen_cache_sizes = set() + for index, value in enumerate(cache_usage): + cache = require_object( + value, + f"cacheUsage[{index}]", + { + "inventorySize", + "beforeFirstProcess", + "afterFirstRefresh", + "afterWarmRefresh", + "diskCacheAvailableForNewProcess", + }, + ) + size = require_integer(cache["inventorySize"], "inventorySize", minimum=1) + if size not in expected_sizes or size in seen_cache_sizes: + raise MetricsError(f"cacheUsage[{index}] has an invalid or duplicate size") + seen_cache_sizes.add(size) + before = validate_cache_value( + cache["beforeFirstProcess"], f"cacheUsage[{index}].beforeFirstProcess" + ) + after_first = validate_cache_value( + cache["afterFirstRefresh"], f"cacheUsage[{index}].afterFirstRefresh" + ) + validate_cache_value( + cache["afterWarmRefresh"], f"cacheUsage[{index}].afterWarmRefresh" + ) + if before != {"files": 0, "bytes": 0}: + raise MetricsError("first-process cache must be empty") + available = cache["diskCacheAvailableForNewProcess"] + if not isinstance(available, bool) or available != (after_first["bytes"] > 0): + raise MetricsError( + "diskCacheAvailableForNewProcess must exactly reflect nonzero cached bytes" + ) + + cache_control = metrics["persistentCacheResolveSamples"] + if ( + not isinstance(cache_control, list) + or len(cache_control) != len(PERSISTENT_CACHE_SCENARIOS) + ): + raise MetricsError( + "persistentCacheResolveSamples must contain cold, disk-warm, and same-process samples" + ) + seen_cache_control_scenarios = set() + for index, value in enumerate(cache_control): + sample = require_object( + value, + f"persistentCacheResolveSamples[{index}]", + {"scenario", "sampleCount", "latencyUs", "interpreterProcessesStarted"}, + ) + scenario = sample["scenario"] + if ( + not isinstance(scenario, str) + or scenario not in PERSISTENT_CACHE_SCENARIOS + or scenario in seen_cache_control_scenarios + ): + raise MetricsError( + f"persistentCacheResolveSamples[{index}] has an invalid or duplicate scenario" + ) + seen_cache_control_scenarios.add(scenario) + if require_integer(sample["sampleCount"], "sampleCount", minimum=1) != 1: + raise MetricsError( + f"persistentCacheResolveSamples[{index}].sampleCount must equal 1" + ) + require_integer_list( + sample["latencyUs"], + f"persistentCacheResolveSamples[{index}].latencyUs", + 1, + ) + processes = require_integer( + sample["interpreterProcessesStarted"], "interpreterProcessesStarted" + ) + if processes != PERSISTENT_CACHE_SCENARIOS[scenario]: + raise MetricsError( + f"persistentCacheResolveSamples[{index}] has an invalid interpreter process count" + ) + persistent_cache = validate_cache_value( + metrics["persistentCacheAfterColdResolve"], + "persistentCacheAfterColdResolve", + ) + if persistent_cache["files"] == 0 or persistent_cache["bytes"] == 0: + raise MetricsError("cold real resolve must produce a nonempty persistent cache") + + require_integer(metrics["resolveConcurrency"], "resolveConcurrency", minimum=1) + require_integer(metrics["coldResolveBatches"], "coldResolveBatches", minimum=1) + if metrics["resolveConcurrency"] != expected_concurrency: + raise MetricsError(f"resolveConcurrency must equal {expected_concurrency}") + if metrics["coldResolveBatches"] != expected_batches: + raise MetricsError(f"coldResolveBatches must equal {expected_batches}") + require_integer_list( + metrics["resolveLatencyUs"], + "resolveLatencyUs", + expected_concurrency * expected_batches, + ) + for name in ( + "overlapProcessesStarted", + "overlapAmbientEnvironmentCount", + "overlapAmbientManagerCount", + "maxAmbientEnvironmentCount", + "maxAmbientManagerCount", + "latencyProcessesStarted", + ): + require_integer(metrics[name], name) + if metrics["overlapProcessesStarted"] != expected_concurrency: + raise MetricsError("overlapProcessesStarted does not match resolveConcurrency") + if metrics["latencyProcessesStarted"] != expected_concurrency * expected_batches: + raise MetricsError("latencyProcessesStarted does not match resolve batch work") + for overlap_name, maximum_name in ( + ("overlapAmbientEnvironmentCount", "maxAmbientEnvironmentCount"), + ("overlapAmbientManagerCount", "maxAmbientManagerCount"), + ): + if metrics[overlap_name] > metrics[maximum_name]: + raise MetricsError( + f"{overlap_name} must not exceed {maximum_name}; " + f"got {metrics[overlap_name]} and {metrics[maximum_name]}" + ) + + inventory_resources = metrics["inventoryResourceSamples"] + if not isinstance(inventory_resources, list) or len(inventory_resources) != len(expected_sizes): + raise MetricsError("inventoryResourceSamples must contain one entry per size") + seen_resource_sizes = set() + resource_samples = [] + for index, value in enumerate(inventory_resources): + sample = require_object( + value, + f"inventoryResourceSamples[{index}]", + {"inventorySize", "resources"}, + ) + size = require_integer(sample["inventorySize"], "inventorySize", minimum=1) + if size not in expected_sizes or size in seen_resource_sizes: + raise MetricsError( + f"inventoryResourceSamples[{index}] has an invalid or duplicate size" + ) + seen_resource_sizes.add(size) + resource_samples.append( + validate_resource( + sample["resources"], f"inventoryResourceSamples[{index}].resources" + ) + ) + + batch_resources = metrics["resolveBatchResourceSamples"] + if not isinstance(batch_resources, list) or len(batch_resources) != expected_batches: + raise MetricsError("resolveBatchResourceSamples must contain one entry per batch") + seen_batches = set() + for index, value in enumerate(batch_resources): + sample = require_object( + value, + f"resolveBatchResourceSamples[{index}]", + {"batch", "resources"}, + ) + batch = require_integer(sample["batch"], "batch") + if batch not in range(expected_batches) or batch in seen_batches: + raise MetricsError( + f"resolveBatchResourceSamples[{index}] has an invalid or duplicate batch" + ) + seen_batches.add(batch) + resource_samples.append( + validate_resource( + sample["resources"], f"resolveBatchResourceSamples[{index}].resources" + ) + ) + + for name in ( + "preResolveResources", + "barrierObservedResources", + "postOverlapResources", + "resourceAfter", + ): + resource_samples.append(validate_resource(metrics[name], name)) + reported_peak = validate_resource( + metrics["observedResourcePeak"], "observedResourcePeak" + ) + expected_peak = observed_resource_peak(resource_samples) + if reported_peak != expected_peak: + raise MetricsError( + f"observedResourcePeak must equal {expected_peak}, got {reported_peak}" + ) + + reported_delta = require_integer( + metrics["rssDeltaFromPreResolveBytes"], + "rssDeltaFromPreResolveBytes", + minimum=-(2**63), + ) + expected_delta = ( + metrics["resourceAfter"]["residentBytes"] + - metrics["preResolveResources"]["residentBytes"] + ) + if reported_delta != expected_delta: + raise MetricsError( + "rssDeltaFromPreResolveBytes must equal " + f"{expected_delta}, got {reported_delta}" + ) + return metrics + + +def failed_metrics(mode: str, reason: str) -> dict[str, Any]: + return { + "status": "failed", + "mode": mode, + "metricsProduced": False, + "failureReason": reason, + "measurementCounts": {key: 0 for key in sorted(COUNT_KEYS)}, + } + + +def write_json(path: Path, value: dict[str, Any]) -> None: + with tempfile.NamedTemporaryFile( + mode="w", encoding="utf-8", dir=path.parent, prefix=f".{path.name}.", delete=False + ) as temporary: + temporary_path = Path(temporary.name) + try: + temporary.write(json.dumps(value, indent=2, sort_keys=True) + "\n") + except (OSError, TypeError, ValueError): + temporary.close() + temporary_path.unlink(missing_ok=True) + raise + try: + os.replace(temporary_path, path) + finally: + temporary_path.unlink(missing_ok=True) + + +def extract_metrics( + input_path: Path, + output_path: Path, + benchmark_status: int, + expected_mode: str, +) -> dict[str, Any]: + try: + payloads = [ + line.split(METRICS_PREFIX, 1)[1] + for line in input_path.read_text(encoding="utf-8").splitlines() + if METRICS_PREFIX in line + ] + if len(payloads) != 1: + raise MetricsError( + f"expected exactly one SESSION_METRICS payload, found {len(payloads)}" + ) + try: + decoded = json.loads(payloads[0]) + except json.JSONDecodeError as error: + raise MetricsError("SESSION_METRICS payload is not valid JSON") from error + metrics = validate_success_metrics(decoded, expected_mode).copy() + except (MetricsError, OSError, UnicodeError) as error: + write_json(output_path, failed_metrics(expected_mode, str(error))) + raise MetricsError(str(error)) from error + + metrics["status"] = "passed" if benchmark_status == 0 else "failed" + metrics["metricsProduced"] = True + write_json(output_path, metrics) + return metrics + + +def validate_artifact(value: Any, mode: str) -> None: + if not isinstance(value, dict) or not isinstance(value.get("metricsProduced"), bool): + raise MetricsError("artifact must declare whether metrics were produced") + if value["metricsProduced"]: + if value.get("status") not in ("passed", "failed"): + raise MetricsError("artifact status must be passed or failed") + metrics = value.copy() + del metrics["metricsProduced"] + metrics["status"] = "passed" + validate_success_metrics(metrics, mode) + else: + require_object(value, "failed artifact", { + "status", "mode", "metricsProduced", "failureReason", "measurementCounts" + }) + if value["status"] != "failed" or value["mode"] != mode: + raise MetricsError("failed artifact has invalid status or mode") + if not isinstance(value["failureReason"], str) or not value["failureReason"]: + raise MetricsError("failed artifact must include a failure reason") + counts = require_object(value["measurementCounts"], "measurementCounts", COUNT_KEYS) + for key, count in counts.items(): + if require_integer(count, f"measurementCounts.{key}") != 0: + raise MetricsError("failed artifact without metrics must have zero counts") + + +def ensure_failure_metrics(output_path: Path, mode: str, reason: str) -> None: + if output_path.exists(): + try: + validate_artifact(json.loads(output_path.read_text(encoding="utf-8")), mode) + return + except (MetricsError, OSError, UnicodeError, json.JSONDecodeError) as error: + print(f"Replacing invalid session metrics artifact: {error}", file=sys.stderr) + write_json(output_path, failed_metrics(mode, reason)) + + +def main(argv: Sequence[str] | None = None) -> int: + parser = argparse.ArgumentParser() + subparsers = parser.add_subparsers(dest="command", required=True) + + extract = subparsers.add_parser("extract") + extract.add_argument("--input", type=Path, required=True) + extract.add_argument("--output", type=Path, required=True) + extract.add_argument("--benchmark-status", type=int, required=True) + extract.add_argument("--mode", choices=("fast", "stress"), required=True) + + ensure = subparsers.add_parser("ensure-failure") + ensure.add_argument("--output", type=Path, required=True) + ensure.add_argument("--mode", choices=("fast", "stress"), required=True) + ensure.add_argument("--reason", required=True) + + args = parser.parse_args(argv) + if args.command == "ensure-failure": + ensure_failure_metrics(args.output, args.mode, args.reason) + return 0 + try: + extract_metrics( + args.input, + args.output, + args.benchmark_status, + args.mode, + ) + except MetricsError as error: + parser.exit(1, f"session metrics validation failed: {error}\n") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/scripts/tests/test_session_metrics.py b/scripts/tests/test_session_metrics.py new file mode 100644 index 00000000..6f3ece2f --- /dev/null +++ b/scripts/tests/test_session_metrics.py @@ -0,0 +1,401 @@ +# Copyright (c) Microsoft Corporation. +# Licensed under the MIT License. + +import copy +import json +import os +import runpy +import sys +import tempfile +import unittest +from pathlib import Path +from unittest.mock import patch + + +ROOT = Path(__file__).resolve().parents[2] +sys.path.insert(0, str(ROOT / "scripts")) + +from session_metrics import (MetricsError, ensure_failure_metrics, extract_metrics, + failed_metrics, write_json) # noqa: E402 + + +def resource(resident=100, threads=2, descriptors=3): + return { + "residentBytes": resident, + "threads": threads, + "handlesOrDescriptors": descriptors, + } + + +def valid_metrics(): + sizes = [1, 10, 100] + scenarios = [] + for size in sizes: + for scenario, count in ( + ("firstProcessEmptyDiskCache", 1), + ("newProcessAfterFirstRefresh", 1), + ("sameProcessWarm", 3), + ): + scenarios.append( + { + "scenario": scenario, + "inventorySize": size, + "sampleCount": count, + "roundTripUs": [20] * count, + "ttfeUs": [10] * count, + } + ) + return { + "status": "passed", + "mode": "fast", + "sizes": sizes, + "samplesPerSize": 3, + "measurementCounts": { + "firstProcessEmptyDiskCache": 3, + "newProcessAfterFirstRefresh": 3, + "sameProcessWarm": 9, + "persistentCacheResolve": 3, + "cacheUsage": 3, + "inventoryResources": 3, + "resolveLatency": 8, + "resolveBatchResources": 2, + }, + "refreshScenarioSamples": scenarios, + "cacheUsage": [ + { + "inventorySize": size, + "beforeFirstProcess": {"files": 0, "bytes": 0}, + "afterFirstRefresh": {"files": 0, "bytes": 0}, + "afterWarmRefresh": {"files": 0, "bytes": 0}, + "diskCacheAvailableForNewProcess": False, + } + for size in sizes + ], + "persistentCacheResolveSamples": [ + { + "scenario": scenario, + "sampleCount": 1, + "latencyUs": [50], + "interpreterProcessesStarted": processes, + } + for scenario, processes in ( + ("cold", 1), + ("diskWarm", 0), + ("sameProcessWarm", 0), + ) + ], + "persistentCacheAfterColdResolve": {"files": 1, "bytes": 100}, + "resolveConcurrency": 4, + "coldResolveBatches": 2, + "resolveLatencyUs": [50] * 8, + "overlapProcessesStarted": 4, + "overlapAmbientEnvironmentCount": 1, + "overlapAmbientManagerCount": 1, + "maxAmbientEnvironmentCount": 1, + "maxAmbientManagerCount": 1, + "latencyProcessesStarted": 8, + "inventoryResourceSamples": [ + {"inventorySize": size, "resources": resource()} for size in sizes + ], + "resolveBatchResourceSamples": [ + {"batch": batch, "resources": resource()} for batch in range(2) + ], + "preResolveResources": resource(resident=100), + "barrierObservedResources": resource(), + "postOverlapResources": resource(), + "observedResourcePeak": resource(), + "resourceAfter": resource(resident=90), + "rssDeltaFromPreResolveBytes": -10, + } + + +class SessionBarrierTests(unittest.TestCase): + def setUp(self): + self.temp = tempfile.TemporaryDirectory() + self.addCleanup(self.temp.cleanup) + self.root = Path(self.temp.name) + self.resolve_root = self.root / "resolve-environments" + self.prefix = self.resolve_root / "venv" + self.prefix.mkdir(parents=True) + self.barrier = self.root / "barrier" + self.barrier.mkdir() + (self.barrier / "release").write_text("release", encoding="ascii") + + def run_fixture(self, prefix, resolve_root): + with patch.dict(os.environ, { + "PET_SESSION_RESOLVE_BARRIER": str(self.barrier), + "PET_SESSION_RESOLVE_ROOT": str(resolve_root), + }), patch.object(sys, "prefix", str(prefix)): + runpy.run_path(str(ROOT / "crates/pet/tests/fixtures/session_sitecustomize.py")) + return (self.barrier / f"entered-{os.getpid()}").exists() + + def test_barrier_accepts_only_fixture_environment_prefixes(self): + self.assertFalse(self.run_fixture(self.root / "ambient", self.resolve_root)) + self.assertFalse(self.run_fixture(self.root / "resolve-environments-other", self.resolve_root)) + self.assertTrue(self.run_fixture(self.prefix, self.resolve_root)) + self.assertTrue((self.barrier / f"released-{os.getpid()}").is_file()) + + @unittest.skipIf(os.name == "nt", "creating directory symlinks requires Windows privileges") + def test_barrier_matches_real_and_symlinked_temp_directory_spellings(self): + alias = self.root / "alias" + alias.symlink_to(self.resolve_root, target_is_directory=True) + self.assertTrue(self.run_fixture(self.prefix, alias)) + (self.barrier / f"entered-{os.getpid()}").unlink() + self.assertTrue(self.run_fixture(alias / "venv", self.resolve_root)) + + +class SessionMetricsTests(unittest.TestCase): + def setUp(self): + self.temp = tempfile.TemporaryDirectory() + self.root = Path(self.temp.name) + self.input = self.root / "output.txt" + self.output = self.root / "metrics.json" + + def tearDown(self): + self.temp.cleanup() + + def write_payloads(self, *payloads): + self.input.write_text( + "\n".join(f"test benchmark ... SESSION_METRICS {payload}" for payload in payloads), + encoding="utf-8", + ) + + def artifact(self): + return json.loads(self.output.read_text(encoding="utf-8")) + + def assert_metrics_rejected(self, metrics, message): + self.write_payloads(json.dumps(metrics)) + with self.assertRaisesRegex(MetricsError, message): + extract_metrics(self.input, self.output, 0, "fast") + self.assertEqual(self.artifact()["status"], "failed") + self.assertFalse(self.artifact()["metricsProduced"]) + + def test_no_metrics_fails_closed_and_persists_failure(self): + self.input.write_text("benchmark output only\n", encoding="utf-8") + with self.assertRaisesRegex(MetricsError, "exactly one"): + extract_metrics(self.input, self.output, 0, "fast") + self.assertEqual(self.artifact()["status"], "failed") + self.assertFalse(self.artifact()["metricsProduced"]) + + def test_malformed_metrics_fails_closed(self): + self.write_payloads("{not-json") + with self.assertRaisesRegex(MetricsError, "not valid JSON"): + extract_metrics(self.input, self.output, 0, "fast") + self.assertEqual(self.artifact()["status"], "failed") + + def test_duplicate_metrics_fails_closed(self): + payload = json.dumps(valid_metrics()) + self.write_payloads(payload, payload) + with self.assertRaisesRegex(MetricsError, "found 2"): + extract_metrics(self.input, self.output, 0, "fast") + + def test_count_mismatch_fails_closed(self): + metrics = valid_metrics() + metrics["measurementCounts"]["sameProcessWarm"] = 8 + self.write_payloads(json.dumps(metrics)) + with self.assertRaisesRegex(MetricsError, "measurementCounts"): + extract_metrics(self.input, self.output, 0, "fast") + + def test_persistent_cache_control_requires_cache_hit_proof(self): + for field, invalid in ( + (("persistentCacheAfterColdResolve", "bytes"), 0), + (("persistentCacheResolveSamples", 1, "interpreterProcessesStarted"), 1), + (("persistentCacheResolveSamples", 0, "interpreterProcessesStarted"), 0), + ): + with self.subTest(field=field): + metrics = copy.deepcopy(valid_metrics()) + container = metrics + for key in field[:-1]: + container = container[key] + container[field[-1]] = invalid + self.write_payloads(json.dumps(metrics)) + with self.assertRaises(MetricsError): + extract_metrics(self.input, self.output, 0, "fast") + + def test_inconsistent_mode_fails_closed(self): + self.write_payloads(json.dumps(valid_metrics())) + with self.assertRaisesRegex(MetricsError, "metrics.mode"): + extract_metrics(self.input, self.output, 0, "stress") + + def test_failed_cargo_with_valid_metrics_persists_real_failed_metrics(self): + self.write_payloads(json.dumps(valid_metrics())) + metrics = extract_metrics(self.input, self.output, 101, "fast") + self.assertEqual(metrics["status"], "failed") + self.assertTrue(metrics["metricsProduced"]) + self.assertEqual(metrics["measurementCounts"]["sameProcessWarm"], 9) + + def test_success_preserves_validated_metrics(self): + self.write_payloads(json.dumps(valid_metrics())) + metrics = extract_metrics(self.input, self.output, 0, "fast") + self.assertEqual(metrics["status"], "passed") + self.assertTrue(metrics["metricsProduced"]) + self.assertEqual(len(metrics["refreshScenarioSamples"]), 9) + + def test_ambient_overlap_counts_allow_equal_or_larger_maxima(self): + for maximum_offset in (0, 1): + with self.subTest(maximum_offset=maximum_offset): + metrics = valid_metrics() + metrics["maxAmbientEnvironmentCount"] += maximum_offset + metrics["maxAmbientManagerCount"] += maximum_offset + self.write_payloads(json.dumps(metrics)) + extract_metrics(self.input, self.output, 0, "fast") + + def test_ambient_overlap_counts_reject_smaller_maxima(self): + for overlap_name, maximum_name in ( + ("overlapAmbientEnvironmentCount", "maxAmbientEnvironmentCount"), + ("overlapAmbientManagerCount", "maxAmbientManagerCount"), + ): + with self.subTest(overlap_name=overlap_name): + metrics = valid_metrics() + metrics[overlap_name] = 2 + metrics[maximum_name] = 1 + self.assert_metrics_rejected( + metrics, f"{overlap_name} must not exceed {maximum_name}" + ) + + def test_resource_peak_rejects_low_and_high_values_for_every_field(self): + for field in ("residentBytes", "threads", "handlesOrDescriptors"): + for difference in (-1, 1): + with self.subTest(field=field, difference=difference): + metrics = valid_metrics() + metrics["observedResourcePeak"][field] += difference + self.assert_metrics_rejected(metrics, "observedResourcePeak") + + def test_resource_peak_includes_every_sample_category(self): + categories = ( + ("inventory", lambda metrics: metrics["inventoryResourceSamples"][0]["resources"]), + ("pre-resolve", lambda metrics: metrics["preResolveResources"]), + ("barrier", lambda metrics: metrics["barrierObservedResources"]), + ("post-overlap", lambda metrics: metrics["postOverlapResources"]), + ("resolve-batch", lambda metrics: metrics["resolveBatchResourceSamples"][0]["resources"]), + ("after", lambda metrics: metrics["resourceAfter"]), + ) + for name, select_sample in categories: + with self.subTest(category=name): + metrics = valid_metrics() + sample = select_sample(metrics) + sample.update( + residentBytes=1000, + threads=1001, + handlesOrDescriptors=1002, + ) + metrics["observedResourcePeak"] = resource(1000, 1001, 1002) + metrics["rssDeltaFromPreResolveBytes"] = ( + metrics["resourceAfter"]["residentBytes"] + - metrics["preResolveResources"]["residentBytes"] + ) + self.write_payloads(json.dumps(metrics)) + extract_metrics(self.input, self.output, 0, "fast") + + def test_resource_peak_optional_fields_support_all_null_and_mixed_samples(self): + metrics = valid_metrics() + samples = [ + *(entry["resources"] for entry in metrics["inventoryResourceSamples"]), + metrics["preResolveResources"], + metrics["barrierObservedResources"], + metrics["postOverlapResources"], + *(entry["resources"] for entry in metrics["resolveBatchResourceSamples"]), + metrics["resourceAfter"], + ] + for sample in samples: + sample["threads"] = None + sample["handlesOrDescriptors"] = None + metrics["observedResourcePeak"]["threads"] = None + metrics["observedResourcePeak"]["handlesOrDescriptors"] = None + self.write_payloads(json.dumps(metrics)) + extract_metrics(self.input, self.output, 0, "fast") + + samples[0]["threads"] = 7 + samples[-1]["threads"] = 5 + samples[1]["handlesOrDescriptors"] = 8 + samples[-2]["handlesOrDescriptors"] = 6 + metrics["observedResourcePeak"]["threads"] = 7 + metrics["observedResourcePeak"]["handlesOrDescriptors"] = 8 + self.write_payloads(json.dumps(metrics)) + extract_metrics(self.input, self.output, 0, "fast") + + def test_rss_delta_accepts_positive_negative_and_zero_values(self): + for after, expected_delta in ((110, 10), (90, -10), (100, 0)): + with self.subTest(expected_delta=expected_delta): + metrics = valid_metrics() + metrics["resourceAfter"]["residentBytes"] = after + metrics["rssDeltaFromPreResolveBytes"] = expected_delta + if after > metrics["observedResourcePeak"]["residentBytes"]: + metrics["observedResourcePeak"]["residentBytes"] = after + self.write_payloads(json.dumps(metrics)) + extract_metrics(self.input, self.output, 0, "fast") + + def test_inconsistent_rss_delta_fails_closed_with_failure_artifact(self): + metrics = valid_metrics() + metrics["rssDeltaFromPreResolveBytes"] = -9 + self.assert_metrics_rejected(metrics, "rssDeltaFromPreResolveBytes") + + def test_malformed_nested_types_fail_closed_with_failure_artifacts(self): + fields = [ + ("sizes", 0), ("samplesPerSize",), ("resolveConcurrency",), + ("coldResolveBatches",), ("refreshScenarioSamples", 0, "inventorySize"), + ("refreshScenarioSamples", 0, "sampleCount"), + ("refreshScenarioSamples", 0, "scenario"), ("cacheUsage", 0, "inventorySize"), + ("persistentCacheResolveSamples", 0, "sampleCount"), + ("persistentCacheResolveSamples", 0, "scenario"), + ("persistentCacheAfterColdResolve", "files"), + ("inventoryResourceSamples", 0, "inventorySize"), + ("resolveBatchResourceSamples", 0, "batch"), + ] + for field in fields: + for invalid in [True, False, 1.0, 3.0, 4.0, [], {}]: + with self.subTest(field=field, invalid=invalid): + metrics = copy.deepcopy(valid_metrics()) + container = metrics + for key in field[:-1]: + container = container[key] + container[field[-1]] = invalid + self.write_payloads(json.dumps(metrics)) + with self.assertRaises(MetricsError): + extract_metrics(self.input, self.output, 0, "fast") + self.assertEqual(self.artifact()["status"], "failed") + self.assertFalse(self.artifact()["metricsProduced"]) + + def test_timeout_fallback_replaces_corrupt_artifacts(self): + for text in ['{', 'null', '{}', json.dumps(valid_metrics())]: + with self.subTest(text=text): + self.output.write_text(text, encoding="utf-8") + ensure_failure_metrics(self.output, "fast", "interrupted") + self.assertEqual(self.artifact(), failed_metrics("fast", "interrupted")) + self.output.unlink() + ensure_failure_metrics(self.output, "fast", "interrupted") + self.assertEqual(self.artifact(), failed_metrics("fast", "interrupted")) + + def test_timeout_fallback_preserves_valid_success_and_failure_artifacts(self): + for benchmark_status in [0, 101]: + with self.subTest(benchmark_status=benchmark_status): + self.write_payloads(json.dumps(valid_metrics())) + extract_metrics(self.input, self.output, benchmark_status, "fast") + before = self.output.read_bytes() + ensure_failure_metrics(self.output, "fast", "interrupted") + self.assertEqual(self.output.read_bytes(), before) + failure = failed_metrics("fast", "original failure") + write_json(self.output, failure) + ensure_failure_metrics(self.output, "fast", "interrupted") + self.assertEqual(self.artifact(), failure) + + def test_atomic_write_preserves_previous_artifact_if_replace_fails(self): + write_json(self.output, failed_metrics("fast", "original failure")) + before = self.output.read_bytes() + with patch("session_metrics.os.replace", side_effect=OSError("replace failed")): + with self.assertRaisesRegex(OSError, "replace failed"): + write_json(self.output, failed_metrics("fast", "new failure")) + self.assertEqual(self.output.read_bytes(), before) + self.assertEqual(list(self.root.glob(".metrics.json.*")), []) + + def test_workflow_uses_validator_and_timeout_fallback(self): + workflow = (ROOT / ".github/workflows/session-benchmarks.yml").read_text( + encoding="utf-8" + ) + self.assertIn("python -B scripts/session_metrics.py extract", workflow) + self.assertIn("python -B scripts/session_metrics.py ensure-failure", workflow) + self.assertIn("if: always()", workflow) + + +if __name__ == "__main__": + unittest.main()