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(
194 &self,
195 script_path: &Path,
196 output_dir: Option<&Path>,
197 verbose: bool,
198 ) -> Result<K6Results> {
199 self.execute_with_port(script_path, output_dir, verbose, None).await
200 }
201
202 pub async fn execute_with_port(
204 &self,
205 script_path: &Path,
206 output_dir: Option<&Path>,
207 verbose: bool,
208 api_port: Option<u16>,
209 ) -> Result<K6Results> {
210 println!("Starting load test...\n");
211
212 let mut cmd = TokioCommand::new(&self.k6_path);
213 cmd.arg("run");
214
215 if let Some(port) = api_port {
218 cmd.arg("--address").arg(format!("localhost:{}", port));
219 }
220
221 if !self.local_ips.is_empty() {
226 cmd.arg("--local-ips").arg(&self.local_ips);
227 }
228
229 if self.discard_response_bodies {
232 cmd.env("K6_DISCARD_RESPONSE_BODIES", "true");
233 }
234
235 if verbose {
242 cmd.arg("--verbose");
243 }
244
245 let abs_script =
247 std::fs::canonicalize(script_path).unwrap_or_else(|_| script_path.to_path_buf());
248 cmd.arg(&abs_script);
249
250 if let Some(dir) = output_dir {
253 cmd.current_dir(dir);
254 }
255
256 cmd.stdout(Stdio::piped());
257 cmd.stderr(Stdio::piped());
258
259 let mut child = cmd.spawn().map_err(|e| BenchError::K6ExecutionFailed(e.to_string()))?;
260
261 let stdout = child
262 .stdout
263 .take()
264 .ok_or_else(|| BenchError::K6ExecutionFailed("Failed to capture stdout".to_string()))?;
265
266 let stderr = child
267 .stderr
268 .take()
269 .ok_or_else(|| BenchError::K6ExecutionFailed("Failed to capture stderr".to_string()))?;
270
271 let stdout_reader = BufReader::new(stdout);
273 let stderr_reader = BufReader::new(stderr);
274
275 let mut stdout_lines = stdout_reader.lines();
276 let mut stderr_lines = stderr_reader.lines();
277
278 let spinner = ProgressBar::new_spinner();
280 spinner.set_style(
281 ProgressStyle::default_spinner().template("{spinner:.green} {msg}").unwrap(),
282 );
283 spinner.set_message("Running load test...");
284
285 let failure_details: Arc<tokio::sync::Mutex<Vec<String>>> =
288 Arc::new(tokio::sync::Mutex::new(Vec::new()));
289 let fd_stdout = Arc::clone(&failure_details);
290 let fd_stderr = Arc::clone(&failure_details);
291
292 let exchange_details: Arc<tokio::sync::Mutex<Vec<String>>> =
294 Arc::new(tokio::sync::Mutex::new(Vec::new()));
295 let ex_stdout = Arc::clone(&exchange_details);
296 let ex_stderr = Arc::clone(&exchange_details);
297
298 let network_events: Arc<tokio::sync::Mutex<Vec<String>>> =
303 Arc::new(tokio::sync::Mutex::new(Vec::new()));
304 let ne_stdout = Arc::clone(&network_events);
305 let ne_stderr = Arc::clone(&network_events);
306
307 let log_lines: Arc<tokio::sync::Mutex<Vec<String>>> =
309 Arc::new(tokio::sync::Mutex::new(Vec::new()));
310 let log_stdout = Arc::clone(&log_lines);
311 let log_stderr = Arc::clone(&log_lines);
312
313 let stdout_handle = tokio::spawn(async move {
315 while let Ok(Some(line)) = stdout_lines.next_line().await {
316 log_stdout.lock().await.push(format!("[stdout] {}", line));
317 if let Some(json_str) = extract_failure_json(&line) {
318 fd_stdout.lock().await.push(json_str);
319 } else if let Some(json_str) = extract_exchange_json(&line) {
320 ex_stdout.lock().await.push(json_str);
321 } else if let Some(json_str) = extract_network_event_json(&line) {
322 ne_stdout.lock().await.push(json_str);
323 } else {
324 spinner.set_message(line.clone());
325 if !line.is_empty() && !line.contains("running") && !line.contains("default") {
326 println!("{}", line);
327 }
328 }
329 }
330 spinner.finish_and_clear();
331 });
332
333 let stderr_handle = tokio::spawn(async move {
335 while let Ok(Some(line)) = stderr_lines.next_line().await {
336 if !line.is_empty() {
337 log_stderr.lock().await.push(format!("[stderr] {}", line));
338 if let Some(json_str) = extract_failure_json(&line) {
339 fd_stderr.lock().await.push(json_str);
340 } else if let Some(json_str) = extract_exchange_json(&line) {
341 ex_stderr.lock().await.push(json_str);
342 } else if let Some(json_str) = extract_network_event_json(&line) {
343 ne_stderr.lock().await.push(json_str);
344 } else {
345 eprintln!("{}", line);
346 }
347 }
348 }
349 });
350
351 let status =
353 child.wait().await.map_err(|e| BenchError::K6ExecutionFailed(e.to_string()))?;
354
355 let _ = stdout_handle.await;
357 let _ = stderr_handle.await;
358
359 let exit_code = status.code().unwrap_or(-1);
362 if !status.success() && exit_code != 99 {
363 return Err(BenchError::K6ExecutionFailed(format!(
364 "k6 exited with status: {}",
365 status
366 )));
367 }
368 if exit_code == 99 {
369 tracing::warn!("k6 thresholds crossed (exit code 99) — results will still be parsed");
370 }
371
372 if let Some(dir) = output_dir {
374 let details = failure_details.lock().await;
375 if !details.is_empty() {
376 let failure_path = dir.join("conformance-failure-details.json");
377 let parsed: Vec<serde_json::Value> =
378 details.iter().filter_map(|s| serde_json::from_str(s).ok()).collect();
379 if let Ok(json) = serde_json::to_string_pretty(&parsed) {
380 let _ = std::fs::write(&failure_path, json);
381 }
382 }
383
384 let exchanges = exchange_details.lock().await;
386 if !exchanges.is_empty() {
387 let exchange_path = dir.join("conformance-requests.json");
388 let parsed: Vec<serde_json::Value> =
389 exchanges.iter().filter_map(|s| serde_json::from_str(s).ok()).collect();
390 if let Ok(json) = serde_json::to_string_pretty(&parsed) {
391 let _ = std::fs::write(&exchange_path, json);
392 tracing::info!(
393 "Exported {} request/response pairs to {}",
394 parsed.len(),
395 exchange_path.display()
396 );
397 }
398 }
399
400 let net_events = network_events.lock().await;
405 let net_path = dir.join("conformance-network-events.json");
406 let parsed: Vec<serde_json::Value> =
407 net_events.iter().filter_map(|s| serde_json::from_str(s).ok()).collect();
408 if let Ok(json) = serde_json::to_string_pretty(&parsed) {
409 let _ = std::fs::write(&net_path, json);
410 if !parsed.is_empty() {
411 tracing::warn!(
412 "Recorded {} wire-level network event(s) to {}",
413 parsed.len(),
414 net_path.display()
415 );
416 }
417 }
418
419 let lines = log_lines.lock().await;
421 if !lines.is_empty() {
422 let log_path = dir.join("k6-output.log");
423 let _ = std::fs::write(&log_path, lines.join("\n"));
424 println!("k6 output log saved to: {}", log_path.display());
425 }
426 }
427
428 let results = if let Some(dir) = output_dir {
430 Self::parse_results(dir)?
431 } else {
432 K6Results::default()
433 };
434
435 Ok(results)
436 }
437
438 fn parse_results(output_dir: &Path) -> Result<K6Results> {
440 let summary_path = output_dir.join("summary.json");
441
442 if !summary_path.exists() {
443 return Ok(K6Results::default());
444 }
445
446 let content = std::fs::read_to_string(summary_path)
447 .map_err(|e| BenchError::ResultsParseError(e.to_string()))?;
448
449 let json: serde_json::Value = serde_json::from_str(&content)
450 .map_err(|e| BenchError::ResultsParseError(e.to_string()))?;
451
452 let duration_values = &json["metrics"]["http_req_duration"]["values"];
453
454 let server_latency = &json["metrics"]["mockforge_server_injected_latency_ms"]["values"];
455 let server_jitter = &json["metrics"]["mockforge_server_injected_jitter_ms"]["values"];
456 let server_fault = &json["metrics"]["mockforge_server_fault_total"]["values"]["count"];
457
458 let tcp_connecting = &json["metrics"]["http_req_connecting"]["values"];
469 let tls_handshake = &json["metrics"]["http_req_tls_handshaking"]["values"];
470 let mf_conns_opened = &json["metrics"]["mockforge_connections_opened"]["values"]["count"];
471
472 Ok(K6Results {
473 total_requests: json["metrics"]["http_reqs"]["values"]["count"].as_u64().unwrap_or(0),
474 failed_requests: json["metrics"]["http_req_failed"]["values"]["passes"]
478 .as_u64()
479 .unwrap_or(0),
480 avg_duration_ms: duration_values["avg"].as_f64().unwrap_or(0.0),
481 p95_duration_ms: duration_values["p(95)"].as_f64().unwrap_or(0.0),
482 p99_duration_ms: duration_values["p(99)"].as_f64().unwrap_or(0.0),
483 rps: json["metrics"]["http_reqs"]["values"]["rate"].as_f64().unwrap_or(0.0),
484 vus_max: json["metrics"]["vus_max"]["values"]["value"].as_u64().unwrap_or(0) as u32,
485 min_duration_ms: duration_values["min"].as_f64().unwrap_or(0.0),
486 max_duration_ms: duration_values["max"].as_f64().unwrap_or(0.0),
487 med_duration_ms: duration_values["med"].as_f64().unwrap_or(0.0),
488 p90_duration_ms: duration_values["p(90)"].as_f64().unwrap_or(0.0),
489 server_injected_latency_samples: server_latency["count"].as_u64().unwrap_or(0),
490 server_injected_latency_avg_ms: server_latency["avg"].as_f64().unwrap_or(0.0),
491 server_injected_latency_max_ms: server_latency["max"].as_f64().unwrap_or(0.0),
492 server_injected_jitter_samples: server_jitter["count"].as_u64().unwrap_or(0),
493 server_injected_jitter_avg_ms: server_jitter["avg"].as_f64().unwrap_or(0.0),
494 server_reported_faults: server_fault.as_u64().unwrap_or(0),
495 tcp_connect_samples: mf_conns_opened.as_u64().unwrap_or(0),
498 tcp_connect_avg_ms: tcp_connecting["avg"].as_f64().unwrap_or(0.0),
499 tcp_connect_max_ms: tcp_connecting["max"].as_f64().unwrap_or(0.0),
500 tls_handshake_samples: if tls_handshake["avg"].as_f64().unwrap_or(0.0) > 0.0 {
502 mf_conns_opened.as_u64().unwrap_or(0)
505 } else {
506 0
507 },
508 tls_handshake_avg_ms: tls_handshake["avg"].as_f64().unwrap_or(0.0),
509 tls_handshake_max_ms: tls_handshake["max"].as_f64().unwrap_or(0.0),
510 iterations_completed: json["metrics"]["iterations"]["values"]["count"]
511 .as_u64()
512 .unwrap_or(0),
513 })
514 }
515}
516
517impl Default for K6Executor {
518 fn default() -> Self {
519 Self::new().expect("k6 not found")
520 }
521}
522
523#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
525pub struct K6Results {
526 pub total_requests: u64,
527 pub failed_requests: u64,
528 pub avg_duration_ms: f64,
529 pub p95_duration_ms: f64,
530 pub p99_duration_ms: f64,
531 pub rps: f64,
532 pub vus_max: u32,
533 pub min_duration_ms: f64,
534 pub max_duration_ms: f64,
535 pub med_duration_ms: f64,
536 pub p90_duration_ms: f64,
537 #[serde(default)]
542 pub server_injected_latency_samples: u64,
543 #[serde(default)]
544 pub server_injected_latency_avg_ms: f64,
545 #[serde(default)]
546 pub server_injected_latency_max_ms: f64,
547 #[serde(default)]
548 pub server_injected_jitter_samples: u64,
549 #[serde(default)]
550 pub server_injected_jitter_avg_ms: f64,
551 #[serde(default)]
553 pub server_reported_faults: u64,
554 #[serde(default)]
560 pub tcp_connect_samples: u64,
561 #[serde(default)]
562 pub tcp_connect_avg_ms: f64,
563 #[serde(default)]
564 pub tcp_connect_max_ms: f64,
565 #[serde(default)]
568 pub tls_handshake_samples: u64,
569 #[serde(default)]
570 pub tls_handshake_avg_ms: f64,
571 #[serde(default)]
572 pub tls_handshake_max_ms: f64,
573 #[serde(default)]
579 pub iterations_completed: u64,
580}
581
582impl K6Results {
583 pub fn error_rate(&self) -> f64 {
585 if self.total_requests == 0 {
586 return 0.0;
587 }
588 (self.failed_requests as f64 / self.total_requests as f64) * 100.0
589 }
590
591 pub fn success_rate(&self) -> f64 {
593 100.0 - self.error_rate()
594 }
595}
596
597#[cfg(test)]
598mod tests {
599 use super::*;
600
601 #[test]
602 fn test_k6_results_error_rate() {
603 let results = K6Results {
604 total_requests: 100,
605 failed_requests: 5,
606 avg_duration_ms: 100.0,
607 p95_duration_ms: 200.0,
608 p99_duration_ms: 300.0,
609 ..Default::default()
610 };
611
612 assert_eq!(results.error_rate(), 5.0);
613 assert_eq!(results.success_rate(), 95.0);
614 }
615
616 #[test]
617 fn test_k6_results_zero_requests() {
618 let results = K6Results::default();
619 assert_eq!(results.error_rate(), 0.0);
620 }
621
622 #[test]
623 fn discard_response_bodies_defaults_off_and_builder_flips_it() {
624 let exec = K6Executor {
628 k6_path: "k6".to_string(),
629 local_ips: String::new(),
630 discard_response_bodies: false,
631 };
632 assert!(!exec.discard_response_bodies);
633 let exec = exec.with_discard_response_bodies(true);
634 assert!(exec.discard_response_bodies);
635 }
636
637 #[test]
638 fn test_extract_failure_json_raw() {
639 let line = r#"MOCKFORGE_FAILURE:{"check":"test","expected":"status === 200"}"#;
640 let result = extract_failure_json(line).unwrap();
641 let parsed: serde_json::Value = serde_json::from_str(&result).unwrap();
642 assert_eq!(parsed["check"], "test");
643 }
644
645 #[test]
646 fn test_extract_failure_json_logfmt() {
647 let line = r#"time="2026-01-01T00:00:00Z" level=info msg="MOCKFORGE_FAILURE:{\"check\":\"test\",\"response\":{\"body\":\"{\\\"key\\\":\\\"val\\\"}\"}} " source=console"#;
648 let result = extract_failure_json(line).unwrap();
649 let parsed: serde_json::Value = serde_json::from_str(&result).unwrap();
650 assert_eq!(parsed["check"], "test");
651 assert_eq!(parsed["response"]["body"], r#"{"key":"val"}"#);
652 }
653
654 #[test]
655 fn test_extract_failure_json_no_marker() {
656 assert!(extract_failure_json("just a regular log line").is_none());
657 }
658
659 #[test]
665 fn test_extract_exchange_logfmt_with_backslash_escapes() {
666 let line = r#"time="2026-06-26T10:00:00Z" level=info msg="MOCKFORGE_EXCHANGE:{\"check\":\"u\",\"request\":{\"body\":\"--bnd\\r\\n\\u001a\"}}" source=console"#;
670 let result = extract_exchange_json(line).unwrap();
671 let parsed: serde_json::Value = serde_json::from_str(&result).unwrap();
672 assert_eq!(parsed["check"], "u");
673 assert_eq!(parsed["request"]["body"], "--bnd\r\n\u{001a}");
676 }
677
678 #[test]
679 fn test_extract_exchange_raw_no_logfmt_wrapping() {
680 let line =
681 r#"MOCKFORGE_EXCHANGE:{"check":"x","request":{"body":""},"response":{"status":200}}"#;
682 let result = extract_exchange_json(line).unwrap();
683 let parsed: serde_json::Value = serde_json::from_str(&result).unwrap();
684 assert_eq!(parsed["check"], "x");
685 assert_eq!(parsed["response"]["status"], 200);
686 }
687
688 #[test]
692 fn test_extract_exchange_logfmt_tolerates_extra_trailing_fields() {
693 let line = r#"msg="MOCKFORGE_EXCHANGE:{\"check\":\"t\"}" source=console vu=1 iter=0"#;
694 let result = extract_exchange_json(line).unwrap();
695 let parsed: serde_json::Value = serde_json::from_str(&result).unwrap();
696 assert_eq!(parsed["check"], "t");
697 }
698
699 #[test]
703 fn test_extract_exchange_double_backslash_followed_by_quote() {
704 let line = r#"msg="MOCKFORGE_EXCHANGE:{\"k\":\"a\\\\\\\"x\\\"\"}" source=console"#;
707 let result = extract_exchange_json(line).unwrap();
708 let parsed: serde_json::Value = serde_json::from_str(&result).unwrap();
709 assert_eq!(parsed["k"], r#"a\"x""#);
710 }
711}