| 1 | //! Measurement-only persistence backlog tests. |
| 2 | |
| 3 | use super::*; |
| 4 | use std::time::{Duration, Instant}; |
| 5 | |
| 6 | use crate::models::{ContentBlock, Message}; |
| 7 | |
| 8 | const BACKLOG_RECEIPT_PATH_ENV: &str = "CODEWHALE_TEST_PERSISTENCE_BACKLOG_RECEIPT_PATH"; |
| 9 | const BACKLOG_SOURCE_SHA_ENV: &str = "CODEWHALE_TEST_PERSISTENCE_BACKLOG_SOURCE_SHA"; |
| 10 | const BACKLOG_SOURCE_DIRTY_ENV: &str = "CODEWHALE_TEST_PERSISTENCE_BACKLOG_SOURCE_DIRTY"; |
| 11 | const BACKLOG_RUSTC_VERSION_ENV: &str = "CODEWHALE_TEST_PERSISTENCE_BACKLOG_RUSTC_VERSION"; |
| 12 | const BACKLOG_CARGO_VERSION_ENV: &str = "CODEWHALE_TEST_PERSISTENCE_BACKLOG_CARGO_VERSION"; |
| 13 | const BACKLOG_FIXTURE_ID: &str = "paused-production-channel-session-snapshot-v1"; |
| 14 | const BACKLOG_REQUESTS_ATTEMPTED: usize = 128; |
| 15 | const BACKLOG_CONTENT_BYTES_PER_REQUEST: usize = 64 * 1024; |
| 16 | const BACKLOG_EXPECTED_APPLIED_VERSION: usize = BACKLOG_REQUESTS_ATTEMPTED - 1; |
| 17 | const BACKLOG_SESSION_ID: &str = "persistence-backlog-single-session"; |
| 18 | const BACKLOG_PAYLOAD_ESTIMATOR: &str = "retained-saved-session-json-bytes-v1"; |
| 19 | const BACKLOG_REQUEST_VARIANT: &str = "session_snapshot"; |
| 20 | |
| 21 | #[derive(Debug)] |
| 22 | struct PersistenceBacklogObservation { |
| 23 | accepted_requests: usize, |
| 24 | retained_queued_requests: usize, |
| 25 | estimated_retained_payload_bytes: usize, |
| 26 | applied_version: Option<usize>, |
| 27 | enqueue_elapsed_ns: u128, |
| 28 | rss_before_bytes: Option<u64>, |
| 29 | rss_during_bytes: Option<u64>, |
| 30 | rss_after_bytes: Option<u64>, |
| 31 | } |
| 32 | |
| 33 | #[derive(Debug)] |
| 34 | struct PersistenceBacklogProvenance { |
| 35 | source_sha: String, |
| 36 | source_dirty: bool, |
| 37 | rustc_version: String, |
| 38 | cargo_version: String, |
| 39 | } |
| 40 | |
| 41 | fn persistence_backlog_receipt( |
| 42 | observation: &PersistenceBacklogObservation, |
| 43 | provenance: &PersistenceBacklogProvenance, |
| 44 | ) -> serde_json::Value { |
| 45 | let rss_supported = observation.rss_before_bytes.is_some() |
| 46 | && observation.rss_during_bytes.is_some() |
| 47 | && observation.rss_after_bytes.is_some(); |
| 48 | serde_json::json!({ |
| 49 | "document_kind": "codewhale.persistence_backlog_receipt", |
| 50 | "schema_version": 2, |
| 51 | "source_sha": provenance.source_sha, |
| 52 | "source_dirty": provenance.source_dirty, |
| 53 | "rustc_version": provenance.rustc_version, |
| 54 | "cargo_version": provenance.cargo_version, |
| 55 | "build_profile": "test", |
| 56 | "sample_count": 1, |
| 57 | "fixture_id": BACKLOG_FIXTURE_ID, |
| 58 | "platform": std::env::consts::OS, |
| 59 | "request_variant": BACKLOG_REQUEST_VARIANT, |
| 60 | "payload_estimator": BACKLOG_PAYLOAD_ESTIMATOR, |
| 61 | "paused_consumer": true, |
| 62 | "requests_attempted": BACKLOG_REQUESTS_ATTEMPTED, |
| 63 | "content_bytes_per_request": BACKLOG_CONTENT_BYTES_PER_REQUEST, |
| 64 | "single_session_id": true, |
| 65 | "expected_applied_version": BACKLOG_EXPECTED_APPLIED_VERSION, |
| 66 | "accepted_requests": observation.accepted_requests, |
| 67 | "retained_queued_requests": observation.retained_queued_requests, |
| 68 | "estimated_retained_payload_bytes": observation.estimated_retained_payload_bytes, |
| 69 | "applied_version": observation.applied_version, |
| 70 | "final_version_applied": observation.applied_version == Some(BACKLOG_EXPECTED_APPLIED_VERSION), |
| 71 | "enqueue_elapsed_ns": observation.enqueue_elapsed_ns, |
| 72 | "rss_supported": rss_supported, |
| 73 | "rss_before_bytes": observation.rss_before_bytes, |
| 74 | "rss_during_bytes": observation.rss_during_bytes, |
| 75 | "rss_after_bytes": observation.rss_after_bytes, |
| 76 | "rss_during_delta_bytes": observation.rss_during_bytes.zip(observation.rss_before_bytes) |
| 77 | .map(|(during, before)| during.saturating_sub(before)), |
| 78 | "rss_after_delta_bytes": observation.rss_after_bytes.zip(observation.rss_before_bytes) |
| 79 | .map(|(after, before)| after.saturating_sub(before)), |
| 80 | "limitations": [ |
| 81 | "current-process RSS samples are available on macOS only", |
| 82 | "serialized SavedSession bytes estimate retained heap payload after draining the paused production receiver; allocator and channel overhead are observed only through RSS", |
| 83 | "one bounded sample characterizes backlog retention but does not prove an asymptotic growth rate", |
| 84 | ], |
| 85 | }) |
| 86 | } |
| 87 | |
| 88 | #[cfg(target_os = "macos")] |
| 89 | fn current_process_rss_bytes() -> Option<u64> { |
| 90 | let output = std::process::Command::new("/bin/ps") |
| 91 | .args(["-o", "rss=", "-p", &std::process::id().to_string()]) |
| 92 | .output() |
| 93 | .ok()?; |
| 94 | if !output.status.success() { |
| 95 | return None; |
| 96 | } |
| 97 | let kib = String::from_utf8(output.stdout) |
| 98 | .ok()? |
| 99 | .trim() |
| 100 | .parse::<u64>() |
| 101 | .ok()?; |
| 102 | kib.checked_mul(1024) |
| 103 | } |
| 104 | |
| 105 | #[cfg(not(target_os = "macos"))] |
| 106 | fn current_process_rss_bytes() -> Option<u64> { |
| 107 | None |
| 108 | } |
| 109 | |
| 110 | fn backlog_session(workspace: &std::path::Path, index: usize) -> SavedSession { |
| 111 | let messages = [Message { |
| 112 | role: "assistant".to_string(), |
| 113 | content: vec![ContentBlock::Text { |
| 114 | text: "x".repeat(BACKLOG_CONTENT_BYTES_PER_REQUEST), |
| 115 | cache_control: None, |
| 116 | }], |
| 117 | }]; |
| 118 | let mut session = crate::session_manager::create_saved_session_with_mode( |
| 119 | &messages, |
| 120 | "measurement-model", |
| 121 | workspace, |
| 122 | 0, |
| 123 | None, |
| 124 | Some("agent"), |
| 125 | ); |
| 126 | session.metadata.id = BACKLOG_SESSION_ID.to_string(); |
| 127 | session.metadata.title = format!("Persistence backlog version {index:04}"); |
| 128 | session |
| 129 | } |
| 130 | |
| 131 | fn retained_request_payload_bytes(request: &PersistRequest) -> usize { |
| 132 | let PersistRequest::SessionSnapshot(session) = request else { |
| 133 | panic!("backlog fixture retained an unexpected request variant") |
| 134 | }; |
| 135 | serde_json::to_vec(session) |
| 136 | .expect("retained measurement session must serialize") |
| 137 | .len() |
| 138 | } |
| 139 | |
| 140 | fn saved_session_version(session: &SavedSession) -> Option<usize> { |
| 141 | session |
| 142 | .metadata |
| 143 | .title |
| 144 | .strip_prefix("Persistence backlog version ") |
| 145 | .and_then(|value| value.parse::<usize>().ok()) |
| 146 | } |
| 147 | |
| 148 | fn drain_retained_requests(receiver: &mut PersistRequestReceiver) -> (usize, usize, Option<usize>) { |
| 149 | let mut retained_queued_requests = 0; |
| 150 | let mut estimated_retained_payload_bytes = 0; |
| 151 | let mut pending = PendingState::default(); |
| 152 | while let Ok(request) = receiver.try_recv() { |
| 153 | retained_queued_requests += 1; |
| 154 | estimated_retained_payload_bytes += retained_request_payload_bytes(&request); |
| 155 | assert!(matches!(pending.absorb(request), Control::Continue)); |
| 156 | } |
| 157 | let applied_version = pending |
| 158 | .sessions |
| 159 | .get(BACKLOG_SESSION_ID) |
| 160 | .and_then(saved_session_version); |
| 161 | ( |
| 162 | retained_queued_requests, |
| 163 | estimated_retained_payload_bytes, |
| 164 | applied_version, |
| 165 | ) |
| 166 | } |
| 167 | |
| 168 | fn measure_paused_persistence_backlog() -> PersistenceBacklogObservation { |
| 169 | let tmp = tempfile::tempdir().expect("isolated measurement home"); |
| 170 | let _env_lock = crate::test_support::lock_test_env(); |
| 171 | let _home = crate::test_support::EnvVarGuard::set("HOME", tmp.path()); |
| 172 | let _userprofile = crate::test_support::EnvVarGuard::set("USERPROFILE", tmp.path()); |
| 173 | let _codewhale_home = |
| 174 | crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", tmp.path().join(".codewhale")); |
| 175 | |
| 176 | let (tx, mut receiver) = persistence_request_channel(); |
| 177 | let handle = PersistActorHandle { tx }; |
| 178 | let rss_before_bytes = current_process_rss_bytes(); |
| 179 | let mut accepted_requests = 0; |
| 180 | let mut enqueue_elapsed_ns = 0; |
| 181 | |
| 182 | for index in 0..BACKLOG_REQUESTS_ATTEMPTED { |
| 183 | let session = backlog_session(tmp.path(), index); |
| 184 | let started = Instant::now(); |
| 185 | let accepted = handle.try_send(PersistRequest::SessionSnapshot(session)); |
| 186 | enqueue_elapsed_ns += started.elapsed().as_nanos(); |
| 187 | if accepted { |
| 188 | accepted_requests += 1; |
| 189 | } |
| 190 | } |
| 191 | |
| 192 | // The receiver has deliberately never been polled: RSS is sampled while |
| 193 | // the production channel still owns every representation it retained. |
| 194 | std::hint::black_box(&receiver); |
| 195 | let rss_during_bytes = current_process_rss_bytes(); |
| 196 | |
| 197 | let (retained_queued_requests, estimated_retained_payload_bytes, applied_version) = |
| 198 | drain_retained_requests(&mut receiver); |
| 199 | drop(handle); |
| 200 | drop(receiver); |
| 201 | std::thread::sleep(Duration::from_millis(50)); |
| 202 | let rss_after_bytes = current_process_rss_bytes(); |
| 203 | |
| 204 | PersistenceBacklogObservation { |
| 205 | accepted_requests, |
| 206 | retained_queued_requests, |
| 207 | estimated_retained_payload_bytes, |
| 208 | applied_version, |
| 209 | enqueue_elapsed_ns, |
| 210 | rss_before_bytes, |
| 211 | rss_during_bytes, |
| 212 | rss_after_bytes, |
| 213 | } |
| 214 | } |
| 215 | |
| 216 | #[test] |
| 217 | fn persistence_backlog_receipt_contract_keeps_required_fields() { |
| 218 | let receipt = persistence_backlog_receipt( |
| 219 | &PersistenceBacklogObservation { |
| 220 | accepted_requests: 2, |
| 221 | retained_queued_requests: 1, |
| 222 | estimated_retained_payload_bytes: 4_096, |
| 223 | applied_version: Some(BACKLOG_EXPECTED_APPLIED_VERSION), |
| 224 | enqueue_elapsed_ns: 1_000, |
| 225 | rss_before_bytes: Some(10_000), |
| 226 | rss_during_bytes: Some(14_000), |
| 227 | rss_after_bytes: Some(11_000), |
| 228 | }, |
| 229 | &PersistenceBacklogProvenance { |
| 230 | source_sha: "0123456789abcdef0123456789abcdef01234567".to_string(), |
| 231 | source_dirty: false, |
| 232 | rustc_version: "rustc test".to_string(), |
| 233 | cargo_version: "cargo test".to_string(), |
| 234 | }, |
| 235 | ); |
| 236 | let object = receipt.as_object().expect("receipt object"); |
| 237 | for field in [ |
| 238 | "document_kind", |
| 239 | "schema_version", |
| 240 | "source_sha", |
| 241 | "source_dirty", |
| 242 | "rustc_version", |
| 243 | "cargo_version", |
| 244 | "build_profile", |
| 245 | "sample_count", |
| 246 | "fixture_id", |
| 247 | "platform", |
| 248 | "request_variant", |
| 249 | "payload_estimator", |
| 250 | "paused_consumer", |
| 251 | "requests_attempted", |
| 252 | "content_bytes_per_request", |
| 253 | "single_session_id", |
| 254 | "expected_applied_version", |
| 255 | "accepted_requests", |
| 256 | "retained_queued_requests", |
| 257 | "estimated_retained_payload_bytes", |
| 258 | "applied_version", |
| 259 | "final_version_applied", |
| 260 | "enqueue_elapsed_ns", |
| 261 | "rss_supported", |
| 262 | "rss_before_bytes", |
| 263 | "rss_during_bytes", |
| 264 | "rss_after_bytes", |
| 265 | "rss_during_delta_bytes", |
| 266 | "rss_after_delta_bytes", |
| 267 | "limitations", |
| 268 | ] { |
| 269 | assert!(object.contains_key(field), "receipt lost `{field}`"); |
| 270 | } |
| 271 | assert_eq!(receipt["final_version_applied"], true); |
| 272 | assert_eq!(receipt["rss_during_delta_bytes"], 4_000); |
| 273 | assert_eq!(receipt["rss_after_delta_bytes"], 1_000); |
| 274 | } |
| 275 | |
| 276 | #[test] |
| 277 | fn applied_version_follows_production_drain_order_instead_of_numeric_maximum() { |
| 278 | let tmp = tempfile::tempdir().expect("isolated measurement workspace"); |
| 279 | let (tx, mut receiver) = persistence_request_channel(); |
| 280 | tx.send(PersistRequest::SessionSnapshot(backlog_session( |
| 281 | tmp.path(), |
| 282 | BACKLOG_EXPECTED_APPLIED_VERSION, |
| 283 | ))) |
| 284 | .expect("send final version first"); |
| 285 | tx.send(PersistRequest::SessionSnapshot(backlog_session( |
| 286 | tmp.path(), |
| 287 | BACKLOG_EXPECTED_APPLIED_VERSION - 1, |
| 288 | ))) |
| 289 | .expect("send stale version last"); |
| 290 | |
| 291 | let (_, _, applied_version) = drain_retained_requests(&mut receiver); |
| 292 | assert_eq!( |
| 293 | applied_version, |
| 294 | Some(BACKLOG_EXPECTED_APPLIED_VERSION - 1), |
| 295 | "production PendingState must expose stale last-drained ordering instead of hiding it behind max(version)" |
| 296 | ); |
| 297 | } |
| 298 | |
| 299 | /// Exact, ignored child-process measurement invoked by |
| 300 | /// `scripts/measure-persistence-backlog.py`. |
| 301 | #[test] |
| 302 | #[ignore = "one-shot metric for scripts/measure-persistence-backlog.py"] |
| 303 | fn write_paused_persistence_backlog_measurement_receipt() { |
| 304 | let Ok(path) = std::env::var(BACKLOG_RECEIPT_PATH_ENV) else { |
| 305 | return; |
| 306 | }; |
| 307 | let provenance = PersistenceBacklogProvenance { |
| 308 | source_sha: std::env::var(BACKLOG_SOURCE_SHA_ENV) |
| 309 | .expect("measurement wrapper must provide the exact source SHA"), |
| 310 | source_dirty: std::env::var(BACKLOG_SOURCE_DIRTY_ENV) |
| 311 | .expect("measurement wrapper must provide source dirty state") |
| 312 | .parse::<bool>() |
| 313 | .expect("source dirty state must be true or false"), |
| 314 | rustc_version: std::env::var(BACKLOG_RUSTC_VERSION_ENV) |
| 315 | .expect("measurement wrapper must provide rustc version"), |
| 316 | cargo_version: std::env::var(BACKLOG_CARGO_VERSION_ENV) |
| 317 | .expect("measurement wrapper must provide cargo version"), |
| 318 | }; |
| 319 | let receipt = persistence_backlog_receipt(&measure_paused_persistence_backlog(), &provenance); |
| 320 | std::fs::write( |
| 321 | path, |
| 322 | serde_json::to_vec_pretty(&receipt).expect("serialize backlog receipt"), |
| 323 | ) |
| 324 | .expect("write backlog receipt"); |
| 325 | } |
| 326 |