Skip to main content

previa_engine/execution/
http_step.rs

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}