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 dns_policy: String,
142}
143
144impl K6Executor {
145 pub fn new() -> Result<Self> {
147 let k6_path = which::which("k6")
148 .map_err(|_| BenchError::K6NotFound)?
149 .to_string_lossy()
150 .to_string();
151
152 Ok(Self {
153 k6_path,
154 local_ips: String::new(),
155 discard_response_bodies: false,
156 dns_policy: String::new(),
157 })
158 }
159
160 pub fn with_local_ips(mut self, local_ips: impl Into<String>) -> Self {
164 self.local_ips = local_ips.into();
165 self
166 }
167
168 pub fn with_discard_response_bodies(mut self, discard: bool) -> Self {
172 self.discard_response_bodies = discard;
173 self
174 }
175
176 pub fn with_dns_policy(mut self, policy: impl Into<String>) -> Self {
180 self.dns_policy = policy.into();
181 self
182 }
183
184 pub fn is_k6_installed() -> bool {
186 which::which("k6").is_ok()
187 }
188
189 pub async fn get_version(&self) -> Result<String> {
191 let output = TokioCommand::new(&self.k6_path)
192 .arg("version")
193 .output()
194 .await
195 .map_err(|e| BenchError::K6ExecutionFailed(e.to_string()))?;
196
197 if !output.status.success() {
198 return Err(BenchError::K6ExecutionFailed("Failed to get k6 version".to_string()));
199 }
200
201 Ok(String::from_utf8_lossy(&output.stdout).trim().to_string())
202 }
203
204 pub async fn execute(
214 &self,
215 script_path: &Path,
216 output_dir: Option<&Path>,
217 verbose: bool,
218 ) -> Result<K6Results> {
219 self.execute_with_port(script_path, output_dir, verbose, None).await
220 }
221
222 pub async fn execute_with_port(
224 &self,
225 script_path: &Path,
226 output_dir: Option<&Path>,
227 verbose: bool,
228 api_port: Option<u16>,
229 ) -> Result<K6Results> {
230 println!("Starting load test...\n");
231
232 let mut cmd = TokioCommand::new(&self.k6_path);
233 cmd.arg("run");
234
235 if let Some(port) = api_port {
238 cmd.arg("--address").arg(format!("localhost:{}", port));
239 }
240
241 if !self.local_ips.is_empty() {
246 cmd.arg("--local-ips").arg(&self.local_ips);
247 }
248
249 if self.discard_response_bodies {
252 cmd.env("K6_DISCARD_RESPONSE_BODIES", "true");
253 }
254
255 if !self.dns_policy.is_empty() {
260 cmd.arg("--dns").arg(format!("policy={}", self.dns_policy));
261 }
262
263 if verbose {
270 cmd.arg("--verbose");
271 }
272
273 let abs_script =
275 std::fs::canonicalize(script_path).unwrap_or_else(|_| script_path.to_path_buf());
276 cmd.arg(&abs_script);
277
278 if let Some(dir) = output_dir {
281 cmd.current_dir(dir);
282 }
283
284 cmd.stdout(Stdio::piped());
285 cmd.stderr(Stdio::piped());
286
287 let mut child = cmd.spawn().map_err(|e| BenchError::K6ExecutionFailed(e.to_string()))?;
288
289 let stdout = child
290 .stdout
291 .take()
292 .ok_or_else(|| BenchError::K6ExecutionFailed("Failed to capture stdout".to_string()))?;
293
294 let stderr = child
295 .stderr
296 .take()
297 .ok_or_else(|| BenchError::K6ExecutionFailed("Failed to capture stderr".to_string()))?;
298
299 let stdout_reader = BufReader::new(stdout);
301 let stderr_reader = BufReader::new(stderr);
302
303 let mut stdout_lines = stdout_reader.lines();
304 let mut stderr_lines = stderr_reader.lines();
305
306 let spinner = ProgressBar::new_spinner();
308 spinner.set_style(
309 ProgressStyle::default_spinner().template("{spinner:.green} {msg}").unwrap(),
310 );
311 spinner.set_message("Running load test...");
312
313 let failure_details: Arc<tokio::sync::Mutex<Vec<String>>> =
316 Arc::new(tokio::sync::Mutex::new(Vec::new()));
317 let fd_stdout = Arc::clone(&failure_details);
318 let fd_stderr = Arc::clone(&failure_details);
319
320 let exchange_details: Arc<tokio::sync::Mutex<Vec<String>>> =
322 Arc::new(tokio::sync::Mutex::new(Vec::new()));
323 let ex_stdout = Arc::clone(&exchange_details);
324 let ex_stderr = Arc::clone(&exchange_details);
325
326 let network_events: Arc<tokio::sync::Mutex<Vec<String>>> =
331 Arc::new(tokio::sync::Mutex::new(Vec::new()));
332 let ne_stdout = Arc::clone(&network_events);
333 let ne_stderr = Arc::clone(&network_events);
334
335 let log_lines: Arc<tokio::sync::Mutex<Vec<String>>> =
337 Arc::new(tokio::sync::Mutex::new(Vec::new()));
338 let log_stdout = Arc::clone(&log_lines);
339 let log_stderr = Arc::clone(&log_lines);
340
341 let stdout_handle = tokio::spawn(async move {
343 while let Ok(Some(line)) = stdout_lines.next_line().await {
344 log_stdout.lock().await.push(format!("[stdout] {}", line));
345 if let Some(json_str) = extract_failure_json(&line) {
346 fd_stdout.lock().await.push(json_str);
347 } else if let Some(json_str) = extract_exchange_json(&line) {
348 ex_stdout.lock().await.push(json_str);
349 } else if let Some(json_str) = extract_network_event_json(&line) {
350 ne_stdout.lock().await.push(json_str);
351 } else {
352 spinner.set_message(line.clone());
353 if !line.is_empty() && !line.contains("running") && !line.contains("default") {
354 println!("{}", line);
355 }
356 }
357 }
358 spinner.finish_and_clear();
359 });
360
361 let stderr_handle = tokio::spawn(async move {
363 while let Ok(Some(line)) = stderr_lines.next_line().await {
364 if !line.is_empty() {
365 log_stderr.lock().await.push(format!("[stderr] {}", line));
366 if let Some(json_str) = extract_failure_json(&line) {
367 fd_stderr.lock().await.push(json_str);
368 } else if let Some(json_str) = extract_exchange_json(&line) {
369 ex_stderr.lock().await.push(json_str);
370 } else if let Some(json_str) = extract_network_event_json(&line) {
371 ne_stderr.lock().await.push(json_str);
372 } else {
373 eprintln!("{}", line);
374 }
375 }
376 }
377 });
378
379 let status =
381 child.wait().await.map_err(|e| BenchError::K6ExecutionFailed(e.to_string()))?;
382
383 let _ = stdout_handle.await;
385 let _ = stderr_handle.await;
386
387 let exit_code = status.code().unwrap_or(-1);
390 if !status.success() && exit_code != 99 {
391 return Err(BenchError::K6ExecutionFailed(format!(
392 "k6 exited with status: {}",
393 status
394 )));
395 }
396 if exit_code == 99 {
397 tracing::warn!("k6 thresholds crossed (exit code 99) — results will still be parsed");
398 }
399
400 if let Some(dir) = output_dir {
402 let details = failure_details.lock().await;
403 if !details.is_empty() {
404 let failure_path = dir.join("conformance-failure-details.json");
405 let parsed: Vec<serde_json::Value> =
406 details.iter().filter_map(|s| serde_json::from_str(s).ok()).collect();
407 if let Ok(json) = serde_json::to_string_pretty(&parsed) {
408 let _ = std::fs::write(&failure_path, json);
409 }
410 }
411
412 let exchanges = exchange_details.lock().await;
414 if !exchanges.is_empty() {
415 let exchange_path = dir.join("conformance-requests.json");
416 let parsed: Vec<serde_json::Value> =
417 exchanges.iter().filter_map(|s| serde_json::from_str(s).ok()).collect();
418 if let Ok(json) = serde_json::to_string_pretty(&parsed) {
419 let _ = std::fs::write(&exchange_path, json);
420 tracing::info!(
421 "Exported {} request/response pairs to {}",
422 parsed.len(),
423 exchange_path.display()
424 );
425 }
426 }
427
428 let net_events = network_events.lock().await;
433 let net_path = dir.join("conformance-network-events.json");
434 let parsed: Vec<serde_json::Value> =
435 net_events.iter().filter_map(|s| serde_json::from_str(s).ok()).collect();
436 if let Ok(json) = serde_json::to_string_pretty(&parsed) {
437 let _ = std::fs::write(&net_path, json);
438 if !parsed.is_empty() {
439 tracing::warn!(
440 "Recorded {} wire-level network event(s) to {}",
441 parsed.len(),
442 net_path.display()
443 );
444 }
445 }
446
447 let lines = log_lines.lock().await;
449 if !lines.is_empty() {
450 let log_path = dir.join("k6-output.log");
451 let _ = std::fs::write(&log_path, lines.join("\n"));
452 println!("k6 output log saved to: {}", log_path.display());
453 }
454 }
455
456 let results = if let Some(dir) = output_dir {
458 Self::parse_results(dir)?
459 } else {
460 K6Results::default()
461 };
462
463 Ok(results)
464 }
465
466 fn parse_results(output_dir: &Path) -> Result<K6Results> {
468 let summary_path = output_dir.join("summary.json");
469
470 if !summary_path.exists() {
471 return Ok(K6Results::default());
472 }
473
474 let content = std::fs::read_to_string(summary_path)
475 .map_err(|e| BenchError::ResultsParseError(e.to_string()))?;
476
477 let json: serde_json::Value = serde_json::from_str(&content)
478 .map_err(|e| BenchError::ResultsParseError(e.to_string()))?;
479
480 let duration_values = &json["metrics"]["http_req_duration"]["values"];
481
482 let server_latency = &json["metrics"]["mockforge_server_injected_latency_ms"]["values"];
483 let server_jitter = &json["metrics"]["mockforge_server_injected_jitter_ms"]["values"];
484 let server_fault = &json["metrics"]["mockforge_server_fault_total"]["values"]["count"];
485
486 let tcp_connecting = &json["metrics"]["http_req_connecting"]["values"];
497 let tls_handshake = &json["metrics"]["http_req_tls_handshaking"]["values"];
498 let mf_conns_opened = &json["metrics"]["mockforge_connections_opened"]["values"]["count"];
499
500 Ok(K6Results {
501 total_requests: json["metrics"]["http_reqs"]["values"]["count"].as_u64().unwrap_or(0),
502 failed_requests: json["metrics"]["http_req_failed"]["values"]["passes"]
506 .as_u64()
507 .unwrap_or(0),
508 avg_duration_ms: duration_values["avg"].as_f64().unwrap_or(0.0),
509 p95_duration_ms: duration_values["p(95)"].as_f64().unwrap_or(0.0),
510 p99_duration_ms: duration_values["p(99)"].as_f64().unwrap_or(0.0),
511 rps: json["metrics"]["http_reqs"]["values"]["rate"].as_f64().unwrap_or(0.0),
512 vus_max: json["metrics"]["vus_max"]["values"]["value"].as_u64().unwrap_or(0) as u32,
513 min_duration_ms: duration_values["min"].as_f64().unwrap_or(0.0),
514 max_duration_ms: duration_values["max"].as_f64().unwrap_or(0.0),
515 med_duration_ms: duration_values["med"].as_f64().unwrap_or(0.0),
516 p90_duration_ms: duration_values["p(90)"].as_f64().unwrap_or(0.0),
517 server_injected_latency_samples: server_latency["count"].as_u64().unwrap_or(0),
518 server_injected_latency_avg_ms: server_latency["avg"].as_f64().unwrap_or(0.0),
519 server_injected_latency_max_ms: server_latency["max"].as_f64().unwrap_or(0.0),
520 server_injected_jitter_samples: server_jitter["count"].as_u64().unwrap_or(0),
521 server_injected_jitter_avg_ms: server_jitter["avg"].as_f64().unwrap_or(0.0),
522 server_reported_faults: server_fault.as_u64().unwrap_or(0),
523 tcp_connect_samples: mf_conns_opened.as_u64().unwrap_or(0),
526 tcp_connect_avg_ms: tcp_connecting["avg"].as_f64().unwrap_or(0.0),
527 tcp_connect_max_ms: tcp_connecting["max"].as_f64().unwrap_or(0.0),
528 tls_handshake_samples: if tls_handshake["avg"].as_f64().unwrap_or(0.0) > 0.0 {
530 mf_conns_opened.as_u64().unwrap_or(0)
533 } else {
534 0
535 },
536 tls_handshake_avg_ms: tls_handshake["avg"].as_f64().unwrap_or(0.0),
537 tls_handshake_max_ms: tls_handshake["max"].as_f64().unwrap_or(0.0),
538 iterations_completed: json["metrics"]["iterations"]["values"]["count"]
539 .as_u64()
540 .unwrap_or(0),
541 })
542 }
543}
544
545impl Default for K6Executor {
546 fn default() -> Self {
547 Self::new().expect("k6 not found")
548 }
549}
550
551#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
553pub struct K6Results {
554 pub total_requests: u64,
555 pub failed_requests: u64,
556 pub avg_duration_ms: f64,
557 pub p95_duration_ms: f64,
558 pub p99_duration_ms: f64,
559 pub rps: f64,
560 pub vus_max: u32,
561 pub min_duration_ms: f64,
562 pub max_duration_ms: f64,
563 pub med_duration_ms: f64,
564 pub p90_duration_ms: f64,
565 #[serde(default)]
570 pub server_injected_latency_samples: u64,
571 #[serde(default)]
572 pub server_injected_latency_avg_ms: f64,
573 #[serde(default)]
574 pub server_injected_latency_max_ms: f64,
575 #[serde(default)]
576 pub server_injected_jitter_samples: u64,
577 #[serde(default)]
578 pub server_injected_jitter_avg_ms: f64,
579 #[serde(default)]
581 pub server_reported_faults: u64,
582 #[serde(default)]
588 pub tcp_connect_samples: u64,
589 #[serde(default)]
590 pub tcp_connect_avg_ms: f64,
591 #[serde(default)]
592 pub tcp_connect_max_ms: f64,
593 #[serde(default)]
596 pub tls_handshake_samples: u64,
597 #[serde(default)]
598 pub tls_handshake_avg_ms: f64,
599 #[serde(default)]
600 pub tls_handshake_max_ms: f64,
601 #[serde(default)]
607 pub iterations_completed: u64,
608}
609
610impl K6Results {
611 pub fn error_rate(&self) -> f64 {
613 if self.total_requests == 0 {
614 return 0.0;
615 }
616 (self.failed_requests as f64 / self.total_requests as f64) * 100.0
617 }
618
619 pub fn success_rate(&self) -> f64 {
621 100.0 - self.error_rate()
622 }
623}
624
625#[cfg(test)]
626mod tests {
627 use super::*;
628
629 #[test]
630 fn test_k6_results_error_rate() {
631 let results = K6Results {
632 total_requests: 100,
633 failed_requests: 5,
634 avg_duration_ms: 100.0,
635 p95_duration_ms: 200.0,
636 p99_duration_ms: 300.0,
637 ..Default::default()
638 };
639
640 assert_eq!(results.error_rate(), 5.0);
641 assert_eq!(results.success_rate(), 95.0);
642 }
643
644 #[test]
645 fn test_k6_results_zero_requests() {
646 let results = K6Results::default();
647 assert_eq!(results.error_rate(), 0.0);
648 }
649
650 #[test]
651 fn discard_response_bodies_defaults_off_and_builder_flips_it() {
652 let exec = K6Executor {
656 k6_path: "k6".to_string(),
657 local_ips: String::new(),
658 discard_response_bodies: false,
659 dns_policy: String::new(),
660 };
661 assert!(!exec.discard_response_bodies);
662 let exec = exec.with_discard_response_bodies(true);
663 assert!(exec.discard_response_bodies);
664 }
665
666 #[test]
667 fn dns_policy_defaults_empty_and_builder_sets_it() {
668 let exec = K6Executor {
672 k6_path: "k6".to_string(),
673 local_ips: String::new(),
674 discard_response_bodies: false,
675 dns_policy: String::new(),
676 };
677 assert!(exec.dns_policy.is_empty());
678 let exec = exec.with_dns_policy("preferIPv6");
679 assert_eq!(exec.dns_policy, "preferIPv6");
680 }
681
682 #[test]
683 fn test_extract_failure_json_raw() {
684 let line = r#"MOCKFORGE_FAILURE:{"check":"test","expected":"status === 200"}"#;
685 let result = extract_failure_json(line).unwrap();
686 let parsed: serde_json::Value = serde_json::from_str(&result).unwrap();
687 assert_eq!(parsed["check"], "test");
688 }
689
690 #[test]
691 fn test_extract_failure_json_logfmt() {
692 let line = r#"time="2026-01-01T00:00:00Z" level=info msg="MOCKFORGE_FAILURE:{\"check\":\"test\",\"response\":{\"body\":\"{\\\"key\\\":\\\"val\\\"}\"}} " source=console"#;
693 let result = extract_failure_json(line).unwrap();
694 let parsed: serde_json::Value = serde_json::from_str(&result).unwrap();
695 assert_eq!(parsed["check"], "test");
696 assert_eq!(parsed["response"]["body"], r#"{"key":"val"}"#);
697 }
698
699 #[test]
700 fn test_extract_failure_json_no_marker() {
701 assert!(extract_failure_json("just a regular log line").is_none());
702 }
703
704 #[test]
710 fn test_extract_exchange_logfmt_with_backslash_escapes() {
711 let line = r#"time="2026-06-26T10:00:00Z" level=info msg="MOCKFORGE_EXCHANGE:{\"check\":\"u\",\"request\":{\"body\":\"--bnd\\r\\n\\u001a\"}}" source=console"#;
715 let result = extract_exchange_json(line).unwrap();
716 let parsed: serde_json::Value = serde_json::from_str(&result).unwrap();
717 assert_eq!(parsed["check"], "u");
718 assert_eq!(parsed["request"]["body"], "--bnd\r\n\u{001a}");
721 }
722
723 #[test]
724 fn test_extract_exchange_raw_no_logfmt_wrapping() {
725 let line =
726 r#"MOCKFORGE_EXCHANGE:{"check":"x","request":{"body":""},"response":{"status":200}}"#;
727 let result = extract_exchange_json(line).unwrap();
728 let parsed: serde_json::Value = serde_json::from_str(&result).unwrap();
729 assert_eq!(parsed["check"], "x");
730 assert_eq!(parsed["response"]["status"], 200);
731 }
732
733 #[test]
737 fn test_extract_exchange_logfmt_tolerates_extra_trailing_fields() {
738 let line = r#"msg="MOCKFORGE_EXCHANGE:{\"check\":\"t\"}" source=console vu=1 iter=0"#;
739 let result = extract_exchange_json(line).unwrap();
740 let parsed: serde_json::Value = serde_json::from_str(&result).unwrap();
741 assert_eq!(parsed["check"], "t");
742 }
743
744 #[test]
748 fn test_extract_exchange_double_backslash_followed_by_quote() {
749 let line = r#"msg="MOCKFORGE_EXCHANGE:{\"k\":\"a\\\\\\\"x\\\"\"}" source=console"#;
752 let result = extract_exchange_json(line).unwrap();
753 let parsed: serde_json::Value = serde_json::from_str(&result).unwrap();
754 assert_eq!(parsed["k"], r#"a\"x""#);
755 }
756}