diff --git a/src/core/src/processor/incremental.rs b/src/core/src/processor/incremental.rs index 46318981..d9dfd487 100644 --- a/src/core/src/processor/incremental.rs +++ b/src/core/src/processor/incremental.rs @@ -1,37 +1,12 @@ -use crate::request_substitution::{ - substitute_functions_in_request, substitute_request_variables_in_request, -}; -use crate::conditions; +use super::incremental_loop::{SyncSleep, block_on, process_requests_incremental}; use crate::parser; use crate::runner; -use crate::types::{HttpRequest, HttpResult, RequestContext}; +use crate::types::{HttpRequest, HttpResult}; use anyhow::Result; -/// Result of processing a single request -#[derive(Debug)] -pub enum RequestProcessingResult { - /// Request was skipped due to conditions or dependencies - Skipped { - request: HttpRequest, - reason: String, - }, - /// Request was executed successfully or with errors - Executed { - request: HttpRequest, - result: HttpResult, - }, - /// Request processing failed with an error - Failed { request: HttpRequest, error: String }, -} +pub use super::incremental_loop::RequestProcessingResult; /// Process HTTP requests from a file with incremental callbacks for UI updates -/// -/// This function processes requests one at a time, maintaining proper context -/// for variable substitution, function evaluation, and condition checking. -/// It calls the provided callback after each request is processed. -/// -/// The callback can return `true` to continue processing or `false` to stop. -/// This allows processing to stop early after reaching a target request. pub fn process_http_file_incremental( file_path: &str, environment: Option<&str>, @@ -53,204 +28,30 @@ where } /// Process HTTP requests from a file with incremental callbacks and a custom executor. -/// -/// Like `process_http_file_incremental`, but accepts a custom executor function -/// for HTTP request execution. This is useful for testing with mock executors. pub fn process_http_file_incremental_with_executor( file_path: &str, environment: Option<&str>, insecure: bool, delay_ms: u64, - mut callback: F, + callback: F, executor: &E, ) -> Result<()> where F: FnMut(usize, usize, RequestProcessingResult) -> bool, E: Fn(&HttpRequest, bool, bool) -> Result, { - // Parse the file let requests = parser::parse_http_file(file_path, environment)?; - let total = requests.len(); - - if requests.is_empty() { - return Ok(()); - } - - let mut request_contexts: Vec = Vec::new(); - - for (idx, mut request) in requests.into_iter().enumerate() { - let request_count = (idx + 1) as u32; - - // Apply delay between requests (not before first request) - if idx > 0 && delay_ms > 0 { - std::thread::sleep(std::time::Duration::from_millis(delay_ms)); - } - - // Check dependencies - if let Some(dep_name) = request.depends_on.as_ref() - && !conditions::check_dependency(&Some(dep_name.clone()), &request_contexts) - { - let reason = format!("Dependency on '{}' not met", dep_name); - let should_continue = callback( - idx, - total, - RequestProcessingResult::Skipped { - request: request.clone(), - reason, - }, - ); - add_request_context(&mut request_contexts, request, None, request_count); - if !should_continue { - break; - } - continue; - } - - // Check conditions - if !request.conditions.is_empty() { - match conditions::evaluate_conditions(&request.conditions, &request_contexts) { - Ok(true) => { - // Conditions met, continue - } - Ok(false) => { - // Conditions not met, skip - let should_continue = callback( - idx, - total, - RequestProcessingResult::Skipped { - request: request.clone(), - reason: "Conditions not met".to_string(), - }, - ); - add_request_context(&mut request_contexts, request, None, request_count); - if !should_continue { - break; - } - continue; - } - Err(e) => { - let should_continue = callback( - idx, - total, - RequestProcessingResult::Failed { - request: request.clone(), - error: format!("Condition evaluation error: {}", e), - }, - ); - add_request_context(&mut request_contexts, request, None, request_count); - if !should_continue { - break; - } - continue; - } - } - } - - // Apply variable substitutions - if let Err(e) = substitute_request_variables_in_request(&mut request, &request_contexts) { - let should_continue = callback( - idx, - total, - RequestProcessingResult::Failed { - request: request.clone(), - error: format!("Variable substitution error: {}", e), - }, - ); - add_request_context(&mut request_contexts, request, None, request_count); - if !should_continue { - break; - } - continue; - } - - // Apply function substitutions - if let Err(e) = substitute_functions_in_request(&mut request) { - let should_continue = callback( - idx, - total, - RequestProcessingResult::Failed { - request: request.clone(), - error: format!("Function substitution error: {}", e), - }, - ); - add_request_context(&mut request_contexts, request, None, request_count); - if !should_continue { - break; - } - continue; - } - // Apply pre-request delay - if let Some(pre_delay) = request.pre_delay_ms - && pre_delay > 0 - { - std::thread::sleep(std::time::Duration::from_millis(pre_delay)); - } - - // Capture post-delay before request is moved - let post_delay = request.post_delay_ms; - - // Execute the request - match executor(&request, false, insecure) { - Ok(result) => { - add_request_context( - &mut request_contexts, - request.clone(), - Some(result.clone()), - request_count, - ); - let should_continue = callback( - idx, - total, - RequestProcessingResult::Executed { request, result }, - ); - if !should_continue { - break; - } - } - Err(e) => { - let should_continue = callback( - idx, - total, - RequestProcessingResult::Failed { - request: request.clone(), - error: e.to_string(), - }, - ); - add_request_context(&mut request_contexts, request, None, request_count); - if !should_continue { - break; - } - } - } - - // Apply post-request delay - if let Some(post_delay) = post_delay - && post_delay > 0 - { - std::thread::sleep(std::time::Duration::from_millis(post_delay)); - } - } - - Ok(()) -} - -fn add_request_context( - contexts: &mut Vec, - request: HttpRequest, - result: Option, - request_count: u32, -) { - let context_name = if let Some(ref name) = request.name { - name.clone() - } else { - // Use same format as executor.rs for consistency - format!("request_{}", request_count) + let wrapped = move |request: HttpRequest, verbose: bool, insecure: bool| { + async move { executor(&request, verbose, insecure) } }; - contexts.push(RequestContext { - name: context_name, - request, - result, - }); + block_on(process_requests_incremental( + requests, + insecure, + delay_ms, + callback, + &wrapped, + SyncSleep, + )) } diff --git a/src/core/src/processor/incremental_loop.rs b/src/core/src/processor/incremental_loop.rs new file mode 100644 index 00000000..8f312226 --- /dev/null +++ b/src/core/src/processor/incremental_loop.rs @@ -0,0 +1,832 @@ +use crate::assertions; +use crate::conditions; +use crate::request_substitution::{ + substitute_functions_in_request, substitute_request_variables_in_request, +}; +use crate::types::{HttpRequest, HttpResult, RequestContext}; +use anyhow::Result; +use std::future::Future; +use std::pin::Pin; +#[cfg(not(target_arch = "wasm32"))] +use std::sync::{ + Arc, Mutex, + atomic::{AtomicBool, Ordering}, +}; +#[cfg(not(target_arch = "wasm32"))] +use std::task::Waker; +use std::time::Duration; + +pub type AsyncRequestFuture<'a> = Pin> + 'a>>; +pub type AsyncRequestExecutor = + dyn for<'a> Fn(&'a HttpRequest, bool, bool) -> AsyncRequestFuture<'a>; + +/// Result of processing a single request during incremental execution. +#[derive(Debug)] +pub enum RequestProcessingResult { + Skipped { + request: HttpRequest, + reason: String, + }, + Executed { + request: HttpRequest, + result: HttpResult, + }, + Failed { request: HttpRequest, error: String }, +} + +/// Abstraction over sleep mechanisms for sync and async paths. +pub trait Sleep { + async fn sleep(&self, duration: Duration); +} + +/// Sync sleep adapter — blocks the current thread. +pub struct SyncSleep; + +impl Sleep for SyncSleep { + async fn sleep(&self, duration: Duration) { + if duration == Duration::ZERO { + return; + } + std::thread::sleep(duration); + } +} + +/// Async sleep adapter — uses platform-appropriate async sleep. +pub struct AsyncSleep; + +impl Sleep for AsyncSleep { + async fn sleep(&self, duration: Duration) { + async_sleep_ms(duration).await; + } +} + +#[cfg(not(target_arch = "wasm32"))] +async fn async_sleep_ms(duration: Duration) { + if duration == Duration::ZERO { + return; + } + NativeSleep::new(duration).await; +} + +#[cfg(target_arch = "wasm32")] +async fn async_sleep_ms(duration: Duration) { + if duration == Duration::ZERO { + return; + } + gloo_timers::future::sleep(duration).await; +} + +#[cfg(not(target_arch = "wasm32"))] +struct NativeSleep { + duration: Option, + state: Arc, +} + +#[cfg(not(target_arch = "wasm32"))] +struct NativeSleepState { + completed: AtomicBool, + waker: Mutex>, +} + +#[cfg(not(target_arch = "wasm32"))] +impl NativeSleep { + fn new(duration: Duration) -> Self { + Self { + duration: Some(duration), + state: Arc::new(NativeSleepState { + completed: AtomicBool::new(false), + waker: Mutex::new(None), + }), + } + } +} + +#[cfg(not(target_arch = "wasm32"))] +impl Future for NativeSleep { + type Output = (); + + fn poll( + self: Pin<&mut Self>, + cx: &mut std::task::Context<'_>, + ) -> std::task::Poll { + let this = self.get_mut(); + + if this.state.completed.load(Ordering::Acquire) { + return std::task::Poll::Ready(()); + } + + { + let mut waker = this + .state + .waker + .lock() + .expect("native sleep waker mutex poisoned"); + *waker = Some(cx.waker().clone()); + } + + if let Some(duration) = this.duration.take() { + let state = Arc::clone(&this.state); + std::thread::spawn(move || { + std::thread::sleep(duration); + state.completed.store(true, Ordering::Release); + if let Some(waker) = state + .waker + .lock() + .expect("native sleep waker mutex poisoned") + .take() + { + waker.wake(); + } + }); + } + + std::task::Poll::Pending + } +} + +pub(crate) fn add_request_context( + contexts: &mut Vec, + request: HttpRequest, + result: Option, + request_count: u32, +) { + let context_name = request + .name + .clone() + .unwrap_or_else(|| format!("request_{}", request_count)); + + contexts.push(RequestContext { + name: context_name, + request, + result, + }); +} + +/// Process requests incrementally with dependency checking, condition evaluation, +/// variable/function substitution, pre/post delays, and callback-driven control flow. +/// +/// The executor is called with an owned `HttpRequest` (the loop clones it before +/// dispatching), so the original remains available for callback and context tracking. +pub async fn process_requests_incremental( + requests: Vec, + insecure: bool, + delay_ms: u64, + mut callback: F, + executor: &impl Fn(HttpRequest, bool, bool) -> Fut, + sleep: S, +) -> Result<()> +where + F: FnMut(usize, usize, RequestProcessingResult) -> bool, + Fut: Future>, + S: Sleep, +{ + let total = requests.len(); + + if requests.is_empty() { + return Ok(()); + } + + let mut request_contexts: Vec = Vec::new(); + + for (idx, mut request) in requests.into_iter().enumerate() { + let request_count = (idx + 1) as u32; + + if idx > 0 && delay_ms > 0 { + sleep.sleep(Duration::from_millis(delay_ms)).await; + } + + if let Some(dep_name) = request.depends_on.as_ref() + && !conditions::check_dependency(&Some(dep_name.clone()), &request_contexts) + { + let should_continue = callback( + idx, + total, + RequestProcessingResult::Skipped { + request: request.clone(), + reason: format!("Dependency on '{}' not met", dep_name), + }, + ); + add_request_context(&mut request_contexts, request, None, request_count); + if !should_continue { + break; + } + continue; + } + + if !request.conditions.is_empty() { + match conditions::evaluate_conditions(&request.conditions, &request_contexts) { + Ok(true) => {} + Ok(false) => { + let should_continue = callback( + idx, + total, + RequestProcessingResult::Skipped { + request: request.clone(), + reason: "Conditions not met".to_string(), + }, + ); + add_request_context(&mut request_contexts, request, None, request_count); + if !should_continue { + break; + } + continue; + } + Err(error) => { + let should_continue = callback( + idx, + total, + RequestProcessingResult::Failed { + request: request.clone(), + error: format!("Condition evaluation error: {}", error), + }, + ); + add_request_context(&mut request_contexts, request, None, request_count); + if !should_continue { + break; + } + continue; + } + } + } + + if let Err(error) = substitute_request_variables_in_request(&mut request, &request_contexts) + { + let should_continue = callback( + idx, + total, + RequestProcessingResult::Failed { + request: request.clone(), + error: format!("Variable substitution error: {}", error), + }, + ); + add_request_context(&mut request_contexts, request, None, request_count); + if !should_continue { + break; + } + continue; + } + + if let Err(error) = substitute_functions_in_request(&mut request) { + let should_continue = callback( + idx, + total, + RequestProcessingResult::Failed { + request: request.clone(), + error: format!("Function substitution error: {}", error), + }, + ); + add_request_context(&mut request_contexts, request, None, request_count); + if !should_continue { + break; + } + continue; + } + + if let Some(pre_delay_ms) = request.pre_delay_ms + && pre_delay_ms > 0 + { + sleep.sleep(Duration::from_millis(pre_delay_ms)).await; + } + + let post_delay_ms = request.post_delay_ms; + + // Clone the request for the executor so the original remains available + // for the callback and context tracking. + match executor(request.clone(), false, insecure).await { + Ok(mut result) => { + if !request.assertions.is_empty() { + let assertion_results = + assertions::evaluate_assertions(&request.assertions, &result); + let all_passed = assertion_results.iter().all(|r| r.passed); + result.success = all_passed; + result.assertion_results = assertion_results; + } + add_request_context( + &mut request_contexts, + request.clone(), + Some(result.clone()), + request_count, + ); + let should_continue = callback( + idx, + total, + RequestProcessingResult::Executed { request, result }, + ); + if !should_continue { + break; + } + } + Err(error) => { + let should_continue = callback( + idx, + total, + RequestProcessingResult::Failed { + request: request.clone(), + error: error.to_string(), + }, + ); + add_request_context(&mut request_contexts, request, None, request_count); + if !should_continue { + break; + } + } + } + + if let Some(post_delay_ms) = post_delay_ms + && post_delay_ms > 0 + { + sleep.sleep(Duration::from_millis(post_delay_ms)).await; + } + } + + Ok(()) +} + +/// Block on a future using a no-op waker. +pub(crate) fn block_on(future: F) -> F::Output { + use std::task::{Context, Poll}; + + let waker = std::task::Waker::noop(); + let mut context = Context::from_waker(waker); + let mut future = Box::pin(future); + + loop { + match future.as_mut().poll(&mut context) { + Poll::Ready(output) => return output, + Poll::Pending => std::thread::yield_now(), + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::types::{Assertion, AssertionType, Condition, ConditionType}; + use std::pin::Pin; + use std::sync::{ + Arc, Mutex, + atomic::{AtomicBool, Ordering}, + }; + use std::task::{Context, Poll}; + + struct MockSleep { + calls: Arc>>, + } + + impl MockSleep { + fn new() -> Self { + Self { + calls: Arc::new(Mutex::new(Vec::new())), + } + } + } + + impl Sleep for MockSleep { + async fn sleep(&self, duration: Duration) { + self.calls.lock().unwrap().push(duration); + } + } + + fn make_result(name: Option) -> HttpResult { + HttpResult { + request_name: name, + status_code: 200, + success: true, + error_message: None, + duration_ms: 10, + response_headers: None, + response_body: None, + assertion_results: vec![], + } + } + + fn make_request(name: &str) -> HttpRequest { + HttpRequest { + name: Some(name.to_string()), + method: "GET".to_string(), + url: "https://example.com/test".to_string(), + headers: vec![], + body: None, + assertions: vec![], + variables: vec![], + timeout: None, + connection_timeout: None, + depends_on: None, + conditions: vec![], + pre_delay_ms: None, + post_delay_ms: None, + } + } + + fn ok_executor() -> impl Fn(HttpRequest, bool, bool) -> Pin>>> { + |request: HttpRequest, _verbose: bool, _insecure: bool| { + let name = request.name.clone(); + Box::pin(async move { Ok(make_result(name)) }) + } + } + + // --- Sleep trait tests --- + + #[test] + fn test_sync_sleep_zero_duration() { + let sleep = SyncSleep; + block_on(sleep.sleep(Duration::ZERO)); + } + + #[test] + fn test_sync_sleep_nonzero() { + let sleep = SyncSleep; + let start = std::time::Instant::now(); + block_on(sleep.sleep(Duration::from_millis(5))); + let elapsed = start.elapsed(); + assert!(elapsed >= Duration::from_millis(5)); + } + + // --- block_on tests --- + + #[test] + fn test_block_on_ready() { + let result = block_on(async { 42 }); + assert_eq!(result, 42); + } + + #[test] + fn test_block_on_pending_once() { + struct YieldOnce { + yielded: bool, + } + impl Future for YieldOnce { + type Output = &'static str; + fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll { + if self.yielded { + Poll::Ready("done") + } else { + self.yielded = true; + cx.waker().wake_by_ref(); + Poll::Pending + } + } + } + let result = block_on(YieldOnce { yielded: false }); + assert_eq!(result, "done"); + } + + // --- process_requests_incremental tests --- + + #[test] + fn test_process_requests_empty() { + let executor = ok_executor(); + block_on(process_requests_incremental( + vec![], + false, + 0, + |_, _, _| true, + &executor, + MockSleep::new(), + )) + .unwrap(); + } + + #[test] + fn test_process_requests_single_execution() { + let requests = vec![make_request("req1")]; + let executor = ok_executor(); + let results = Arc::new(Mutex::new(Vec::new())); + let r = Arc::clone(&results); + + block_on(process_requests_incremental( + requests, + false, + 0, + |idx, total, result| { + r.lock().unwrap().push((idx, total, result)); + true + }, + &executor, + MockSleep::new(), + )) + .unwrap(); + + let results = results.lock().unwrap(); + assert_eq!(results.len(), 1); + assert_eq!(results[0].0, 0); + assert_eq!(results[0].1, 1); + assert!(matches!(results[0].2, RequestProcessingResult::Executed { .. })); + } + + #[test] + fn test_process_requests_multiple_executions() { + let requests = vec![make_request("a"), make_request("b"), make_request("c")]; + let executor = ok_executor(); + let results = Arc::new(Mutex::new(Vec::new())); + let r = Arc::clone(&results); + + block_on(process_requests_incremental( + requests, + false, + 0, + |idx, total, result| { + r.lock().unwrap().push((idx, total, result)); + true + }, + &executor, + MockSleep::new(), + )) + .unwrap(); + + let results = results.lock().unwrap(); + assert_eq!(results.len(), 3); + for (i, (idx, total, _)) in results.iter().enumerate() { + assert_eq!(*idx, i); + assert_eq!(*total, 3); + } + } + + #[test] + fn test_process_requests_callback_stops_early() { + let requests = vec![make_request("a"), make_request("b"), make_request("c")]; + let executor = ok_executor(); + let results = Arc::new(Mutex::new(Vec::new())); + let r = Arc::clone(&results); + + block_on(process_requests_incremental( + requests, + false, + 0, + |idx, _total, result| { + r.lock().unwrap().push(result); + idx < 1 + }, + &executor, + MockSleep::new(), + )) + .unwrap(); + + assert_eq!(results.lock().unwrap().len(), 2); + } + + #[test] + fn test_process_requests_inter_request_delay() { + let requests = vec![make_request("a"), make_request("b")]; + let executor = ok_executor(); + let calls = Arc::new(Mutex::new(Vec::new())); + let mock_sleep = MockSleep { calls: Arc::clone(&calls) }; + + block_on(process_requests_incremental( + requests, + false, + 50, + |_, _, _| true, + &executor, + mock_sleep, + )) + .unwrap(); + + assert_eq!(calls.lock().unwrap().len(), 1); + } + + #[test] + fn test_process_requests_dependency_skip() { + let requests = vec![HttpRequest { + depends_on: Some("missing".to_string()), + ..make_request("req1") + }]; + let executor = ok_executor(); + let results = Arc::new(Mutex::new(Vec::new())); + let r = Arc::clone(&results); + + block_on(process_requests_incremental( + requests, + false, + 0, + |_idx, _total, result| { + r.lock().unwrap().push(result); + true + }, + &executor, + MockSleep::new(), + )) + .unwrap(); + + let results = results.lock().unwrap(); + assert_eq!(results.len(), 1); + assert!(matches!(&results[0], RequestProcessingResult::Skipped { reason, .. } + if reason.contains("Dependency"))); + } + + #[test] + fn test_process_requests_condition_skip() { + let requests = vec![HttpRequest { + conditions: vec![Condition { + request_name: "other".to_string(), + condition_type: ConditionType::Status, + expected_value: "200".to_string(), + negate: false, + }], + ..make_request("req1") + }]; + let executor = ok_executor(); + let results = Arc::new(Mutex::new(Vec::new())); + let r = Arc::clone(&results); + + block_on(process_requests_incremental( + requests, + false, + 0, + |_idx, _total, result| { + r.lock().unwrap().push(result); + true + }, + &executor, + MockSleep::new(), + )) + .unwrap(); + + let results = results.lock().unwrap(); + assert_eq!(results.len(), 1); + assert!(matches!(&results[0], RequestProcessingResult::Skipped { reason, .. } + if reason == "Conditions not met")); + } + + #[test] + fn test_process_requests_executor_error() { + let requests = vec![make_request("req1")]; + let results = Arc::new(Mutex::new(Vec::new())); + let r = Arc::clone(&results); + + let err_executor = |_: HttpRequest, _: bool, _: bool| { + Box::pin(async move { Err(anyhow::anyhow!("executor failure")) }) + as Pin>>> + }; + + block_on(process_requests_incremental( + requests, + false, + 0, + |_idx, _total, result| { + r.lock().unwrap().push(result); + true + }, + &err_executor, + MockSleep::new(), + )) + .unwrap(); + + let results = results.lock().unwrap(); + assert_eq!(results.len(), 1); + assert!(matches!(&results[0], RequestProcessingResult::Failed { error, .. } + if error == "executor failure")); + } + + #[test] + fn test_process_requests_assertion_failure() { + let requests = vec![HttpRequest { + assertions: vec![Assertion { + assertion_type: AssertionType::Status, + expected_value: "404".to_string(), + }], + ..make_request("req1") + }]; + let executor = ok_executor(); + let results = Arc::new(Mutex::new(Vec::new())); + let r = Arc::clone(&results); + + block_on(process_requests_incremental( + requests, + false, + 0, + |_idx, _total, result| { + r.lock().unwrap().push(result); + true + }, + &executor, + MockSleep::new(), + )) + .unwrap(); + + let results = results.lock().unwrap(); + assert_eq!(results.len(), 1); + match &results[0] { + RequestProcessingResult::Executed { result, .. } => { + assert!(!result.success); + assert_eq!(result.assertion_results.len(), 1); + assert!(!result.assertion_results[0].passed); + } + _ => panic!("expected Executed result"), + } + } + + #[test] + fn test_process_requests_no_delay_when_zero() { + let requests = vec![make_request("a")]; + let executor = ok_executor(); + let calls = Arc::new(Mutex::new(Vec::new())); + let mock_sleep = MockSleep { calls: Arc::clone(&calls) }; + + block_on(process_requests_incremental( + requests, + false, + 0, + |_, _, _| true, + &executor, + mock_sleep, + )) + .unwrap(); + + assert_eq!(calls.lock().unwrap().len(), 0); + } + + #[test] + fn test_process_requests_pre_delay() { + let requests = vec![HttpRequest { + pre_delay_ms: Some(5), + ..make_request("req1") + }]; + let calls = Arc::new(Mutex::new(Vec::new())); + let mock_sleep = MockSleep { calls: Arc::clone(&calls) }; + let executed = Arc::new(AtomicBool::new(false)); + let exec_flag = Arc::clone(&executed); + + let exec = move |_: HttpRequest, _: bool, _: bool| { + exec_flag.store(true, Ordering::SeqCst); + Box::pin(async move { Ok(make_result(None)) }) + as Pin>>> + }; + + block_on(process_requests_incremental( + requests, + false, + 0, + |_, _, _| true, + &exec, + mock_sleep, + )) + .unwrap(); + + assert!(executed.load(Ordering::SeqCst)); + let calls = calls.lock().unwrap(); + assert!(calls.iter().any(|d| *d == Duration::from_millis(5))); + } + + #[test] + fn test_process_requests_post_delay() { + let requests = vec![HttpRequest { + post_delay_ms: Some(5), + ..make_request("req1") + }]; + let calls = Arc::new(Mutex::new(Vec::new())); + let mock_sleep = MockSleep { calls: Arc::clone(&calls) }; + let executor = ok_executor(); + + block_on(process_requests_incremental( + requests, + false, + 0, + |_, _, _| true, + &executor, + mock_sleep, + )) + .unwrap(); + + let calls = calls.lock().unwrap(); + assert!(calls.iter().any(|d| *d == Duration::from_millis(5))); + } + + #[test] + fn test_process_requests_callback_receives_correct_indices() { + let executor = ok_executor(); + let results = Arc::new(Mutex::new(Vec::new())); + let r = Arc::clone(&results); + + let requests = vec![ + make_request("a"), + HttpRequest { + depends_on: Some("nonexistent".to_string()), + ..make_request("b") + }, + make_request("c"), + ]; + + block_on(process_requests_incremental( + requests, + false, + 0, + |idx, total, result| { + r.lock().unwrap().push((idx, total, result)); + true + }, + &executor, + MockSleep::new(), + )) + .unwrap(); + + let results = results.lock().unwrap(); + assert_eq!(results.len(), 3); + assert_eq!(results[0].0, 0); + assert_eq!(results[1].0, 1); + assert_eq!(results[2].0, 2); + assert_eq!(results[0].1, 3); + assert!(matches!(&results[1].2, RequestProcessingResult::Skipped { .. })); + assert!(matches!(&results[2].2, RequestProcessingResult::Executed { .. })); + } +} diff --git a/src/core/src/processor/mod.rs b/src/core/src/processor/mod.rs index 963ec952..ed421005 100644 --- a/src/core/src/processor/mod.rs +++ b/src/core/src/processor/mod.rs @@ -1,6 +1,7 @@ mod executor; mod formatter; mod incremental; +pub(crate) mod incremental_loop; mod output; pub use executor::{ diff --git a/src/core/src/runner/incremental_async.rs b/src/core/src/runner/incremental_async.rs index a8296492..9ef893ce 100644 --- a/src/core/src/runner/incremental_async.rs +++ b/src/core/src/runner/incremental_async.rs @@ -1,324 +1,48 @@ -use crate::assertions; -use crate::conditions; -use crate::request_substitution::{ - substitute_functions_in_request, substitute_request_variables_in_request, +use crate::processor::incremental_loop::{ + AsyncSleep, process_requests_incremental, RequestProcessingResult, }; -use crate::types::{HttpRequest, HttpResult, RequestContext}; use anyhow::Result; -use std::future::Future; -use std::pin::Pin; -#[cfg(not(target_arch = "wasm32"))] -use std::sync::{ - Arc, Mutex, - atomic::{AtomicBool, Ordering}, + +/// Re-export types from the unified loop for the async path. +pub use crate::processor::incremental_loop::{ + AsyncRequestExecutor, AsyncRequestFuture, }; -#[cfg(not(target_arch = "wasm32"))] -use std::task::Waker; -use std::time::Duration; - -pub type AsyncRequestFuture<'a> = Pin> + 'a>>; -pub type AsyncRequestExecutor = - dyn for<'a> Fn(&'a HttpRequest, bool, bool) -> AsyncRequestFuture<'a>; - -/// Result of processing a single request during async incremental execution. -#[derive(Debug)] -pub enum AsyncRequestProcessingResult { - /// Request was skipped due to conditions or dependencies. - Skipped { - request: HttpRequest, - reason: String, - }, - /// Request was executed successfully or with errors. - Executed { - request: HttpRequest, - result: HttpResult, - }, - /// Request processing failed before execution. - Failed { request: HttpRequest, error: String }, -} +pub use crate::processor::RequestProcessingResult as AsyncRequestProcessingResult; -/// Process already-parsed HTTP requests incrementally while preserving the same -/// dependency, condition, and substitution semantics as the native processor. +/// Process already-parsed HTTP requests incrementally with async sleep. pub async fn process_http_requests_incremental_async( - requests: Vec, + requests: Vec, insecure: bool, delay_ms: u64, - mut callback: F, + callback: F, executor: &AsyncRequestExecutor, ) -> Result<()> where - F: FnMut(usize, usize, AsyncRequestProcessingResult) -> bool, + F: FnMut(usize, usize, RequestProcessingResult) -> bool, { - let total = requests.len(); - - if requests.is_empty() { - return Ok(()); - } - - let mut request_contexts: Vec = Vec::new(); - - for (idx, mut request) in requests.into_iter().enumerate() { - let request_count = (idx + 1) as u32; - - if idx > 0 && delay_ms > 0 { - sleep_ms(delay_ms).await; - } - - if let Some(dep_name) = request.depends_on.as_ref() - && !conditions::check_dependency(&Some(dep_name.clone()), &request_contexts) - { - let should_continue = callback( - idx, - total, - AsyncRequestProcessingResult::Skipped { - request: request.clone(), - reason: format!("Dependency on '{}' not met", dep_name), - }, - ); - add_request_context(&mut request_contexts, request, None, request_count); - if !should_continue { - break; - } - continue; - } - - if !request.conditions.is_empty() { - match conditions::evaluate_conditions(&request.conditions, &request_contexts) { - Ok(true) => {} - Ok(false) => { - let should_continue = callback( - idx, - total, - AsyncRequestProcessingResult::Skipped { - request: request.clone(), - reason: "Conditions not met".to_string(), - }, - ); - add_request_context(&mut request_contexts, request, None, request_count); - if !should_continue { - break; - } - continue; - } - Err(error) => { - let should_continue = callback( - idx, - total, - AsyncRequestProcessingResult::Failed { - request: request.clone(), - error: format!("Condition evaluation error: {}", error), - }, - ); - add_request_context(&mut request_contexts, request, None, request_count); - if !should_continue { - break; - } - continue; - } - } - } - - if let Err(error) = substitute_request_variables_in_request(&mut request, &request_contexts) - { - let should_continue = callback( - idx, - total, - AsyncRequestProcessingResult::Failed { - request: request.clone(), - error: format!("Variable substitution error: {}", error), - }, - ); - add_request_context(&mut request_contexts, request, None, request_count); - if !should_continue { - break; - } - continue; - } - - if let Err(error) = substitute_functions_in_request(&mut request) { - let should_continue = callback( - idx, - total, - AsyncRequestProcessingResult::Failed { - request: request.clone(), - error: format!("Function substitution error: {}", error), - }, - ); - add_request_context(&mut request_contexts, request, None, request_count); - if !should_continue { - break; - } - continue; - } - - if let Some(pre_delay_ms) = request.pre_delay_ms - && pre_delay_ms > 0 - { - sleep_ms(pre_delay_ms).await; - } - - let post_delay_ms = request.post_delay_ms; - - match executor(&request, false, insecure).await { - Ok(mut result) => { - if !request.assertions.is_empty() { - let assertion_results = - assertions::evaluate_assertions(&request.assertions, &result); - let all_passed = assertion_results.iter().all(|r| r.passed); - result.success = all_passed; - result.assertion_results = assertion_results; - } - add_request_context( - &mut request_contexts, - request.clone(), - Some(result.clone()), - request_count, - ); - let should_continue = callback( - idx, - total, - AsyncRequestProcessingResult::Executed { request, result }, - ); - if !should_continue { - break; - } - } - Err(error) => { - let should_continue = callback( - idx, - total, - AsyncRequestProcessingResult::Failed { - request: request.clone(), - error: error.to_string(), - }, - ); - add_request_context(&mut request_contexts, request, None, request_count); - if !should_continue { - break; - } - } - } - - if let Some(post_delay_ms) = post_delay_ms - && post_delay_ms > 0 - { - sleep_ms(post_delay_ms).await; - } - } - - Ok(()) -} - -fn add_request_context( - contexts: &mut Vec, - request: HttpRequest, - result: Option, - request_count: u32, -) { - let context_name = request - .name - .clone() - .unwrap_or_else(|| format!("request_{}", request_count)); - - contexts.push(RequestContext { - name: context_name, - request, - result, - }); -} - -#[cfg(target_arch = "wasm32")] -async fn sleep_ms(delay_ms: u64) { - if delay_ms == 0 { - return; - } - - gloo_timers::future::sleep(Duration::from_millis(delay_ms)).await; -} - -#[cfg(not(target_arch = "wasm32"))] -async fn sleep_ms(delay_ms: u64) { - if delay_ms == 0 { - return; - } - - NativeSleep::new(Duration::from_millis(delay_ms)).await; -} - -#[cfg(not(target_arch = "wasm32"))] -struct NativeSleep { - duration: Option, - state: Arc, -} - -#[cfg(not(target_arch = "wasm32"))] -struct NativeSleepState { - completed: AtomicBool, - waker: Mutex>, -} - -#[cfg(not(target_arch = "wasm32"))] -impl NativeSleep { - fn new(duration: Duration) -> Self { - Self { - duration: Some(duration), - state: Arc::new(NativeSleepState { - completed: AtomicBool::new(false), - waker: Mutex::new(None), - }), - } - } -} - -#[cfg(not(target_arch = "wasm32"))] -impl Future for NativeSleep { - type Output = (); - - fn poll( - self: Pin<&mut Self>, - cx: &mut std::task::Context<'_>, - ) -> std::task::Poll { - let this = self.get_mut(); - - if this.state.completed.load(Ordering::Acquire) { - return std::task::Poll::Ready(()); - } - - { - let mut waker = this - .state - .waker - .lock() - .expect("native sleep waker mutex poisoned"); - *waker = Some(cx.waker().clone()); - } - - if let Some(duration) = this.duration.take() { - let state = Arc::clone(&this.state); - std::thread::spawn(move || { - std::thread::sleep(duration); - state.completed.store(true, Ordering::Release); - if let Some(waker) = state - .waker - .lock() - .expect("native sleep waker mutex poisoned") - .take() - { - waker.wake(); - } - }); - } - - std::task::Poll::Pending - } + let wrapped = |request: crate::types::HttpRequest, verbose: bool, insecure: bool| { + async move { executor(&request, verbose, insecure).await } + }; + process_requests_incremental( + requests, + insecure, + delay_ms, + callback, + &wrapped, + AsyncSleep, + ) + .await } #[cfg(test)] mod tests { use super::*; - use crate::types::{Assertion, AssertionType, Condition, ConditionType, Header}; + use crate::types::{ + Assertion, AssertionType, Condition, ConditionType, Header, HttpRequest, HttpResult, + }; use std::collections::HashMap; - use std::task::{Context, Poll, RawWaker, RawWakerVTable, Waker}; + use std::future::Future; + use std::task::{Context, Poll}; #[test] fn test_async_incremental_preserves_request_variable_context() { @@ -366,7 +90,7 @@ mod tests { 0, |idx, _total, result| { if idx == 1 { - if let AsyncRequestProcessingResult::Executed { request, .. } = result { + if let RequestProcessingResult::Executed { request, .. } = result { captured_request = Some(request); } return false; @@ -429,7 +153,7 @@ mod tests { 0, |idx, _total, result| { if idx == 1 { - if let AsyncRequestProcessingResult::Skipped { reason, .. } = result { + if let RequestProcessingResult::Skipped { reason, .. } = result { dependency_reason = Some(reason); } return false; @@ -495,7 +219,7 @@ mod tests { 0, |idx, _total, result| { if idx == 1 { - if let AsyncRequestProcessingResult::Skipped { reason, .. } = result { + if let RequestProcessingResult::Skipped { reason, .. } = result { skipped_reason = Some(reason); } return false; @@ -545,7 +269,7 @@ mod tests { false, 0, |_idx, _total, result| { - if let AsyncRequestProcessingResult::Failed { error, .. } = result { + if let RequestProcessingResult::Failed { error, .. } = result { captured = Some(error); } false @@ -583,7 +307,7 @@ mod tests { false, 0, |_idx, _total, result| { - if let AsyncRequestProcessingResult::Executed { result, .. } = result { + if let RequestProcessingResult::Executed { result, .. } = result { captured_result = Some(result); } false @@ -624,7 +348,7 @@ mod tests { false, 0, |_idx, _total, result| { - if let AsyncRequestProcessingResult::Executed { result, .. } = result { + if let RequestProcessingResult::Executed { result, .. } = result { captured_result = Some(result); } false @@ -741,8 +465,8 @@ mod tests { } fn block_on(future: F) -> F::Output { - let waker = unsafe { Waker::from_raw(dummy_raw_waker()) }; - let mut context = Context::from_waker(&waker); + let waker = std::task::Waker::noop(); + let mut context = Context::from_waker(waker); let mut future = Box::pin(future); loop { @@ -752,17 +476,4 @@ mod tests { } } } - - unsafe fn dummy_clone(_: *const ()) -> RawWaker { - dummy_raw_waker() - } - - unsafe fn dummy_no_op(_: *const ()) {} - - static DUMMY_WAKER_VTABLE: RawWakerVTable = - RawWakerVTable::new(dummy_clone, dummy_no_op, dummy_no_op, dummy_no_op); - - fn dummy_raw_waker() -> RawWaker { - RawWaker::new(std::ptr::null(), &DUMMY_WAKER_VTABLE) - } } diff --git a/src/core/src/runner/mod.rs b/src/core/src/runner/mod.rs index 109c7684..14234cae 100644 --- a/src/core/src/runner/mod.rs +++ b/src/core/src/runner/mod.rs @@ -16,5 +16,7 @@ pub use incremental_async::{ process_http_requests_incremental_async, }; +pub use crate::processor::RequestProcessingResult; + #[cfg(target_arch = "wasm32")] pub use executor_async::execute_http_request_async; diff --git a/src/gui/src/results_view.rs b/src/gui/src/results_view.rs index 407797d3..a1dabf97 100644 --- a/src/gui/src/results_view.rs +++ b/src/gui/src/results_view.rs @@ -1,6 +1,4 @@ -#[cfg(not(target_arch = "wasm32"))] use httprunner_core::processor::RequestProcessingResult; -use httprunner_core::runner::AsyncRequestProcessingResult; #[cfg(not(target_arch = "wasm32"))] use httprunner_core::telemetry; use httprunner_core::types::AssertionResult; @@ -890,17 +888,17 @@ fn request_result_is_failure(result: &RequestProcessingResult) -> bool { /// Async (WASM) counterpart of [`should_continue_after`]. pub(crate) fn should_continue_after_async( - result: &AsyncRequestProcessingResult, + result: &RequestProcessingResult, fail_fast: bool, ) -> bool { !(fail_fast && async_result_is_failure(result)) } -fn async_result_is_failure(result: &AsyncRequestProcessingResult) -> bool { +fn async_result_is_failure(result: &RequestProcessingResult) -> bool { match result { - AsyncRequestProcessingResult::Executed { result, .. } => !result.success, - AsyncRequestProcessingResult::Failed { .. } => true, - AsyncRequestProcessingResult::Skipped { .. } => false, + RequestProcessingResult::Executed { result, .. } => !result.success, + RequestProcessingResult::Failed { .. } => true, + RequestProcessingResult::Skipped { .. } => false, } } @@ -991,19 +989,19 @@ mod tests { #[test] fn should_continue_after_async_matches_sync_semantics() { - let failed = AsyncRequestProcessingResult::Executed { + let failed = RequestProcessingResult::Executed { request: sample_request(), result: sample_result(false), }; - let skipped = AsyncRequestProcessingResult::Skipped { + let skipped = RequestProcessingResult::Skipped { request: sample_request(), reason: "dependency".to_string(), }; - let processing_failed = AsyncRequestProcessingResult::Failed { + let processing_failed = RequestProcessingResult::Failed { request: sample_request(), error: "boom".to_string(), }; - let ok = AsyncRequestProcessingResult::Executed { + let ok = RequestProcessingResult::Executed { request: sample_request(), result: sample_result(true), }; diff --git a/src/gui/src/results_view_async.rs b/src/gui/src/results_view_async.rs index e72f5389..be3adc76 100644 --- a/src/gui/src/results_view_async.rs +++ b/src/gui/src/results_view_async.rs @@ -4,8 +4,9 @@ use crate::results_view::{ }; use futures_util::FutureExt; use httprunner_core::parser; +use httprunner_core::processor::RequestProcessingResult; use httprunner_core::runner::{ - AsyncRequestFuture, AsyncRequestProcessingResult, process_http_requests_incremental_async, + AsyncRequestFuture, process_http_requests_incremental_async, }; use httprunner_core::types::{HttpRequest, HttpResult}; use std::any::Any; @@ -248,19 +249,19 @@ fn make_async_executor( } } -fn map_process_result(process_result: AsyncRequestProcessingResult) -> ExecutionResult { +fn map_process_result(process_result: RequestProcessingResult) -> ExecutionResult { match process_result { - AsyncRequestProcessingResult::Skipped { request, reason } => { + RequestProcessingResult::Skipped { request, reason } => { ExecutionResult::Failure(FailureResult::simple( format!("⏭️ {}", request.method), request.url, format!("Skipped: {}", reason), )) } - AsyncRequestProcessingResult::Executed { request, result } => { + RequestProcessingResult::Executed { request, result } => { map_http_result(request, result) } - AsyncRequestProcessingResult::Failed { request, error } => { + RequestProcessingResult::Failed { request, error } => { ExecutionResult::Failure(FailureResult::simple(request.method, request.url, error)) } }