From 8e516454db72d273e29e0e3f6265485da320d25d Mon Sep 17 00:00:00 2001 From: Karthik Nadig Date: Wed, 30 Sep 2026 09:49:06 -0700 Subject: [PATCH 1/4] test: benchmark long-lived sessions and inventory scale (Fixes #533) Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- .github/workflows/session-benchmarks.yml | 108 ++ .../tests/fixtures/session_sitecustomize.py | 28 + crates/pet/tests/jsonrpc_client.rs | 196 +++- crates/pet/tests/jsonrpc_server_test.rs | 67 ++ crates/pet/tests/session_performance.rs | 1019 +++++++++++++++++ docs/SESSION_BENCHMARKS.md | 73 ++ scripts/session_metrics.py | 413 +++++++ scripts/tests/test_session_metrics.py | 224 ++++ 8 files changed, 2105 insertions(+), 23 deletions(-) create mode 100644 .github/workflows/session-benchmarks.yml create mode 100644 crates/pet/tests/fixtures/session_sitecustomize.py create mode 100644 crates/pet/tests/session_performance.rs create mode 100644 docs/SESSION_BENCHMARKS.md create mode 100644 scripts/session_metrics.py create mode 100644 scripts/tests/test_session_metrics.py 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..e7c4c980 --- /dev/null +++ b/crates/pet/tests/fixtures/session_sitecustomize.py @@ -0,0 +1,28 @@ +# 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: + executable = os.path.normcase(os.path.abspath(sys.executable)) + root_executable_prefix = os.path.normcase(os.path.abspath(resolve_root)) + os.sep + if not executable.startswith(root_executable_prefix): + 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..73b23e14 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,13 +514,68 @@ 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, + }); } }) } +#[cfg(test)] +pub(crate) fn controlled_pending_request( + client: &PetJsonRpcClient, + submitted_at: Instant, +) -> (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) +} + +#[cfg(test)] +pub(crate) fn pending_request_registered( + client: &PetJsonRpcClient, + request: &PendingRequest, +) -> bool { + client + .inner + .state + .pending + .lock() + .unwrap() + .contains_key(&request.id) +} + +#[cfg(test)] +pub(crate) fn pending_request_count(client: &PetJsonRpcClient) -> usize { + client.inner.state.pending.lock().unwrap().len() +} + fn spawn_stderr_reader( stderr: impl Read + Send + 'static, state: Arc, diff --git a/crates/pet/tests/jsonrpc_server_test.rs b/crates/pet/tests/jsonrpc_server_test.rs index 7123928c..1d1f05e1 100644 --- a/crates/pet/tests/jsonrpc_server_test.rs +++ b/crates/pet/tests/jsonrpc_server_test.rs @@ -17,6 +17,73 @@ mod jsonrpc_client; use jsonrpc_client::{EnvironmentNotification, PetJsonRpcClient}; +#[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) = + jsonrpc_client::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!(jsonrpc_client::pending_request_count(&client), 0); + + let consumed_after_deadline_at = Instant::now() - Duration::from_secs(31); + let (consumed_after_deadline, send_timely) = + jsonrpc_client::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) = + jsonrpc_client::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!(jsonrpc_client::pending_request_count(&client), 0); + + let expired_at = Instant::now() - Duration::from_secs(31); + let (expired, _keep_sender_connected) = + jsonrpc_client::controlled_pending_request(&client, expired_at); + assert!(jsonrpc_client::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!(jsonrpc_client::pending_request_count(&client), 0); + + let (dropped, _keep_sender_connected) = + jsonrpc_client::controlled_pending_request(&client, Instant::now()); + assert!(jsonrpc_client::pending_request_registered( + &client, &dropped + )); + drop(dropped); + assert_eq!( + jsonrpc_client::pending_request_count(&client), + 0, + "dropping a pending request must remove its map entry" + ); +} + fn frame_with_headers(payload: &[u8], headers: &[(&str, &str)], line_ending: &[u8]) -> Vec { let mut frame = Vec::new(); for (name, value) in headers { diff --git a/crates/pet/tests/session_performance.rs b/crates/pet/tests/session_performance.rs new file mode 100644 index 00000000..8f28631b --- /dev/null +++ b/crates/pet/tests/session_performance.rs @@ -0,0 +1,1019 @@ +// Copyright (c) Microsoft Corporation. +// Licensed under the MIT License. + +use pet_fs::path::norm_case; +use serde_json::json; +use std::fs; +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, 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, +} + +struct Fixture { + _root: TempDir, + workspace: PathBuf, + cache: PathBuf, + barrier: PathBuf, + python_path: PathBuf, +} + +impl Fixture { + fn new() -> Self { + 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, + } + } + + 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 { + let python = std::env::var_os("PET_SESSION_PYTHON").unwrap_or_else(|| { + if cfg!(windows) { + "python".into() + } else { + "python3".into() + } + }); + (0..count) + .map(|index| { + let prefix = self.resolve_root().join(format!("resolve-{index}")); + 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 {index}: {}", + String::from_utf8_lossy(&output.stderr) + ); + python_executable(&bin_directory(&prefix), false) + }) + .collect() + } + + 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() + ); + } +} + +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) -> 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" + ); + RefreshMeasurement { + inventory, + round_trip_us: timing.round_trip.as_micros(), + ttfe_us: ttfe.as_micros(), + ambient_environment_count, + ambient_manager_count, + } +} + +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 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 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") + .map(|entry| entry.expect("failed to read process descriptor entry")) + .count(), + ), + } +} + +#[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 max_ambient_environment_count = 0; + let mut max_ambient_manager_count = 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); + assert_eq!( + first.inventory, expected, + "first-process refresh changed fixture identities at size {size}" + ); + max_ambient_environment_count = + max_ambient_environment_count.max(first.ambient_environment_count); + max_ambient_manager_count = max_ambient_manager_count.max(first.ambient_manager_count); + 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); + assert_eq!( + reused.inventory, expected, + "new-process refresh changed fixture identities at size {size}" + ); + max_ambient_environment_count = + max_ambient_environment_count.max(reused.ambient_environment_count); + max_ambient_manager_count = max_ambient_manager_count.max(reused.ambient_manager_count); + 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); + assert_eq!( + warm.inventory, expected, + "same-process warm refresh changed fixture identities at size {size}" + ); + max_ambient_environment_count = + max_ambient_environment_count.max(warm.ambient_environment_count); + max_ambient_manager_count = max_ambient_manager_count.max(warm.ambient_manager_count); + 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 barrier = fixture.barrier.as_os_str(); + let python_path = fixture.python_path.as_os_str(); + let resolve_root = fixture.resolve_root(); + 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); + assert_eq!(initial.inventory, churn_expected); + max_ambient_environment_count = + max_ambient_environment_count.max(initial.ambient_environment_count); + max_ambient_manager_count = max_ambient_manager_count.max(initial.ambient_manager_count); + + 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); + 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); + 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); + 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 = client + .submit_resolve( + warmup_executable + .to_str() + .expect("warm-up resolve fixture path was not UTF-8"), + ) + .expect("failed to submit warm-up resolve"); + fixture.wait_for_entered(1); + warmup_release.release(); + warmup + .wait(REQUEST_TIMEOUT) + .expect("warm-up resolve failed"); + 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 overlap_workspace = fixture._root.path().join("overlap-workspace"); + let overlap_cache = fixture._root.path().join("overlap-cache"); + fs::create_dir_all(&overlap_cache).expect("failed to create overlap cache"); + 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(|executable| { + client + .submit_resolve( + executable + .to_str() + .expect("resolve fixture path was not UTF-8"), + ) + .expect("failed to submit concurrent resolve") + }) + .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": &overlap_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" + ); + max_ambient_environment_count = + max_ambient_environment_count.max(overlap_ambient_environment_count); + max_ambient_manager_count = max_ambient_manager_count.max(overlap_ambient_manager_count); + 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 { + let (result, _) = request + .wait(REQUEST_TIMEOUT) + .expect("barrier-proven concurrent resolve failed"); + assert!(!result.is_null(), "resolve returned no environment"); + } + 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(|executable| { + client + .submit_resolve( + executable + .to_str() + .expect("resolve fixture path was not UTF-8"), + ) + .expect("failed to submit resolve latency request") + }) + .collect::>(); + resolve_latency_us.extend(latency_pending.into_iter().map(|request| { + let (result, latency) = request + .wait(REQUEST_TIMEOUT) + .expect("concurrent resolve failed"); + assert!(!result.is_null(), "resolve returned no environment"); + 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, + "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, + "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": max_ambient_environment_count, + "maxAmbientManagerCount": max_ambient_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..fb76ea27 --- /dev/null +++ b/docs/SESSION_BENCHMARKS.md @@ -0,0 +1,73 @@ +# 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 +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. The new-process scenario records whether +the first refresh actually wrote nonzero cache bytes; a zero-byte directory is +not reported as a disk-cache hit. + +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. 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: +every distinct interpreter process records entry before the client issues +`info`, reconfigures to a new workspace and cache, and refreshes that known +inventory. 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. +Two fast or five stress pre-released batches then use fresh, distinct +interpreters for unobstructed client-latency and resource-cycling samples. +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..2cb1b4a9 --- /dev/null +++ b/scripts/session_metrics.py @@ -0,0 +1,413 @@ +#!/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", + "cacheUsage", + "inventoryResources", + "resolveLatency", + "resolveBatchResources", +} +SUCCESS_KEYS = { + "status", + "mode", + "sizes", + "samplesPerSize", + "measurementCounts", + "refreshScenarioSamples", + "cacheUsage", + "resolveConcurrency", + "coldResolveBatches", + "resolveLatencyUs", + "overlapProcessesStarted", + "overlapAmbientEnvironmentCount", + "overlapAmbientManagerCount", + "maxAmbientEnvironmentCount", + "maxAmbientManagerCount", + "latencyProcessesStarted", + "inventoryResourceSamples", + "resolveBatchResourceSamples", + "preResolveResources", + "barrierObservedResources", + "postOverlapResources", + "observedResourcePeak", + "resourceAfter", + "rssDeltaFromPreResolveBytes", +} +SCENARIOS = { + "firstProcessEmptyDiskCache", + "newProcessAfterFirstRefresh", + "sameProcessWarm", +} + + +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) -> None: + 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}") + + +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, + "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" + ) + + 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") + + 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() + 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) + 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) + validate_resource(sample["resources"], f"resolveBatchResourceSamples[{index}].resources") + + for name in ( + "preResolveResources", + "barrierObservedResources", + "postOverlapResources", + "observedResourcePeak", + "resourceAfter", + ): + validate_resource(metrics[name], name) + require_integer(metrics["rssDeltaFromPreResolveBytes"], "rssDeltaFromPreResolveBytes", minimum=-(2**63)) + 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..638c155d --- /dev/null +++ b/scripts/tests/test_session_metrics.py @@ -0,0 +1,224 @@ +# Copyright (c) Microsoft Corporation. +# Licensed under the MIT License. + +import copy +import json +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(): + return { + "residentBytes": 100, + "threads": 2, + "handlesOrDescriptors": 3, + } + + +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, + "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 + ], + "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(), + "barrierObservedResources": resource(), + "postOverlapResources": resource(), + "observedResourcePeak": resource(), + "resourceAfter": resource(), + "rssDeltaFromPreResolveBytes": -10, + } + + +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 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_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_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"), + ("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() From c5155fd921b65622939348456704de04b722500a Mon Sep 17 00:00:00 2001 From: Karthik Nadig Date: Wed, 30 Sep 2026 10:25:54 -0700 Subject: [PATCH 2/4] test: prove persistent cache reuse in session benchmarks Address PR #563 review and CI: isolate test helpers, scope interpreter barriers by real prefixes, and prove cold/disk-warm/same-process cache behavior without ignored reconfiguration. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- .../tests/fixtures/session_sitecustomize.py | 8 +- crates/pet/tests/jsonrpc_client.rs | 161 ++++++++++----- crates/pet/tests/jsonrpc_server_test.rs | 69 +------ crates/pet/tests/session_performance.rs | 186 +++++++++++++++--- docs/SESSION_BENCHMARKS.md | 28 ++- scripts/session_metrics.py | 57 ++++++ scripts/tests/test_session_metrics.py | 71 +++++++ 7 files changed, 428 insertions(+), 152 deletions(-) diff --git a/crates/pet/tests/fixtures/session_sitecustomize.py b/crates/pet/tests/fixtures/session_sitecustomize.py index e7c4c980..30258709 100644 --- a/crates/pet/tests/fixtures/session_sitecustomize.py +++ b/crates/pet/tests/fixtures/session_sitecustomize.py @@ -10,10 +10,12 @@ barrier = os.environ.get("PET_SESSION_RESOLVE_BARRIER") resolve_root = os.environ.get("PET_SESSION_RESOLVE_ROOT") if barrier and resolve_root: - executable = os.path.normcase(os.path.abspath(sys.executable)) - root_executable_prefix = os.path.normcase(os.path.abspath(resolve_root)) + os.sep - if not executable.startswith(root_executable_prefix): + 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) diff --git a/crates/pet/tests/jsonrpc_client.rs b/crates/pet/tests/jsonrpc_client.rs index 73b23e14..57260516 100644 --- a/crates/pet/tests/jsonrpc_client.rs +++ b/crates/pet/tests/jsonrpc_client.rs @@ -525,57 +525,6 @@ fn spawn_stdout_reader(stdout: ChildStdout, state: Arc) -> JoinHand }) } -#[cfg(test)] -pub(crate) fn controlled_pending_request( - client: &PetJsonRpcClient, - submitted_at: Instant, -) -> (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) -} - -#[cfg(test)] -pub(crate) fn pending_request_registered( - client: &PetJsonRpcClient, - request: &PendingRequest, -) -> bool { - client - .inner - .state - .pending - .lock() - .unwrap() - .contains_key(&request.id) -} - -#[cfg(test)] -pub(crate) fn pending_request_count(client: &PetJsonRpcClient) -> usize { - client.inner.state.pending.lock().unwrap().len() -} - fn spawn_stderr_reader( stderr: impl Read + Send + 'static, state: Arc, @@ -649,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 1d1f05e1..98c27079 100644 --- a/crates/pet/tests/jsonrpc_server_test.rs +++ b/crates/pet/tests/jsonrpc_server_test.rs @@ -13,77 +13,10 @@ use std::thread::JoinHandle; use std::time::{Duration, Instant}; use tempfile::TempDir; -mod jsonrpc_client; +pub mod jsonrpc_client; use jsonrpc_client::{EnvironmentNotification, PetJsonRpcClient}; -#[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) = - jsonrpc_client::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!(jsonrpc_client::pending_request_count(&client), 0); - - let consumed_after_deadline_at = Instant::now() - Duration::from_secs(31); - let (consumed_after_deadline, send_timely) = - jsonrpc_client::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) = - jsonrpc_client::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!(jsonrpc_client::pending_request_count(&client), 0); - - let expired_at = Instant::now() - Duration::from_secs(31); - let (expired, _keep_sender_connected) = - jsonrpc_client::controlled_pending_request(&client, expired_at); - assert!(jsonrpc_client::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!(jsonrpc_client::pending_request_count(&client), 0); - - let (dropped, _keep_sender_connected) = - jsonrpc_client::controlled_pending_request(&client, Instant::now()); - assert!(jsonrpc_client::pending_request_registered( - &client, &dropped - )); - drop(dropped); - assert_eq!( - jsonrpc_client::pending_request_count(&client), - 0, - "dropping a pending request must remove its map entry" - ); -} - fn frame_with_headers(payload: &[u8], headers: &[(&str, &str)], line_ending: &[u8]) -> Vec { let mut frame = Vec::new(); for (name, value) in headers { diff --git a/crates/pet/tests/session_performance.rs b/crates/pet/tests/session_performance.rs index 8f28631b..eb99916d 100644 --- a/crates/pet/tests/session_performance.rs +++ b/crates/pet/tests/session_performance.rs @@ -125,6 +125,16 @@ impl Fixture { } 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) -> PathBuf { let python = std::env::var_os("PET_SESSION_PYTHON").unwrap_or_else(|| { if cfg!(windows) { "python".into() @@ -132,22 +142,18 @@ impl Fixture { "python3".into() } }); - (0..count) - .map(|index| { - let prefix = self.resolve_root().join(format!("resolve-{index}")); - 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 {index}: {}", - String::from_utf8_lossy(&output.stderr) - ); - python_executable(&bin_directory(&prefix), false) - }) - .collect() + 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) + ); + python_executable(&bin_directory(prefix), false) } fn resolve_root(&self) -> PathBuf { @@ -460,6 +466,46 @@ fn cache_json((files, bytes): (usize, u64)) -> serde_json::Value { }) } +fn resolve_and_measure( + client: &PetJsonRpcClient, + executable: &Path, +) -> (EnvironmentIdentity, u128) { + let executable = executable + .to_str() + .expect("resolve fixture path was not UTF-8"); + let (result, latency) = client + .submit_resolve(executable) + .expect("failed to submit cache-control resolve") + .wait(REQUEST_TIMEOUT) + .expect("cache-control resolve failed"); + let environment: EnvironmentNotification = + serde_json::from_value(result).expect("resolve returned an invalid environment"); + assert_eq!( + environment.error, None, + "cache-control resolve reported an error" + ); + ( + EnvironmentIdentity { + executable: norm_case( + environment + .executable + .expect("resolved environment had no executable"), + ), + prefix: norm_case( + environment + .prefix + .expect("resolved environment had no prefix"), + ), + kind: environment + .kind + .expect("resolved environment had no classification"), + name: environment.name, + version: environment.version, + }, + latency.as_micros(), + ) +} + fn shutdown_client(client: &PetJsonRpcClient, context: &str) { let status = client .shutdown(Duration::from_secs(10)) @@ -565,6 +611,94 @@ fn long_lived_session_benchmark() { let mut max_ambient_environment_count = 0; let mut max_ambient_manager_count = 0; + 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(); @@ -660,9 +794,6 @@ fn long_lived_session_benchmark() { })); } - let barrier = fixture.barrier.as_os_str(); - let python_path = fixture.python_path.as_os_str(); - let resolve_root = fixture.resolve_root(); let client = PetJsonRpcClient::spawn_with_environment(&[ ("PET_SESSION_RESOLVE_BARRIER", barrier), ("PET_SESSION_RESOLVE_ROOT", resolve_root.as_os_str()), @@ -770,9 +901,12 @@ fn long_lived_session_benchmark() { fixture.clear_barrier(); let pre_resolve_resources = observe_resources(client.process_id()); + let original_cache_before_overlap = directory_usage(&fixture.cache); + assert!( + original_cache_before_overlap.0 > 0 && original_cache_before_overlap.1 > 0, + "warm-up resolve did not populate the original process cache" + ); let overlap_workspace = fixture._root.path().join("overlap-workspace"); - let overlap_cache = fixture._root.path().join("overlap-cache"); - fs::create_dir_all(&overlap_cache).expect("failed to create overlap cache"); let mut overlap_expected = fixture.reset_inventory_at(&overlap_workspace, 3); overlap_expected.sort_unstable(); let overlap_executables = resolve_executables @@ -810,7 +944,7 @@ fn long_lived_session_benchmark() { client .configure(json!({ "workspaceDirectories": [&overlap_workspace], - "cacheDirectory": &overlap_cache, + "cacheDirectory": &fixture.cache, })) .expect("configure was not responsive while resolves were blocked"); client.clear_notifications(); @@ -838,6 +972,11 @@ fn long_lived_session_benchmark() { max_ambient_environment_count = max_ambient_environment_count.max(overlap_ambient_environment_count); max_ambient_manager_count = max_ambient_manager_count.max(overlap_ambient_manager_count); + assert_eq!( + directory_usage(&fixture.cache), + original_cache_before_overlap, + "overlap reconfiguration did not retain the original cache contents" + ); assert_eq!( fixture.released_count(), 0, @@ -976,6 +1115,7 @@ fn long_lived_session_benchmark() { "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(), @@ -983,6 +1123,8 @@ fn long_lived_session_benchmark() { }, "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, diff --git a/docs/SESSION_BENCHMARKS.md b/docs/SESSION_BENCHMARKS.md index fb76ea27..7fe08032 100644 --- a/docs/SESSION_BENCHMARKS.md +++ b/docs/SESSION_BENCHMARKS.md @@ -2,12 +2,20 @@ `session_performance` runs deterministic virtual-environment fixtures against real PET servers. For every inventory size it labels and measures three distinct -refresh scenarios: a first server process with an explicitly empty cache -directory, a new process after that first refresh using the same cache +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. The new-process scenario records whether -the first refresh actually wrote nonzero cache bytes; a zero-byte directory is -not reported as a disk-cache hit. +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 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, @@ -50,10 +58,14 @@ 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: +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 and cache, and refreshes that known -inventory. The barrier remains held until those responsiveness checks complete. +`info`, reconfigures to a new workspace while retaining the process's original +cache directory, and refreshes that known inventory. 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 diff --git a/scripts/session_metrics.py b/scripts/session_metrics.py index 2cb1b4a9..e6566070 100644 --- a/scripts/session_metrics.py +++ b/scripts/session_metrics.py @@ -19,6 +19,7 @@ "firstProcessEmptyDiskCache", "newProcessAfterFirstRefresh", "sameProcessWarm", + "persistentCacheResolve", "cacheUsage", "inventoryResources", "resolveLatency", @@ -32,6 +33,8 @@ "measurementCounts", "refreshScenarioSamples", "cacheUsage", + "persistentCacheResolveSamples", + "persistentCacheAfterColdResolve", "resolveConcurrency", "coldResolveBatches", "resolveLatencyUs", @@ -55,6 +58,11 @@ "newProcessAfterFirstRefresh", "sameProcessWarm", } +PERSISTENT_CACHE_SCENARIOS = { + "cold": 1, + "diskWarm": 0, + "sameProcessWarm": 0, +} class MetricsError(ValueError): @@ -127,6 +135,7 @@ def validate_success_metrics(value: Any, expected_mode: str) -> dict[str, Any]: "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, @@ -213,6 +222,54 @@ def validate_success_metrics(value: Any, expected_mode: str) -> dict[str, Any]: "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: diff --git a/scripts/tests/test_session_metrics.py b/scripts/tests/test_session_metrics.py index 638c155d..d0fd60f8 100644 --- a/scripts/tests/test_session_metrics.py +++ b/scripts/tests/test_session_metrics.py @@ -3,6 +3,8 @@ import copy import json +import os +import runpy import sys import tempfile import unittest @@ -52,6 +54,7 @@ def valid_metrics(): "firstProcessEmptyDiskCache": 3, "newProcessAfterFirstRefresh": 3, "sameProcessWarm": 9, + "persistentCacheResolve": 3, "cacheUsage": 3, "inventoryResources": 3, "resolveLatency": 8, @@ -68,6 +71,20 @@ def valid_metrics(): } 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, @@ -92,6 +109,41 @@ def valid_metrics(): } +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() @@ -137,6 +189,22 @@ def test_count_mismatch_fails_closed(self): 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"): @@ -162,6 +230,9 @@ def test_malformed_nested_types_fail_closed_with_failure_artifacts(self): ("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"), ] From 8faaabc87477fa0da69bbbf072488138f57083c3 Mon Sep 17 00:00:00 2001 From: Karthik Nadig Date: Wed, 30 Sep 2026 10:43:06 -0700 Subject: [PATCH 3/4] test: validate session resource and cache invariants Allow ambient cache additions while preserving original entries, validate derived resource metrics exactly, and surface Linux descriptor enumeration errors. Address PR #563 feedback and native CI failures. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- crates/pet/tests/session_performance.rs | 69 +++++++++++++++--- docs/SESSION_BENCHMARKS.md | 13 ++-- scripts/session_metrics.py | 55 ++++++++++++-- scripts/tests/test_session_metrics.py | 96 +++++++++++++++++++++++-- 4 files changed, 209 insertions(+), 24 deletions(-) diff --git a/crates/pet/tests/session_performance.rs b/crates/pet/tests/session_performance.rs index eb99916d..4adf9e25 100644 --- a/crates/pet/tests/session_performance.rs +++ b/crates/pet/tests/session_performance.rs @@ -3,7 +3,9 @@ use pet_fs::path::norm_case; use serde_json::json; +use std::collections::BTreeMap; use std::fs; +use std::io; use std::path::{Path, PathBuf}; use std::process::Command; use std::thread; @@ -420,6 +422,55 @@ fn directory_usage(root: &Path) -> (usize, u64) { (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(|_| { @@ -534,8 +585,8 @@ fn sample_process(pid: u32) -> ResourceSample { descriptors: Some( fs::read_dir(format!("/proc/{pid}/fd")) .expect("failed to read process-specific descriptor directory") - .map(|entry| entry.expect("failed to read process descriptor entry")) - .count(), + .try_fold(0, |count, entry| entry.map(|_| count + 1)) + .expect("failed to read process descriptor entry"), ), } } @@ -901,9 +952,10 @@ fn long_lived_session_benchmark() { fixture.clear_barrier(); let pre_resolve_resources = observe_resources(client.process_id()); - let original_cache_before_overlap = directory_usage(&fixture.cache); + let original_cache_before_overlap = + cache_contents(&fixture.cache).expect("failed to capture pre-overlap cache contents"); assert!( - original_cache_before_overlap.0 > 0 && original_cache_before_overlap.1 > 0, + !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"); @@ -972,10 +1024,11 @@ fn long_lived_session_benchmark() { max_ambient_environment_count = max_ambient_environment_count.max(overlap_ambient_environment_count); max_ambient_manager_count = max_ambient_manager_count.max(overlap_ambient_manager_count); - assert_eq!( - directory_usage(&fixture.cache), - original_cache_before_overlap, - "overlap reconfiguration did not retain the original cache contents" + 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(), diff --git a/docs/SESSION_BENCHMARKS.md b/docs/SESSION_BENCHMARKS.md index 7fe08032..dc927944 100644 --- a/docs/SESSION_BENCHMARKS.md +++ b/docs/SESSION_BENCHMARKS.md @@ -53,8 +53,11 @@ 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. Artifact writes are atomic, and the timeout fallback replaces corrupt -artifacts rather than uploading invalid JSON. Failed runs retain validated +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. @@ -64,8 +67,10 @@ 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. The barrier remains held -until those responsiveness checks complete. +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 diff --git a/scripts/session_metrics.py b/scripts/session_metrics.py index e6566070..7d0fb9d2 100644 --- a/scripts/session_metrics.py +++ b/scripts/session_metrics.py @@ -91,7 +91,7 @@ def require_integer_list(value: Any, name: str, expected_length: int) -> list[in return [require_integer(item, f"{name}[{index}]") for index, item in enumerate(value)] -def validate_resource(value: Any, name: str) -> None: +def validate_resource(value: Any, name: str) -> dict[str, Any]: resource = require_object( value, name, {"residentBytes", "threads", "handlesOrDescriptors"} ) @@ -99,6 +99,19 @@ def validate_resource(value: Any, name: str) -> None: 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]: @@ -299,6 +312,7 @@ def validate_success_metrics(value: Any, expected_mode: str) -> dict[str, Any]: 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, @@ -311,7 +325,11 @@ def validate_success_metrics(value: Any, expected_mode: str) -> dict[str, Any]: f"inventoryResourceSamples[{index}] has an invalid or duplicate size" ) seen_resource_sizes.add(size) - validate_resource(sample["resources"], f"inventoryResourceSamples[{index}].resources") + 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: @@ -329,17 +347,42 @@ def validate_success_metrics(value: Any, expected_mode: str) -> dict[str, Any]: f"resolveBatchResourceSamples[{index}] has an invalid or duplicate batch" ) seen_batches.add(batch) - validate_resource(sample["resources"], f"resolveBatchResourceSamples[{index}].resources") + resource_samples.append( + validate_resource( + sample["resources"], f"resolveBatchResourceSamples[{index}].resources" + ) + ) for name in ( "preResolveResources", "barrierObservedResources", "postOverlapResources", - "observedResourcePeak", "resourceAfter", ): - validate_resource(metrics[name], name) - require_integer(metrics["rssDeltaFromPreResolveBytes"], "rssDeltaFromPreResolveBytes", minimum=-(2**63)) + 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 diff --git a/scripts/tests/test_session_metrics.py b/scripts/tests/test_session_metrics.py index d0fd60f8..31711521 100644 --- a/scripts/tests/test_session_metrics.py +++ b/scripts/tests/test_session_metrics.py @@ -19,11 +19,11 @@ failed_metrics, write_json) # noqa: E402 -def resource(): +def resource(resident=100, threads=2, descriptors=3): return { - "residentBytes": 100, - "threads": 2, - "handlesOrDescriptors": 3, + "residentBytes": resident, + "threads": threads, + "handlesOrDescriptors": descriptors, } @@ -100,11 +100,11 @@ def valid_metrics(): "resolveBatchResourceSamples": [ {"batch": batch, "resources": resource()} for batch in range(2) ], - "preResolveResources": resource(), + "preResolveResources": resource(resident=100), "barrierObservedResources": resource(), "postOverlapResources": resource(), "observedResourcePeak": resource(), - "resourceAfter": resource(), + "resourceAfter": resource(resident=90), "rssDeltaFromPreResolveBytes": -10, } @@ -163,6 +163,13 @@ def write_payloads(self, *payloads): 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"): @@ -224,6 +231,83 @@ def test_success_preserves_validated_metrics(self): self.assertTrue(metrics["metricsProduced"]) self.assertEqual(len(metrics["refreshScenarioSamples"]), 9) + 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",), From f8db49d62244d9aa8335a05c540fb18b0a017384 Mon Sep 17 00:00:00 2001 From: Karthik Nadig Date: Wed, 30 Sep 2026 11:25:51 -0700 Subject: [PATCH 4/4] test: account for every ambient refresh observation Include churn observations in ambient maxima and reject overlap counts larger than their maxima. Address PR #563 review feedback without changing workload or coverage budgets. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- crates/pet/tests/session_performance.rs | 82 ++++++++++++++++--------- docs/SESSION_BENCHMARKS.md | 2 + scripts/session_metrics.py | 9 +++ scripts/tests/test_session_metrics.py | 22 +++++++ 4 files changed, 86 insertions(+), 29 deletions(-) diff --git a/crates/pet/tests/session_performance.rs b/crates/pet/tests/session_performance.rs index ec8cb36a..5901765c 100644 --- a/crates/pet/tests/session_performance.rs +++ b/crates/pet/tests/session_performance.rs @@ -47,6 +47,19 @@ struct RefreshMeasurement { 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, @@ -426,7 +439,11 @@ fn first_result_timing_ignores_ambient_environments_and_managers() { ); } -fn refresh_and_measure(client: &PetJsonRpcClient, workspace: &Path) -> RefreshMeasurement { +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] }))) @@ -448,13 +465,30 @@ fn refresh_and_measure(client: &PetJsonRpcClient, workspace: &Path) -> RefreshMe fixture_manager_count, 0, "fixture workspace unexpectedly reported an environment manager" ); - RefreshMeasurement { + 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) { @@ -830,8 +864,7 @@ fn long_lived_session_benchmark() { 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 max_ambient_environment_count = 0; - let mut max_ambient_manager_count = 0; + let mut ambient_count_peak = AmbientCountPeak::default(); let barrier = fixture.barrier.as_os_str(); let python_path = fixture.python_path.as_os_str(); @@ -940,14 +973,12 @@ fn long_lived_session_benchmark() { "cacheDirectory": &fixture.cache, })) .expect("failed to configure first-process scenario server"); - let first = refresh_and_measure(&first_process, &fixture.workspace); + 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}" ); - max_ambient_environment_count = - max_ambient_environment_count.max(first.ambient_environment_count); - max_ambient_manager_count = max_ambient_manager_count.max(first.ambient_manager_count); refresh_scenario_samples.push(json!({ "scenario": "firstProcessEmptyDiskCache", "inventorySize": size, @@ -966,14 +997,11 @@ fn long_lived_session_benchmark() { "cacheDirectory": &fixture.cache, })) .expect("failed to configure reused-cache scenario server"); - let reused = refresh_and_measure(&new_process, &fixture.workspace); + 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}" ); - max_ambient_environment_count = - max_ambient_environment_count.max(reused.ambient_environment_count); - max_ambient_manager_count = max_ambient_manager_count.max(reused.ambient_manager_count); refresh_scenario_samples.push(json!({ "scenario": "newProcessAfterFirstRefresh", "inventorySize": size, @@ -985,14 +1013,12 @@ fn long_lived_session_benchmark() { 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); + 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}" ); - max_ambient_environment_count = - max_ambient_environment_count.max(warm.ambient_environment_count); - max_ambient_manager_count = max_ambient_manager_count.max(warm.ambient_manager_count); warm_round_trip_us.push(warm.round_trip_us); warm_ttfe_us.push(warm.ttfe_us); } @@ -1031,11 +1057,8 @@ fn long_lived_session_benchmark() { let mut churn_expected = fixture.reset_inventory(10); churn_expected.sort_unstable(); - let initial = refresh_and_measure(&client, &fixture.workspace); + let initial = refresh_and_measure(&client, &fixture.workspace, &mut ambient_count_peak); assert_eq!(initial.inventory, churn_expected); - max_ambient_environment_count = - max_ambient_environment_count.max(initial.ambient_environment_count); - max_ambient_manager_count = max_ambient_manager_count.max(initial.ambient_manager_count); let removed_prefix = fixture.workspace.join("env-0000"); let removed_index = churn_expected @@ -1046,7 +1069,7 @@ fn long_lived_session_benchmark() { 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); + 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, @@ -1074,7 +1097,7 @@ fn long_lived_session_benchmark() { version: Some("3.13.2".to_string()), }); churn_expected.sort_unstable(); - let edited = refresh_and_measure(&client, &fixture.workspace); + 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"); @@ -1095,7 +1118,7 @@ fn long_lived_session_benchmark() { ..previous }); churn_expected.sort_unstable(); - let aliased = refresh_and_measure(&client, &fixture.workspace); + let aliased = refresh_and_measure(&client, &fixture.workspace, &mut ambient_count_peak); assert_eq!(aliased.inventory, churn_expected); fixture.clear_barrier(); @@ -1176,9 +1199,10 @@ fn long_lived_session_benchmark() { fixture_manager_count, 0, "overlap fixture unexpectedly reported an environment manager" ); - max_ambient_environment_count = - max_ambient_environment_count.max(overlap_ambient_environment_count); - max_ambient_manager_count = max_ambient_manager_count.max(overlap_ambient_manager_count); + 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!( @@ -1325,8 +1349,8 @@ fn long_lived_session_benchmark() { "overlapProcessesStarted": resolve_concurrency, "overlapAmbientEnvironmentCount": overlap_ambient_environment_count, "overlapAmbientManagerCount": overlap_ambient_manager_count, - "maxAmbientEnvironmentCount": max_ambient_environment_count, - "maxAmbientManagerCount": max_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() diff --git a/docs/SESSION_BENCHMARKS.md b/docs/SESSION_BENCHMARKS.md index 0e22984e..a88034f3 100644 --- a/docs/SESSION_BENCHMARKS.md +++ b/docs/SESSION_BENCHMARKS.md @@ -78,6 +78,8 @@ 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 diff --git a/scripts/session_metrics.py b/scripts/session_metrics.py index 7d0fb9d2..61db2e41 100644 --- a/scripts/session_metrics.py +++ b/scripts/session_metrics.py @@ -307,6 +307,15 @@ def validate_success_metrics(value: Any, expected_mode: str) -> dict[str, Any]: 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): diff --git a/scripts/tests/test_session_metrics.py b/scripts/tests/test_session_metrics.py index 31711521..6f3ece2f 100644 --- a/scripts/tests/test_session_metrics.py +++ b/scripts/tests/test_session_metrics.py @@ -231,6 +231,28 @@ def test_success_preserves_validated_metrics(self): 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):