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