1use reqwest::{Client, Method, Url};
2use serde_json::{Map, Value};
3use std::collections::HashMap;
4use std::future::Future;
5use std::time::Instant;
6
7use crate::assertions::{evaluate_assertions, has_status_assertion};
8use crate::core::types::{
9 PipelineStep, RuntimeEnvGroup, RuntimeSpec, StepExecutionResult, StepRequest, StepResponse,
10};
11use crate::execution::cancel::await_with_cancel;
12use crate::execution::http::{parse_absolute_http_url, parse_method};
13use crate::execution::logging::log_step_response;
14use crate::extractions::evaluate_step_extractions;
15use crate::template::resolve::resolve_template_variables;
16
17#[derive(Debug, Clone)]
18pub struct PreparedHttpStep {
19 pub step_id: String,
20 pub attempt: usize,
21 pub max_attempts: usize,
22 pub method: Method,
23 pub url: Url,
24 pub request: StepRequest,
25 started_at: Instant,
26}
27
28#[derive(Debug)]
29pub struct StartedHttpStep {
30 pub request: StepRequest,
31 pub response: reqwest::Response,
32 started_at: Instant,
33 attempt: usize,
34 max_attempts: usize,
35}
36
37pub fn prepare_http_step(
38 step: &PipelineStep,
39 context: &HashMap<String, StepExecutionResult>,
40 specs: Option<&[RuntimeSpec]>,
41 env_groups: Option<&[RuntimeEnvGroup]>,
42 selected_env_group_slug: Option<&str>,
43 attempt: usize,
44 max_attempts: usize,
45) -> Result<PreparedHttpStep, StepExecutionResult> {
46 let started_at = Instant::now();
47 let resolved_url = resolve_template_variables(
48 &Value::String(step.url.clone()),
49 context,
50 specs,
51 env_groups,
52 selected_env_group_slug,
53 )
54 .as_str()
55 .unwrap_or(step.url.as_str())
56 .to_owned();
57
58 let resolved_headers = resolve_template_variables(
59 &serde_json::to_value(&step.headers).unwrap_or(Value::Object(Map::new())),
60 context,
61 specs,
62 env_groups,
63 selected_env_group_slug,
64 )
65 .as_object()
66 .map(|m| {
67 m.iter()
68 .map(|(k, v)| (k.clone(), v.as_str().unwrap_or_default().to_owned()))
69 .collect::<HashMap<String, String>>()
70 })
71 .unwrap_or_default();
72
73 let resolved_body = step.body.as_ref().map(|body| {
74 resolve_template_variables(body, context, specs, env_groups, selected_env_group_slug)
75 });
76
77 let request = StepRequest {
78 method: step.method.clone(),
79 url: resolved_url.clone(),
80 headers: resolved_headers,
81 body: resolved_body,
82 };
83
84 let method = match parse_method(&step.method) {
85 Ok(method) => method,
86 Err(err) => {
87 return Err(invalid_step_result(
88 step,
89 request,
90 err,
91 started_at,
92 attempt,
93 max_attempts,
94 ));
95 }
96 };
97
98 let url = match parse_absolute_http_url(&resolved_url) {
99 Ok(url) => url,
100 Err(err) => {
101 return Err(invalid_step_result(
102 step,
103 request,
104 err,
105 started_at,
106 attempt,
107 max_attempts,
108 ));
109 }
110 };
111
112 Ok(PreparedHttpStep {
113 step_id: step.id.clone(),
114 attempt,
115 max_attempts,
116 method,
117 url,
118 request,
119 started_at,
120 })
121}
122
123pub async fn send_prepared_http_step<FCancel>(
124 client: &Client,
125 prepared: PreparedHttpStep,
126 step: &PipelineStep,
127 context: &HashMap<String, StepExecutionResult>,
128 specs: Option<&[RuntimeSpec]>,
129 env_groups: Option<&[RuntimeEnvGroup]>,
130 selected_env_group_slug: Option<&str>,
131 should_cancel: FCancel,
132) -> Option<StepExecutionResult>
133where
134 FCancel: FnMut() -> bool,
135{
136 send_prepared_http_step_with_hooks(
137 client,
138 prepared,
139 step,
140 context,
141 specs,
142 env_groups,
143 selected_env_group_slug,
144 should_cancel,
145 || async {},
146 || async {},
147 || async {},
148 )
149 .await
150}
151
152pub async fn send_prepared_http_step_with_hooks<
153 FCancel,
154 FStart,
155 FStartFuture,
156 FSend,
157 FSendFuture,
158 FBody,
159 FBodyFuture,
160>(
161 client: &Client,
162 prepared: PreparedHttpStep,
163 step: &PipelineStep,
164 context: &HashMap<String, StepExecutionResult>,
165 specs: Option<&[RuntimeSpec]>,
166 env_groups: Option<&[RuntimeEnvGroup]>,
167 selected_env_group_slug: Option<&str>,
168 mut should_cancel: FCancel,
169 on_send_started: FStart,
170 on_send_returned: FSend,
171 on_body_completed: FBody,
172) -> Option<StepExecutionResult>
173where
174 FCancel: FnMut() -> bool,
175 FStart: FnMut() -> FStartFuture,
176 FStartFuture: Future<Output = ()>,
177 FSend: FnMut() -> FSendFuture,
178 FSendFuture: Future<Output = ()>,
179 FBody: FnMut() -> FBodyFuture,
180 FBodyFuture: Future<Output = ()>,
181{
182 let started = start_prepared_http_step_with_hooks(
183 client,
184 prepared,
185 step,
186 || should_cancel(),
187 on_send_started,
188 on_send_returned,
189 )
190 .await?;
191
192 match started {
193 Ok(started) => {
194 complete_started_http_step_with_hook(
195 started,
196 step,
197 context,
198 specs,
199 env_groups,
200 selected_env_group_slug,
201 should_cancel,
202 on_body_completed,
203 )
204 .await
205 }
206 Err(result) => Some(result),
207 }
208}
209
210pub async fn start_prepared_http_step_with_hooks<
211 FCancel,
212 FStart,
213 FStartFuture,
214 FSend,
215 FSendFuture,
216>(
217 client: &Client,
218 prepared: PreparedHttpStep,
219 step: &PipelineStep,
220 mut should_cancel: FCancel,
221 mut on_send_started: FStart,
222 mut on_send_returned: FSend,
223) -> Option<Result<StartedHttpStep, StepExecutionResult>>
224where
225 FCancel: FnMut() -> bool,
226 FStart: FnMut() -> FStartFuture,
227 FStartFuture: Future<Output = ()>,
228 FSend: FnMut() -> FSendFuture,
229 FSendFuture: Future<Output = ()>,
230{
231 let mut request_builder = client.request(prepared.method.clone(), prepared.url.clone());
232
233 for (key, value) in &prepared.request.headers {
234 request_builder = request_builder.header(key, value);
235 }
236
237 if let Some(body) = prepared.request.body.as_ref() {
238 if !prepared.request.method.eq_ignore_ascii_case("GET")
239 && !prepared.request.method.eq_ignore_ascii_case("HEAD")
240 {
241 request_builder = request_builder.json(body);
242 }
243 }
244
245 let request = prepared.request.clone();
246 if should_cancel() {
247 return None;
248 }
249 on_send_started().await;
250 let Some(send_result) = await_with_cancel(request_builder.send(), &mut should_cancel).await
251 else {
252 return None;
253 };
254 on_send_returned().await;
255
256 match send_result {
257 Ok(response) => Some(Ok(StartedHttpStep {
258 request,
259 response,
260 started_at: prepared.started_at,
261 attempt: prepared.attempt,
262 max_attempts: prepared.max_attempts,
263 })),
264 Err(err) => {
265 let result = step_result(
266 step,
267 request,
268 None,
269 Some(err.to_string()),
270 "error",
271 prepared.started_at,
272 prepared.attempt,
273 prepared.max_attempts,
274 None,
275 );
276 log_step_response(&step.id, None, result.error.as_deref(), &result.extracts);
277 Some(Err(result))
278 }
279 }
280}
281
282pub async fn complete_started_http_step_with_hook<FCancel, FBody, FBodyFuture>(
283 started: StartedHttpStep,
284 step: &PipelineStep,
285 context: &HashMap<String, StepExecutionResult>,
286 specs: Option<&[RuntimeSpec]>,
287 env_groups: Option<&[RuntimeEnvGroup]>,
288 selected_env_group_slug: Option<&str>,
289 mut should_cancel: FCancel,
290 mut on_body_completed: FBody,
291) -> Option<StepExecutionResult>
292where
293 FCancel: FnMut() -> bool,
294 FBody: FnMut() -> FBodyFuture,
295 FBodyFuture: Future<Output = ()>,
296{
297 let status = started.response.status();
298 let status_text = status.canonical_reason().unwrap_or("").to_owned();
299 let mut headers = HashMap::new();
300 for (key, value) in started.response.headers() {
301 headers.insert(
302 key.as_str().to_owned(),
303 value.to_str().unwrap_or_default().to_owned(),
304 );
305 }
306
307 let content_type = headers
308 .iter()
309 .find(|(key, _)| key.eq_ignore_ascii_case("content-type"))
310 .map(|(_, value)| value.as_str())
311 .unwrap_or("");
312
313 let body = if content_type.contains("application/json") {
314 let Some(body_result) =
315 await_with_cancel(started.response.json::<Value>(), &mut should_cancel).await
316 else {
317 return None;
318 };
319 on_body_completed().await;
320 match body_result {
321 Ok(value) => value,
322 Err(err) => {
323 let result = step_result(
324 step,
325 started.request,
326 None,
327 Some(err.to_string()),
328 "error",
329 started.started_at,
330 started.attempt,
331 started.max_attempts,
332 None,
333 );
334 log_step_response(&step.id, None, result.error.as_deref(), &result.extracts);
335 return Some(result);
336 }
337 }
338 } else {
339 let Some(body_result) =
340 await_with_cancel(started.response.text(), &mut should_cancel).await
341 else {
342 return None;
343 };
344 on_body_completed().await;
345 Value::String(body_result.unwrap_or_default())
346 };
347
348 let http_error =
349 (!status.is_success()).then(|| format!("HTTP {} {}", status.as_u16(), status_text));
350 let mut result = step_result(
351 step,
352 started.request,
353 Some(StepResponse {
354 status: status.as_u16(),
355 status_text: status_text.clone(),
356 headers,
357 body,
358 }),
359 http_error.clone(),
360 "success",
361 started.started_at,
362 started.attempt,
363 started.max_attempts,
364 None,
365 );
366
367 let extraction_failed = match evaluate_step_extractions(step, &result) {
368 Ok(extracts) => {
369 result.extracts = extracts;
370 false
371 }
372 Err(error) => {
373 result.status = "error".to_owned();
374 result.error = Some(match result.error.take() {
375 Some(existing) => format!("{existing} | {error}"),
376 None => error,
377 });
378 true
379 }
380 };
381
382 if !extraction_failed {
383 let has_status_assert = has_status_assertion(step);
384 let assert_results = evaluate_assertions(
385 step,
386 &result,
387 context,
388 specs,
389 env_groups,
390 selected_env_group_slug,
391 );
392 let assertion_failed = assert_results.iter().any(|result| !result.passed);
393 if !assert_results.is_empty() {
394 if assertion_failed {
395 result.status = "error".to_owned();
396 let failed_count = assert_results
397 .iter()
398 .filter(|result| !result.passed)
399 .count();
400 result.error = Some(match result.error {
401 Some(err) => format!("{} | {} assertion(s) failed", err, failed_count),
402 None => format!("{} assertion(s) failed", failed_count),
403 });
404 } else if http_error.is_some() {
405 if has_status_assert {
406 result.status = "success".to_owned();
407 result.error = None;
408 } else {
409 result.status = "error".to_owned();
410 }
411 }
412 result.assert_results = Some(assert_results);
413 } else if http_error.is_some() {
414 result.status = "error".to_owned();
415 }
416 }
417
418 log_step_response(
419 &step.id,
420 result.response.as_ref(),
421 result.error.as_deref(),
422 &result.extracts,
423 );
424 Some(result)
425}
426
427fn invalid_step_result(
428 step: &PipelineStep,
429 request: StepRequest,
430 error: String,
431 started_at: Instant,
432 attempt: usize,
433 max_attempts: usize,
434) -> StepExecutionResult {
435 step_result(
436 step,
437 request,
438 None,
439 Some(error),
440 "error",
441 started_at,
442 attempt,
443 max_attempts,
444 None,
445 )
446}
447
448fn step_result(
449 step: &PipelineStep,
450 request: StepRequest,
451 response: Option<StepResponse>,
452 error: Option<String>,
453 status: &str,
454 started_at: Instant,
455 attempt: usize,
456 max_attempts: usize,
457 assert_results: Option<Vec<crate::core::types::AssertionResult>>,
458) -> StepExecutionResult {
459 StepExecutionResult {
460 step_id: step.id.clone(),
461 status: status.to_owned(),
462 request: Some(request),
463 response,
464 error,
465 duration: Some(started_at.elapsed().as_millis()),
466 attempts: if max_attempts > 1 {
467 Some(attempt)
468 } else {
469 None
470 },
471 attempt: Some(attempt),
472 max_attempts: Some(max_attempts),
473 extracts: HashMap::new(),
474 assert_results,
475 }
476}
477
478#[cfg(test)]
479mod tests {
480 use super::*;
481 use crate::core::types::{PipelineStep, StepExtraction};
482 use httpmock::Method::GET;
483 use serde_json::json;
484 use std::collections::HashMap;
485
486 #[tokio::test]
487 async fn sends_prepared_step_and_returns_success_result() {
488 let server = httpmock::MockServer::start_async().await;
489 server
490 .mock_async(|when, then| {
491 when.method(GET).path("/users");
492 then.status(200)
493 .header("content-type", "application/json")
494 .json_body(json!({"ok": true}));
495 })
496 .await;
497
498 let client = reqwest::Client::new();
499 let step = PipelineStep {
500 id: "get-users".to_owned(),
501 name: "GET users".to_owned(),
502 description: None,
503 method: "GET".to_owned(),
504 url: format!("{}/users", server.base_url()),
505 headers: HashMap::new(),
506 body: None,
507 operation_id: None,
508 delay: None,
509 retry: None,
510 extracts: Vec::new(),
511 asserts: vec![],
512 };
513 let context = HashMap::new();
514
515 let prepared = prepare_http_step(&step, &context, None, None, None, 1, 1)
516 .expect("step should prepare");
517
518 let result =
519 send_prepared_http_step(&client, prepared, &step, &context, None, None, None, || {
520 false
521 })
522 .await
523 .expect("send should not be cancelled");
524
525 assert_eq!(result.step_id, "get-users");
526 assert_eq!(result.status, "success");
527 assert_eq!(result.response.as_ref().map(|r| r.status), Some(200));
528 }
529
530 #[tokio::test]
531 async fn prepared_step_extracts_from_text_response() {
532 let server = httpmock::MockServer::start_async().await;
533 server
534 .mock_async(|when, then| {
535 when.method(GET).path("/message");
536 then.status(200)
537 .header("content-type", "text/html")
538 .body("<strong>123456</strong>");
539 })
540 .await;
541 let step = PipelineStep {
542 id: "message".to_owned(),
543 name: "Read message".to_owned(),
544 description: None,
545 method: "GET".to_owned(),
546 url: format!("{}/message", server.base_url()),
547 headers: HashMap::new(),
548 body: None,
549 operation_id: None,
550 delay: None,
551 retry: None,
552 asserts: Vec::new(),
553 extracts: vec![StepExtraction {
554 name: "code".to_owned(),
555 field: "body".to_owned(),
556 regex: r"<strong>([0-9]{6})</strong>".to_owned(),
557 group: 1,
558 required: true,
559 }],
560 };
561 let context = HashMap::new();
562 let prepared = prepare_http_step(&step, &context, None, None, None, 1, 1)
563 .expect("step should prepare");
564
565 let result = send_prepared_http_step(
566 &reqwest::Client::new(),
567 prepared,
568 &step,
569 &context,
570 None,
571 None,
572 None,
573 || false,
574 )
575 .await
576 .expect("send should not be cancelled");
577
578 assert_eq!(result.extracts.get("code"), Some(&"123456".to_owned()));
579 }
580
581 #[tokio::test]
582 async fn hooks_report_send_started_before_send_returned() {
583 let server = httpmock::MockServer::start_async().await;
584 server
585 .mock_async(|when, then| {
586 when.method(GET).path("/users");
587 then.status(200)
588 .header("content-type", "application/json")
589 .json_body(json!({"ok": true}));
590 })
591 .await;
592
593 let client = reqwest::Client::new();
594 let step = PipelineStep {
595 id: "get-users".to_owned(),
596 name: "GET users".to_owned(),
597 description: None,
598 method: "GET".to_owned(),
599 url: format!("{}/users", server.base_url()),
600 headers: HashMap::new(),
601 body: None,
602 operation_id: None,
603 delay: None,
604 retry: None,
605 extracts: Vec::new(),
606 asserts: vec![],
607 };
608 let context = HashMap::new();
609 let prepared = prepare_http_step(&step, &context, None, None, None, 1, 1)
610 .expect("step should prepare");
611 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::<&'static str>::new()));
612
613 let started_events = std::sync::Arc::clone(&events);
614 let returned_events = std::sync::Arc::clone(&events);
615 let result = send_prepared_http_step_with_hooks(
616 &client,
617 prepared,
618 &step,
619 &context,
620 None,
621 None,
622 None,
623 || false,
624 move || {
625 let events = std::sync::Arc::clone(&started_events);
626 async move {
627 events.lock().expect("events lock").push("started");
628 }
629 },
630 move || {
631 let events = std::sync::Arc::clone(&returned_events);
632 async move {
633 events.lock().expect("events lock").push("returned");
634 }
635 },
636 || async {},
637 )
638 .await
639 .expect("send should not be cancelled");
640
641 assert_eq!(result.status, "success");
642 assert_eq!(
643 events.lock().expect("events lock").as_slice(),
644 ["started", "returned"]
645 );
646 }
647
648 #[tokio::test]
649 async fn split_http_helpers_start_send_before_body_completion() {
650 let server = httpmock::MockServer::start_async().await;
651 server
652 .mock_async(|when, then| {
653 when.method(GET).path("/users");
654 then.status(200)
655 .header("content-type", "application/json")
656 .json_body(json!({"ok": true}));
657 })
658 .await;
659
660 let client = reqwest::Client::new();
661 let step = PipelineStep {
662 id: "get-users".to_owned(),
663 name: "GET users".to_owned(),
664 description: None,
665 method: "GET".to_owned(),
666 url: format!("{}/users", server.base_url()),
667 headers: HashMap::new(),
668 body: None,
669 operation_id: None,
670 delay: None,
671 retry: None,
672 extracts: Vec::new(),
673 asserts: vec![],
674 };
675 let context = HashMap::new();
676 let prepared = prepare_http_step(&step, &context, None, None, None, 1, 1)
677 .expect("step should prepare");
678 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::<&'static str>::new()));
679
680 let started_events = std::sync::Arc::clone(&events);
681 let returned_events = std::sync::Arc::clone(&events);
682 let started = start_prepared_http_step_with_hooks(
683 &client,
684 prepared,
685 &step,
686 || false,
687 move || {
688 let events = std::sync::Arc::clone(&started_events);
689 async move {
690 events.lock().expect("events lock").push("started");
691 }
692 },
693 move || {
694 let events = std::sync::Arc::clone(&returned_events);
695 async move {
696 events.lock().expect("events lock").push("returned");
697 }
698 },
699 )
700 .await
701 .expect("start should not be cancelled")
702 .expect("request should start");
703
704 assert_eq!(
705 events.lock().expect("events lock").as_slice(),
706 ["started", "returned"]
707 );
708
709 let body_events = std::sync::Arc::clone(&events);
710 let result = complete_started_http_step_with_hook(
711 started,
712 &step,
713 &context,
714 None,
715 None,
716 None,
717 || false,
718 move || {
719 let events = std::sync::Arc::clone(&body_events);
720 async move {
721 events.lock().expect("events lock").push("body");
722 }
723 },
724 )
725 .await
726 .expect("complete should not be cancelled");
727
728 assert_eq!(result.status, "success");
729 assert_eq!(
730 events.lock().expect("events lock").as_slice(),
731 ["started", "returned", "body"]
732 );
733 }
734}