1use crate::error::{BenchError, Result};
4use indicatif::{ProgressBar, ProgressStyle};
5use std::path::Path;
6use std::process::Stdio;
7use std::sync::Arc;
8use tokio::io::{AsyncBufReadExt, BufReader};
9use tokio::process::Command as TokioCommand;
10
11fn extract_mockforge_marker_json(line: &str, marker: &str) -> Option<String> {
34 let start = line.find(marker)?;
35 let json_start = start + marker.len();
36 let rest = &line[json_start..];
37
38 let is_logfmt = start >= 5 && line.as_bytes().get(start - 5..start) == Some(b"msg=\"");
42 if is_logfmt {
43 let bytes = rest.as_bytes();
48 let mut i = 0;
49 let mut out = String::with_capacity(rest.len());
50 while i < bytes.len() {
51 let b = bytes[i];
52 if b == b'"' {
53 return Some(out);
55 }
56 if b == b'\\' && i + 1 < bytes.len() {
57 let next = bytes[i + 1];
58 match next {
59 b'"' => out.push('"'),
60 b'\\' => out.push('\\'),
61 other => {
65 out.push('\\');
66 out.push(other as char);
67 }
68 }
69 i += 2;
70 continue;
71 }
72 let ch_start = i;
75 i += 1;
77 while i < bytes.len() && (bytes[i] & 0b1100_0000) == 0b1000_0000 {
78 i += 1;
79 }
80 out.push_str(&rest[ch_start..i]);
81 }
82 if out.is_empty() {
85 None
86 } else {
87 Some(out)
88 }
89 } else {
90 let trimmed = rest.trim();
93 if trimmed.is_empty() {
94 None
95 } else {
96 Some(trimmed.to_string())
97 }
98 }
99}
100
101fn extract_exchange_json(line: &str) -> Option<String> {
103 extract_mockforge_marker_json(line, "MOCKFORGE_EXCHANGE:")
104}
105
106fn extract_failure_json(line: &str) -> Option<String> {
108 extract_mockforge_marker_json(line, "MOCKFORGE_FAILURE:")
109}
110
111fn extract_network_event_json(line: &str) -> Option<String> {
115 extract_mockforge_marker_json(line, "MOCKFORGE_NETWORK_EVENT:")
116}
117
118pub struct K6Executor {
120 k6_path: String,
121 local_ips: String,
127 discard_response_bodies: bool,
134}
135
136impl K6Executor {
137 pub fn new() -> Result<Self> {
139 let k6_path = which::which("k6")
140 .map_err(|_| BenchError::K6NotFound)?
141 .to_string_lossy()
142 .to_string();
143
144 Ok(Self {
145 k6_path,
146 local_ips: String::new(),
147 discard_response_bodies: false,
148 })
149 }
150
151 pub fn with_local_ips(mut self, local_ips: impl Into<String>) -> Self {
155 self.local_ips = local_ips.into();
156 self
157 }
158
159 pub fn with_discard_response_bodies(mut self, discard: bool) -> Self {
163 self.discard_response_bodies = discard;
164 self
165 }
166
167 pub fn is_k6_installed() -> bool {
169 which::which("k6").is_ok()
170 }
171
172 pub async fn get_version(&self) -> Result<String> {
174 let output = TokioCommand::new(&self.k6_path)
175 .arg("version")
176 .output()
177 .await
178 .map_err(|e| BenchError::K6ExecutionFailed(e.to_string()))?;
179
180 if !output.status.success() {
181 return Err(BenchError::K6ExecutionFailed("Failed to get k6 version".to_string()));
182 }
183
184 Ok(String::from_utf8_lossy(&output.stdout).trim().to_string())
185 }
186
187 pub async fn execute(
197 &self,
198 script_path: &Path,
199 output_dir: Option<&Path>,
200 verbose: bool,
201 ) -> Result<K6Results> {
202 self.execute_with_port(script_path, output_dir, verbose, None).await
203 }
204
205 pub async fn execute_with_port(
207 &self,
208 script_path: &Path,
209 output_dir: Option<&Path>,
210 verbose: bool,
211 api_port: Option<u16>,
212 ) -> Result<K6Results> {
213 println!("Starting load test...\n");
214
215 let mut cmd = TokioCommand::new(&self.k6_path);
216 cmd.arg("run");
217
218 if let Some(port) = api_port {
221 cmd.arg("--address").arg(format!("localhost:{}", port));
222 }
223
224 if !self.local_ips.is_empty() {
229 cmd.arg("--local-ips").arg(&self.local_ips);
230 }
231
232 if self.discard_response_bodies {
235 cmd.env("K6_DISCARD_RESPONSE_BODIES", "true");
236 }
237
238 if verbose {
245 cmd.arg("--verbose");
246 }
247
248 let abs_script =
250 std::fs::canonicalize(script_path).unwrap_or_else(|_| script_path.to_path_buf());
251 cmd.arg(&abs_script);
252
253 if let Some(dir) = output_dir {
256 cmd.current_dir(dir);
257 }
258
259 cmd.stdout(Stdio::piped());
260 cmd.stderr(Stdio::piped());
261
262 let mut child = cmd.spawn().map_err(|e| BenchError::K6ExecutionFailed(e.to_string()))?;
263
264 let stdout = child
265 .stdout
266 .take()
267 .ok_or_else(|| BenchError::K6ExecutionFailed("Failed to capture stdout".to_string()))?;
268
269 let stderr = child
270 .stderr
271 .take()
272 .ok_or_else(|| BenchError::K6ExecutionFailed("Failed to capture stderr".to_string()))?;
273
274 let stdout_reader = BufReader::new(stdout);
276 let stderr_reader = BufReader::new(stderr);
277
278 let mut stdout_lines = stdout_reader.lines();
279 let mut stderr_lines = stderr_reader.lines();
280
281 let spinner = ProgressBar::new_spinner();
283 spinner.set_style(
284 ProgressStyle::default_spinner().template("{spinner:.green} {msg}").unwrap(),
285 );
286 spinner.set_message("Running load test...");
287
288 let failure_details: Arc<tokio::sync::Mutex<Vec<String>>> =
291 Arc::new(tokio::sync::Mutex::new(Vec::new()));
292 let fd_stdout = Arc::clone(&failure_details);
293 let fd_stderr = Arc::clone(&failure_details);
294
295 let exchange_details: Arc<tokio::sync::Mutex<Vec<String>>> =
297 Arc::new(tokio::sync::Mutex::new(Vec::new()));
298 let ex_stdout = Arc::clone(&exchange_details);
299 let ex_stderr = Arc::clone(&exchange_details);
300
301 let network_events: Arc<tokio::sync::Mutex<Vec<String>>> =
306 Arc::new(tokio::sync::Mutex::new(Vec::new()));
307 let ne_stdout = Arc::clone(&network_events);
308 let ne_stderr = Arc::clone(&network_events);
309
310 let log_lines: Arc<tokio::sync::Mutex<Vec<String>>> =
312 Arc::new(tokio::sync::Mutex::new(Vec::new()));
313 let log_stdout = Arc::clone(&log_lines);
314 let log_stderr = Arc::clone(&log_lines);
315
316 let stdout_handle = tokio::spawn(async move {
318 while let Ok(Some(line)) = stdout_lines.next_line().await {
319 log_stdout.lock().await.push(format!("[stdout] {}", line));
320 if let Some(json_str) = extract_failure_json(&line) {
321 fd_stdout.lock().await.push(json_str);
322 } else if let Some(json_str) = extract_exchange_json(&line) {
323 ex_stdout.lock().await.push(json_str);
324 } else if let Some(json_str) = extract_network_event_json(&line) {
325 ne_stdout.lock().await.push(json_str);
326 } else {
327 spinner.set_message(line.clone());
328 if !line.is_empty() && !line.contains("running") && !line.contains("default") {
329 println!("{}", line);
330 }
331 }
332 }
333 spinner.finish_and_clear();
334 });
335
336 let stderr_handle = tokio::spawn(async move {
338 while let Ok(Some(line)) = stderr_lines.next_line().await {
339 if !line.is_empty() {
340 log_stderr.lock().await.push(format!("[stderr] {}", line));
341 if let Some(json_str) = extract_failure_json(&line) {
342 fd_stderr.lock().await.push(json_str);
343 } else if let Some(json_str) = extract_exchange_json(&line) {
344 ex_stderr.lock().await.push(json_str);
345 } else if let Some(json_str) = extract_network_event_json(&line) {
346 ne_stderr.lock().await.push(json_str);
347 } else {
348 eprintln!("{}", line);
349 }
350 }
351 }
352 });
353
354 let status =
356 child.wait().await.map_err(|e| BenchError::K6ExecutionFailed(e.to_string()))?;
357
358 let _ = stdout_handle.await;
360 let _ = stderr_handle.await;
361
362 let exit_code = status.code().unwrap_or(-1);
365 if !status.success() && exit_code != 99 {
366 return Err(BenchError::K6ExecutionFailed(format!(
367 "k6 exited with status: {}",
368 status
369 )));
370 }
371 if exit_code == 99 {
372 tracing::warn!("k6 thresholds crossed (exit code 99) — results will still be parsed");
373 }
374
375 if let Some(dir) = output_dir {
377 let details = failure_details.lock().await;
378 if !details.is_empty() {
379 let failure_path = dir.join("conformance-failure-details.json");
380 let parsed: Vec<serde_json::Value> =
381 details.iter().filter_map(|s| serde_json::from_str(s).ok()).collect();
382 if let Ok(json) = serde_json::to_string_pretty(&parsed) {
383 let _ = std::fs::write(&failure_path, json);
384 }
385 }
386
387 let exchanges = exchange_details.lock().await;
389 if !exchanges.is_empty() {
390 let exchange_path = dir.join("conformance-requests.json");
391 let parsed: Vec<serde_json::Value> =
392 exchanges.iter().filter_map(|s| serde_json::from_str(s).ok()).collect();
393 if let Ok(json) = serde_json::to_string_pretty(&parsed) {
394 let _ = std::fs::write(&exchange_path, json);
395 tracing::info!(
396 "Exported {} request/response pairs to {}",
397 parsed.len(),
398 exchange_path.display()
399 );
400 }
401 }
402
403 let net_events = network_events.lock().await;
408 let net_path = dir.join("conformance-network-events.json");
409 let parsed: Vec<serde_json::Value> =
410 net_events.iter().filter_map(|s| serde_json::from_str(s).ok()).collect();
411 if let Ok(json) = serde_json::to_string_pretty(&parsed) {
412 let _ = std::fs::write(&net_path, json);
413 if !parsed.is_empty() {
414 tracing::warn!(
415 "Recorded {} wire-level network event(s) to {}",
416 parsed.len(),
417 net_path.display()
418 );
419 }
420 }
421
422 let lines = log_lines.lock().await;
424 if !lines.is_empty() {
425 let log_path = dir.join("k6-output.log");
426 let _ = std::fs::write(&log_path, lines.join("\n"));
427 println!("k6 output log saved to: {}", log_path.display());
428 }
429 }
430
431 let results = if let Some(dir) = output_dir {
433 Self::parse_results(dir)?
434 } else {
435 K6Results::default()
436 };
437
438 Ok(results)
439 }
440
441 fn parse_results(output_dir: &Path) -> Result<K6Results> {
443 let summary_path = output_dir.join("summary.json");
444
445 if !summary_path.exists() {
446 return Ok(K6Results::default());
447 }
448
449 let content = std::fs::read_to_string(summary_path)
450 .map_err(|e| BenchError::ResultsParseError(e.to_string()))?;
451
452 let json: serde_json::Value = serde_json::from_str(&content)
453 .map_err(|e| BenchError::ResultsParseError(e.to_string()))?;
454
455 let duration_values = &json["metrics"]["http_req_duration"]["values"];
456
457 let server_latency = &json["metrics"]["mockforge_server_injected_latency_ms"]["values"];
458 let server_jitter = &json["metrics"]["mockforge_server_injected_jitter_ms"]["values"];
459 let server_fault = &json["metrics"]["mockforge_server_fault_total"]["values"]["count"];
460
461 let tcp_connecting = &json["metrics"]["http_req_connecting"]["values"];
472 let tls_handshake = &json["metrics"]["http_req_tls_handshaking"]["values"];
473 let mf_conns_opened = &json["metrics"]["mockforge_connections_opened"]["values"]["count"];
474
475 Ok(K6Results {
476 total_requests: json["metrics"]["http_reqs"]["values"]["count"].as_u64().unwrap_or(0),
477 failed_requests: json["metrics"]["http_req_failed"]["values"]["passes"]
481 .as_u64()
482 .unwrap_or(0),
483 avg_duration_ms: duration_values["avg"].as_f64().unwrap_or(0.0),
484 p95_duration_ms: duration_values["p(95)"].as_f64().unwrap_or(0.0),
485 p99_duration_ms: duration_values["p(99)"].as_f64().unwrap_or(0.0),
486 rps: json["metrics"]["http_reqs"]["values"]["rate"].as_f64().unwrap_or(0.0),
487 vus_max: json["metrics"]["vus_max"]["values"]["value"].as_u64().unwrap_or(0) as u32,
488 min_duration_ms: duration_values["min"].as_f64().unwrap_or(0.0),
489 max_duration_ms: duration_values["max"].as_f64().unwrap_or(0.0),
490 med_duration_ms: duration_values["med"].as_f64().unwrap_or(0.0),
491 p90_duration_ms: duration_values["p(90)"].as_f64().unwrap_or(0.0),
492 server_injected_latency_samples: server_latency["count"].as_u64().unwrap_or(0),
493 server_injected_latency_avg_ms: server_latency["avg"].as_f64().unwrap_or(0.0),
494 server_injected_latency_max_ms: server_latency["max"].as_f64().unwrap_or(0.0),
495 server_injected_jitter_samples: server_jitter["count"].as_u64().unwrap_or(0),
496 server_injected_jitter_avg_ms: server_jitter["avg"].as_f64().unwrap_or(0.0),
497 server_reported_faults: server_fault.as_u64().unwrap_or(0),
498 tcp_connect_samples: mf_conns_opened.as_u64().unwrap_or(0),
501 tcp_connect_avg_ms: tcp_connecting["avg"].as_f64().unwrap_or(0.0),
502 tcp_connect_max_ms: tcp_connecting["max"].as_f64().unwrap_or(0.0),
503 tls_handshake_samples: if tls_handshake["avg"].as_f64().unwrap_or(0.0) > 0.0 {
505 mf_conns_opened.as_u64().unwrap_or(0)
508 } else {
509 0
510 },
511 tls_handshake_avg_ms: tls_handshake["avg"].as_f64().unwrap_or(0.0),
512 tls_handshake_max_ms: tls_handshake["max"].as_f64().unwrap_or(0.0),
513 iterations_completed: json["metrics"]["iterations"]["values"]["count"]
514 .as_u64()
515 .unwrap_or(0),
516 })
517 }
518}
519
520impl Default for K6Executor {
521 fn default() -> Self {
522 Self::new().expect("k6 not found")
523 }
524}
525
526#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
528pub struct K6Results {
529 pub total_requests: u64,
530 pub failed_requests: u64,
531 pub avg_duration_ms: f64,
532 pub p95_duration_ms: f64,
533 pub p99_duration_ms: f64,
534 pub rps: f64,
535 pub vus_max: u32,
536 pub min_duration_ms: f64,
537 pub max_duration_ms: f64,
538 pub med_duration_ms: f64,
539 pub p90_duration_ms: f64,
540 #[serde(default)]
545 pub server_injected_latency_samples: u64,
546 #[serde(default)]
547 pub server_injected_latency_avg_ms: f64,
548 #[serde(default)]
549 pub server_injected_latency_max_ms: f64,
550 #[serde(default)]
551 pub server_injected_jitter_samples: u64,
552 #[serde(default)]
553 pub server_injected_jitter_avg_ms: f64,
554 #[serde(default)]
556 pub server_reported_faults: u64,
557 #[serde(default)]
563 pub tcp_connect_samples: u64,
564 #[serde(default)]
565 pub tcp_connect_avg_ms: f64,
566 #[serde(default)]
567 pub tcp_connect_max_ms: f64,
568 #[serde(default)]
571 pub tls_handshake_samples: u64,
572 #[serde(default)]
573 pub tls_handshake_avg_ms: f64,
574 #[serde(default)]
575 pub tls_handshake_max_ms: f64,
576 #[serde(default)]
582 pub iterations_completed: u64,
583}
584
585impl K6Results {
586 pub fn error_rate(&self) -> f64 {
588 if self.total_requests == 0 {
589 return 0.0;
590 }
591 (self.failed_requests as f64 / self.total_requests as f64) * 100.0
592 }
593
594 pub fn success_rate(&self) -> f64 {
596 100.0 - self.error_rate()
597 }
598}
599
600#[cfg(test)]
601mod tests {
602 use super::*;
603
604 #[test]
605 fn test_k6_results_error_rate() {
606 let results = K6Results {
607 total_requests: 100,
608 failed_requests: 5,
609 avg_duration_ms: 100.0,
610 p95_duration_ms: 200.0,
611 p99_duration_ms: 300.0,
612 ..Default::default()
613 };
614
615 assert_eq!(results.error_rate(), 5.0);
616 assert_eq!(results.success_rate(), 95.0);
617 }
618
619 #[test]
620 fn test_k6_results_zero_requests() {
621 let results = K6Results::default();
622 assert_eq!(results.error_rate(), 0.0);
623 }
624
625 #[test]
626 fn discard_response_bodies_defaults_off_and_builder_flips_it() {
627 let exec = K6Executor {
631 k6_path: "k6".to_string(),
632 local_ips: String::new(),
633 discard_response_bodies: false,
634 };
635 assert!(!exec.discard_response_bodies);
636 let exec = exec.with_discard_response_bodies(true);
637 assert!(exec.discard_response_bodies);
638 }
639
640 #[test]
641 fn test_extract_failure_json_raw() {
642 let line = r#"MOCKFORGE_FAILURE:{"check":"test","expected":"status === 200"}"#;
643 let result = extract_failure_json(line).unwrap();
644 let parsed: serde_json::Value = serde_json::from_str(&result).unwrap();
645 assert_eq!(parsed["check"], "test");
646 }
647
648 #[test]
649 fn test_extract_failure_json_logfmt() {
650 let line = r#"time="2026-01-01T00:00:00Z" level=info msg="MOCKFORGE_FAILURE:{\"check\":\"test\",\"response\":{\"body\":\"{\\\"key\\\":\\\"val\\\"}\"}} " source=console"#;
651 let result = extract_failure_json(line).unwrap();
652 let parsed: serde_json::Value = serde_json::from_str(&result).unwrap();
653 assert_eq!(parsed["check"], "test");
654 assert_eq!(parsed["response"]["body"], r#"{"key":"val"}"#);
655 }
656
657 #[test]
658 fn test_extract_failure_json_no_marker() {
659 assert!(extract_failure_json("just a regular log line").is_none());
660 }
661
662 #[test]
668 fn test_extract_exchange_logfmt_with_backslash_escapes() {
669 let line = r#"time="2026-06-26T10:00:00Z" level=info msg="MOCKFORGE_EXCHANGE:{\"check\":\"u\",\"request\":{\"body\":\"--bnd\\r\\n\\u001a\"}}" source=console"#;
673 let result = extract_exchange_json(line).unwrap();
674 let parsed: serde_json::Value = serde_json::from_str(&result).unwrap();
675 assert_eq!(parsed["check"], "u");
676 assert_eq!(parsed["request"]["body"], "--bnd\r\n\u{001a}");
679 }
680
681 #[test]
682 fn test_extract_exchange_raw_no_logfmt_wrapping() {
683 let line =
684 r#"MOCKFORGE_EXCHANGE:{"check":"x","request":{"body":""},"response":{"status":200}}"#;
685 let result = extract_exchange_json(line).unwrap();
686 let parsed: serde_json::Value = serde_json::from_str(&result).unwrap();
687 assert_eq!(parsed["check"], "x");
688 assert_eq!(parsed["response"]["status"], 200);
689 }
690
691 #[test]
695 fn test_extract_exchange_logfmt_tolerates_extra_trailing_fields() {
696 let line = r#"msg="MOCKFORGE_EXCHANGE:{\"check\":\"t\"}" source=console vu=1 iter=0"#;
697 let result = extract_exchange_json(line).unwrap();
698 let parsed: serde_json::Value = serde_json::from_str(&result).unwrap();
699 assert_eq!(parsed["check"], "t");
700 }
701
702 #[test]
706 fn test_extract_exchange_double_backslash_followed_by_quote() {
707 let line = r#"msg="MOCKFORGE_EXCHANGE:{\"k\":\"a\\\\\\\"x\\\"\"}" source=console"#;
710 let result = extract_exchange_json(line).unwrap();
711 let parsed: serde_json::Value = serde_json::from_str(&result).unwrap();
712 assert_eq!(parsed["k"], r#"a\"x""#);
713 }
714}