Skip to main content

camel_processor/
do_try.rs

1//! ## Stop semantics (ADR-0025)
2//!
3//! This segment implements `OutcomePipeline` and propagates `PipelineOutcome::Stopped(ex)` with the exchange state intact (including mutations made inside the segment body before Stop fired). See ADR-0025 §3 (stopped-exchange-state-preservation invariant).
4//!
5//! Catch-block failure envelope (bd rc-zgbqq): if a catch body fails, the catch error stays the main error in every disposition; the original error lives in the `warn` record ("do_try catch block failed; catch error supersedes original", fields `original_error` and `catch_error`) and, when a span is active, in the span event.
6//! See the error-handler spec, requirement "delegate failure propagates the original error" (that rule covers `handled_by` delegates, not catch bodies).
7
8use camel_api::error_handler::ExceptionDisposition;
9use camel_api::exchange::PROPERTY_EXCEPTION_HANDLED;
10use camel_api::{BoxProcessor, CamelError, Exchange, PredicateSource};
11use tower::Service;
12use tower::ServiceExt;
13
14use crate::error_handler::record_span_error;
15
16/// Matcher for a `doCatch` clause.
17#[derive(Clone)]
18pub enum CatchMatcher {
19    /// Match by CamelError variant name (e.g. ["ProcessorError", "Io"]).
20    /// `"*"` matches any variant — equivalent to Camel's `doCatch(Throwable.class)`.
21    ByVariant(Vec<String>),
22    /// Match by a fallible predicate over the Exchange.
23    Predicate(PredicateSource),
24}
25
26impl CatchMatcher {
27    /// Returns whether this matcher matches the given error and exchange.
28    ///
29    /// A failed predicate evaluation yields `Err` (language-value-boundary:
30    /// the caller chains the original error and fails the scope).
31    pub async fn matches(&self, err: &CamelError, ex: &Exchange) -> Result<bool, CamelError> {
32        match self {
33            CatchMatcher::ByVariant(names) => {
34                if names.iter().any(|n| n == "*") {
35                    return Ok(true);
36                }
37                Ok(names.iter().any(|n| n == err.variant_name()))
38            }
39            CatchMatcher::Predicate(source) => source.matches(ex).await,
40        }
41    }
42}
43
44/// Chain the original try error into a failed catch-predicate error
45/// (language-value-boundary: catch `when` / `on_when` failure envelope).
46///
47/// Predicate evaluation errors arrive as `CamelError::ExpressionFailed` with
48/// `cause: None`. The copy keeps every field unchanged except `cause`, which
49/// holds the original error so the chain stays retrievable. Errors that carry
50/// no `cause` slot (or an existing chain) pass through untouched.
51pub(crate) fn chain_predicate_error(
52    mut predicate_err: CamelError,
53    original: CamelError,
54) -> CamelError {
55    if let CamelError::ExpressionFailed { cause, .. } = &mut predicate_err
56        && cause.is_none()
57    {
58        *cause = Some(Box::new(original));
59    }
60    predicate_err
61}
62
63/// A single `doCatch` clause.
64#[derive(Clone)]
65pub struct CatchClause {
66    /// The main matcher (variant-name list or predicate).
67    pub matcher: CatchMatcher,
68    /// Optional sub-predicate evaluated AFTER the main matcher passes.
69    pub on_when: Option<PredicateSource>,
70    /// Sub-pipeline executed when the clause matches.
71    pub steps: Vec<BoxProcessor>,
72    /// ADR-0019 disposition: Handled (default), Propagate, or Continued (rejected at parse time).
73    /// In YAML, use lowercase: `handled`, `propagate`, `continued`.
74    pub disposition: ExceptionDisposition,
75}
76
77/// The `doTry` processor. Wrap with `BoxProcessor::new(DoTryService::new(...))`.
78#[derive(Clone)]
79pub struct DoTryService {
80    /// Steps in the try block.
81    pub try_steps: Vec<BoxProcessor>,
82    /// Catch clauses evaluated first-match-wins.
83    pub catch_clauses: Vec<CatchClause>,
84    /// Steps in the finally block (empty = no finally).
85    pub finally_steps: Vec<BoxProcessor>,
86    /// Optional onWhen predicate for finally.
87    pub finally_on_when: Option<PredicateSource>,
88}
89
90impl DoTryService {
91    /// Create a new `DoTryService` with the given try steps.
92    pub fn new(try_steps: Vec<BoxProcessor>) -> Self {
93        Self {
94            try_steps,
95            catch_clauses: Vec::new(),
96            finally_steps: Vec::new(),
97            finally_on_when: None,
98        }
99    }
100
101    /// Full constructor used by the compile pipeline (Task 10b control_flow.rs).
102    /// Builder API (Task 8) constructs via `new()` + field mutation.
103    pub fn with_catch_and_finally(
104        try_steps: Vec<BoxProcessor>,
105        catch_clauses: Vec<CatchClause>,
106        finally_steps: Vec<BoxProcessor>,
107        finally_on_when: Option<PredicateSource>,
108    ) -> Self {
109        Self {
110            try_steps,
111            catch_clauses,
112            finally_steps,
113            finally_on_when,
114        }
115    }
116}
117
118/// Run a sequence of steps, preserving the last exchange state on error.
119/// Returns `Err(Box<(last_ex, err)>)` so DoTry can populate exception properties.
120async fn run_pipeline(
121    steps: Vec<BoxProcessor>,
122    mut ex: Exchange,
123) -> Result<Exchange, Box<(Exchange, CamelError)>> {
124    for mut svc in steps {
125        match svc.ready().await {
126            Ok(ready) => {
127                let snapshot = ex.clone();
128                match ready.call(ex).await {
129                    Ok(new_ex) => ex = new_ex,
130                    Err(err) => return Err(Box::new((snapshot, err))),
131                }
132            }
133            Err(err) => return Err(Box::new((ex, err))),
134        }
135    }
136    Ok(ex)
137}
138
139/// Run the finally block. Camel parity for finally-throws:
140/// - If finally succeeds (or is skipped): `Completed` with its exchange.
141/// - If finally throws AND there was a previous error: `Restore` (caller
142///   logs and restores the previous error).
143/// - If finally throws AND no previous error: `NoPreviousFail` (caller logs
144///   and propagates finally_err).
145/// - If the finally `on_when` predicate itself fails: `Err` (the typed
146///   predicate error fails the scope).
147///
148/// Logging lives at the CALLERS so the restore record's field names can
149/// match the calling flow (`previous_error` vs `catch_error`, bd rc-zgbqq).
150async fn run_finally(
151    finally_steps: Vec<BoxProcessor>,
152    finally_on_when: Option<PredicateSource>,
153    ex: Exchange,
154    previous_err: Option<CamelError>,
155) -> Result<FinallyOutcome, CamelError> {
156    if finally_steps.is_empty() {
157        return Ok(FinallyOutcome::Completed(ex));
158    }
159    if let Some(on_when) = &finally_on_when
160        && !on_when.matches(&ex).await?
161    {
162        return Ok(FinallyOutcome::Completed(ex));
163    }
164    match run_pipeline(finally_steps, ex).await {
165        Ok(ex) => Ok(FinallyOutcome::Completed(ex)),
166        Err(failed) => {
167            let (_, finally_err) = *failed;
168            Ok(match previous_err {
169                Some(prev) => FinallyOutcome::Restore {
170                    previous: prev,
171                    finally_err,
172                },
173                None => FinallyOutcome::NoPreviousFail(finally_err),
174            })
175        }
176    }
177}
178
179/// Result of `run_finally`. Tower-local to this module — unrelated to the
180/// identically-named enum in `do_try_segment.rs`.
181// The variant shape is fixed by the sealed do_try ruling (bd rc-zgbqq);
182// Exchange simply dwarfs the error payloads, so allow the lint instead of
183// boxing and drifting from the specified envelope.
184#[allow(clippy::large_enum_variant)]
185enum FinallyOutcome {
186    /// Finally ran (or was skipped) successfully; carries its exchange.
187    Completed(Exchange),
188    /// Finally threw with no previous error; the finally error wins.
189    NoPreviousFail(CamelError),
190    /// Finally threw with a previous error; the previous error is restored.
191    Restore {
192        previous: CamelError,
193        finally_err: CamelError,
194    },
195}
196
197/// Run finally after a catch-predicate failure, then surface `failure`.
198///
199/// A failed finally `on_when` predicate supersedes `failure`; a throwing
200/// finally body restores `failure` (Camel parity with the other restore
201/// flows).
202async fn run_finally_for_failure(
203    finally_steps: Vec<BoxProcessor>,
204    finally_on_when: Option<PredicateSource>,
205    ex: Exchange,
206    failure: CamelError,
207) -> CamelError {
208    match run_finally(finally_steps, finally_on_when, ex, Some(failure.clone())).await {
209        Err(pred_err) => pred_err,
210        Ok(FinallyOutcome::Restore {
211            previous,
212            finally_err,
213        }) => {
214            tracing::warn!(
215                finally_error = %finally_err,
216                previous_error = %previous,
217                "doFinally threw after catch predicate failure; restoring previous (Camel parity)"
218            );
219            previous
220        }
221        // NoPreviousFail is unreachable here (previous is Some); Completed
222        // falls through to the chained failure.
223        Ok(_) => failure,
224    }
225}
226
227impl tower::Service<Exchange> for DoTryService {
228    type Response = Exchange;
229    type Error = CamelError;
230    type Future = std::pin::Pin<
231        Box<dyn std::future::Future<Output = Result<Self::Response, Self::Error>> + Send>,
232    >;
233
234    fn poll_ready(
235        &mut self,
236        _cx: &mut std::task::Context<'_>,
237    ) -> std::task::Poll<Result<(), Self::Error>> {
238        std::task::Poll::Ready(Ok(()))
239    }
240
241    fn call(&mut self, mut exchange: Exchange) -> Self::Future {
242        // Clear stale CamelExceptionHandled marker from prior handlers in the same route.
243        // Exchange::clear_error() does NOT touch HANDLED, so direct map access is required.
244        exchange.properties.remove(PROPERTY_EXCEPTION_HANDLED);
245
246        let try_steps = self.try_steps.clone();
247        let catch_clauses = self.catch_clauses.clone();
248        let finally_steps = self.finally_steps.clone();
249        let finally_on_when = self.finally_on_when.clone();
250
251        Box::pin(async move {
252            let try_result = run_pipeline(try_steps, exchange).await;
253            match try_result {
254                Ok(ex) => match run_finally(finally_steps, finally_on_when, ex, None).await {
255                    Ok(FinallyOutcome::Completed(ex)) => Ok(ex),
256                    Ok(FinallyOutcome::NoPreviousFail(fin)) => {
257                        tracing::warn!(error = %fin, "doFinally threw");
258                        Err(fin)
259                    }
260                    // Unreachable: the try-Ok flow passes no previous error.
261                    Ok(FinallyOutcome::Restore { previous, .. }) => Err(previous),
262                    // Failed finally `on_when` predicate: the typed error fails
263                    // the scope.
264                    Err(pred_err) => Err(pred_err),
265                },
266                Err(failed) => {
267                    let (failed_ex, original_err) = *failed;
268                    let mut ex = failed_ex;
269                    ex.set_error(original_err.clone());
270
271                    for clause in catch_clauses {
272                        let CatchClause {
273                            matcher,
274                            on_when,
275                            steps,
276                            disposition,
277                        } = clause;
278                        let matched = match matcher.matches(&original_err, &ex).await {
279                            Ok(m) => m,
280                            // Failed catch `when` predicate: replaces the original
281                            // error (preserved as `cause`), fails the scope after
282                            // finally runs.
283                            Err(pred_err) => {
284                                let chained = chain_predicate_error(pred_err, original_err);
285                                return Err(run_finally_for_failure(
286                                    finally_steps,
287                                    finally_on_when,
288                                    ex,
289                                    chained,
290                                )
291                                .await);
292                            }
293                        };
294                        if !matched {
295                            continue;
296                        }
297                        if let Some(ref on_when) = on_when {
298                            match on_when.matches(&ex).await {
299                                Ok(true) => {}
300                                Ok(false) => continue,
301                                // Failed catch `on_when` predicate: same envelope
302                                // as a failed `when` predicate.
303                                Err(pred_err) => {
304                                    let chained = chain_predicate_error(pred_err, original_err);
305                                    return Err(run_finally_for_failure(
306                                        finally_steps,
307                                        finally_on_when,
308                                        ex,
309                                        chained,
310                                    )
311                                    .await);
312                                }
313                            }
314                        }
315
316                        let catch_result = run_pipeline(steps, ex.clone()).await;
317
318                        return match catch_result {
319                            Ok(ok_ex) => {
320                                // Determine previous-error threading based on disposition.
321                                // IMPORTANT: do NOT call handle_error() before run_finally() —
322                                // handle_error() calls clear_error() which removes
323                                // PROPERTY_EXCEPTION_MESSAGE/KIND/CAUGHT, preventing finally
324                                // steps from inspecting the caught exception.
325                                //
326                                // disposition semantics (ADR-0019 strict):
327                                //   Handled    -> catch output is final, no propagation
328                                //   Propagate  -> catch ran for side-effects, original propagates
329                                //   Continued  -> rejected at parse time (defensive: treat as
330                                //                 Propagate + log if we ever reach runtime)
331                                let prev = match disposition {
332                                    ExceptionDisposition::Handled => None,
333                                    ExceptionDisposition::Continued => {
334                                        tracing::warn!(
335                                            "ExceptionDisposition::Continued reached doTry runtime; \
336                                             treating as Propagate. Should have been rejected at parse time."
337                                        );
338                                        Some(original_err.clone())
339                                    }
340                                    // Propagate and any future variant thread the original error.
341                                    _ => Some(original_err.clone()),
342                                };
343                                let mut ex = match run_finally(
344                                    finally_steps.clone(),
345                                    finally_on_when.clone(),
346                                    ok_ex,
347                                    prev,
348                                )
349                                .await
350                                {
351                                    Ok(FinallyOutcome::Completed(ex)) => ex,
352                                    Ok(FinallyOutcome::NoPreviousFail(fin)) => {
353                                        tracing::warn!(error = %fin, "doFinally threw");
354                                        return Err(fin);
355                                    }
356                                    Ok(FinallyOutcome::Restore {
357                                        previous,
358                                        finally_err,
359                                    }) => {
360                                        tracing::warn!(
361                                            finally_error = %finally_err,
362                                            previous_error = %previous,
363                                            "doFinally threw; restoring previous exception (Camel parity)"
364                                        );
365                                        return Err(previous);
366                                    }
367                                    // Failed finally `on_when` predicate: the typed
368                                    // error fails the scope.
369                                    Err(pred_err) => return Err(pred_err),
370                                };
371                                // AFTER finally has run (and had access to exception props),
372                                // apply handle_error() for Handled disposition to clear the
373                                // error state and set CamelExceptionHandled=true marker.
374                                if matches!(disposition, ExceptionDisposition::Handled) {
375                                    ex.handle_error();
376                                }
377                                match disposition {
378                                    ExceptionDisposition::Handled => Ok(ex),
379                                    _ => Err(original_err),
380                                }
381                            }
382                            Err(failed) => {
383                                // Catch threw. Sealed failure envelope (bd rc-zgbqq):
384                                // the catch error stays the main error (returned
385                                // below and threaded into finally); the original is
386                                // surfaced through the unconditional warn record
387                                // and, best-effort, the active span. The event
388                                // lands on the current span when one is entered —
389                                // no new span, no second event.
390                                let (catch_ex, catch_err) = *failed;
391                                tracing::warn!(
392                                    original_error = %original_err,
393                                    catch_error = %catch_err,
394                                    "do_try catch block failed; catch error supersedes original"
395                                );
396                                record_span_error(&catch_err);
397                                // Run finally with previous=catch_err. Per Camel
398                                // parity, if finally itself throws, catch_err is
399                                // restored.
400                                let outcome = run_finally(
401                                    finally_steps.clone(),
402                                    finally_on_when.clone(),
403                                    catch_ex,
404                                    Some(catch_err.clone()),
405                                )
406                                .await;
407                                match outcome {
408                                    Err(pred_err) => return Err(pred_err),
409                                    Ok(FinallyOutcome::Restore {
410                                        previous,
411                                        finally_err,
412                                    }) => {
413                                        tracing::warn!(
414                                            catch_error = %previous,
415                                            finally_error = %finally_err,
416                                            "doFinally threw after failed catch; restoring catch error"
417                                        );
418                                        return Err(previous);
419                                    }
420                                    Ok(_) => {}
421                                }
422                                Err(catch_err)
423                            }
424                        };
425                    }
426
427                    // No catch matched. Run finally with previous=original. Propagate original.
428                    match run_finally(
429                        finally_steps,
430                        finally_on_when,
431                        ex,
432                        Some(original_err.clone()),
433                    )
434                    .await
435                    {
436                        Ok(FinallyOutcome::Restore {
437                            previous,
438                            finally_err,
439                        }) => {
440                            tracing::warn!(
441                                finally_error = %finally_err,
442                                previous_error = %previous,
443                                "doFinally threw; restoring previous exception (Camel parity)"
444                            );
445                            Err(previous)
446                        }
447                        Ok(_) => Err(original_err),
448                        Err(pred_err) => Err(pred_err),
449                    }
450                }
451            }
452        })
453    }
454}
455
456// ── DoTrySegment (ADR-0025 OutcomePipeline) ──────────────────────────────
457
458/// Compilable segment for a `doCatch` clause within a `DoTrySegment`.
459///
460/// `disposition` controls outcome when the catch body completes:
461#[cfg(test)]
462mod tests {
463    use super::*;
464    use crate::test_log_capture::{capture_debugs_with_span_records, captured_field, record_field};
465    use camel_api::{BoxProcessor, BoxProcessorExt};
466    use std::sync::Arc;
467    use std::sync::atomic::{AtomicU32, Ordering};
468
469    fn passthrough() -> BoxProcessor {
470        BoxProcessor::from_fn(move |ex| Box::pin(async move { Ok(ex) }))
471    }
472
473    fn record_call(flag: Arc<AtomicU32>) -> BoxProcessor {
474        BoxProcessor::from_fn(move |ex| {
475            let f = flag.clone();
476            Box::pin(async move {
477                f.fetch_add(1, Ordering::SeqCst);
478                Ok(ex)
479            })
480        })
481    }
482
483    fn always_fail(err: CamelError) -> BoxProcessor {
484        BoxProcessor::from_fn(move |_ex| {
485            let e = err.clone();
486            Box::pin(async move { Err(e) })
487        })
488    }
489
490    #[tokio::test]
491    async fn happy_path_try_succeeds_finally_runs() {
492        let finally_flag = Arc::new(AtomicU32::new(0));
493        let mut svc = DoTryService::new(vec![passthrough()]);
494        svc.finally_steps = vec![record_call(finally_flag.clone())];
495
496        let mut boxed = BoxProcessor::new(svc);
497        let result = boxed.ready().await.unwrap().call(Exchange::default()).await;
498        assert!(result.is_ok());
499        assert_eq!(finally_flag.load(Ordering::SeqCst), 1);
500    }
501
502    #[tokio::test]
503    async fn catch_by_variant_handled_returns_ok() {
504        let try_step = always_fail(CamelError::ProcessorError("boom".into()));
505        let mut svc = DoTryService::new(vec![try_step]);
506        svc.catch_clauses.push(CatchClause {
507            matcher: CatchMatcher::ByVariant(vec!["ProcessorError".into()]),
508            on_when: None,
509            steps: vec![passthrough()],
510            disposition: ExceptionDisposition::Handled,
511        });
512
513        let mut boxed = BoxProcessor::new(svc);
514        let result = boxed.ready().await.unwrap().call(Exchange::default()).await;
515        assert!(result.is_ok(), "Handled must return Ok");
516        let ex = result.unwrap();
517        assert_eq!(
518            ex.properties.get(PROPERTY_EXCEPTION_HANDLED),
519            Some(&camel_api::Value::Bool(true)),
520            "CamelExceptionHandled must be set via handle_error()"
521        );
522    }
523
524    #[tokio::test]
525    async fn catch_by_variant_propagate_runs_side_effects_and_rethrows() {
526        let original = CamelError::ProcessorError("boom".into());
527        let try_step = always_fail(original.clone());
528        let side_effect = Arc::new(AtomicU32::new(0));
529        let catch_step = record_call(side_effect.clone());
530        let mut svc = DoTryService::new(vec![try_step]);
531        svc.catch_clauses.push(CatchClause {
532            matcher: CatchMatcher::ByVariant(vec!["ProcessorError".into()]),
533            on_when: None,
534            steps: vec![catch_step],
535            disposition: ExceptionDisposition::Propagate,
536        });
537
538        let mut boxed = BoxProcessor::new(svc);
539        let result = boxed.ready().await.unwrap().call(Exchange::default()).await;
540        assert!(result.is_err(), "Propagate must rethrow original");
541        assert!(matches!(result.unwrap_err(), CamelError::ProcessorError(_)));
542        assert_eq!(
543            side_effect.load(Ordering::SeqCst),
544            1,
545            "catch branch must have run for side-effects"
546        );
547    }
548
549    #[tokio::test]
550    async fn catch_by_predicate_matches_via_exception_kind() {
551        let try_step = always_fail(CamelError::Io("disk full".into()));
552        let predicate = FilterPredicate::new(|ex: &Exchange| {
553            ex.properties
554                .get(camel_api::exchange::PROPERTY_EXCEPTION_KIND)
555                .map(|v| matches!(v, camel_api::Value::String(s) if s == "io"))
556                .unwrap_or(false)
557        });
558        let mut svc = DoTryService::new(vec![try_step]);
559        svc.catch_clauses.push(CatchClause {
560            matcher: CatchMatcher::Predicate(PredicateSource::Sync(predicate)),
561            on_when: None,
562            steps: vec![passthrough()],
563            disposition: ExceptionDisposition::Handled,
564        });
565
566        let mut boxed = BoxProcessor::new(svc);
567        let result = boxed.ready().await.unwrap().call(Exchange::default()).await;
568        assert!(
569            result.is_ok(),
570            "Predicate matcher must catch the error and Handled must return Ok"
571        );
572    }
573
574    #[tokio::test]
575    async fn on_when_filters_clause_and_next_evaluated() {
576        let try_step = always_fail(CamelError::ProcessorError("boom".into()));
577        let first_call = Arc::new(AtomicU32::new(0));
578        let second_call = Arc::new(AtomicU32::new(0));
579
580        let mut svc = DoTryService::new(vec![try_step]);
581        svc.catch_clauses.push(CatchClause {
582            matcher: CatchMatcher::ByVariant(vec!["ProcessorError".into()]),
583            on_when: Some(PredicateSource::Sync(FilterPredicate::new(|_ex| false))),
584            steps: vec![record_call(first_call.clone())],
585            disposition: ExceptionDisposition::Handled,
586        });
587        svc.catch_clauses.push(CatchClause {
588            matcher: CatchMatcher::ByVariant(vec!["*".into()]),
589            on_when: None,
590            steps: vec![record_call(second_call.clone())],
591            disposition: ExceptionDisposition::Handled,
592        });
593
594        let mut boxed = BoxProcessor::new(svc);
595        let _ = boxed.ready().await.unwrap().call(Exchange::default()).await;
596        assert_eq!(first_call.load(Ordering::SeqCst), 0);
597        assert_eq!(second_call.load(Ordering::SeqCst), 1);
598    }
599
600    #[tokio::test]
601    async fn first_match_wins_subsequent_clauses_not_evaluated() {
602        let try_step = always_fail(CamelError::Io("err".into()));
603        let first_call = Arc::new(AtomicU32::new(0));
604        let second_call = Arc::new(AtomicU32::new(0));
605
606        let mut svc = DoTryService::new(vec![try_step]);
607        svc.catch_clauses.push(CatchClause {
608            matcher: CatchMatcher::ByVariant(vec!["Io".into()]),
609            on_when: None,
610            steps: vec![record_call(first_call.clone())],
611            disposition: ExceptionDisposition::Handled,
612        });
613        svc.catch_clauses.push(CatchClause {
614            matcher: CatchMatcher::ByVariant(vec!["*".into()]),
615            on_when: None,
616            steps: vec![record_call(second_call.clone())],
617            disposition: ExceptionDisposition::Handled,
618        });
619
620        let mut boxed = BoxProcessor::new(svc);
621        let _ = boxed.ready().await.unwrap().call(Exchange::default()).await;
622        assert_eq!(first_call.load(Ordering::SeqCst), 1);
623        assert_eq!(second_call.load(Ordering::SeqCst), 0);
624    }
625
626    #[tokio::test]
627    async fn no_clause_matches_propagates_original() {
628        let try_step = always_fail(CamelError::CircuitOpen("cb".into()));
629        let mut svc = DoTryService::new(vec![try_step]);
630        svc.catch_clauses.push(CatchClause {
631            matcher: CatchMatcher::ByVariant(vec!["ProcessorError".into()]),
632            on_when: None,
633            steps: vec![passthrough()],
634            disposition: ExceptionDisposition::Handled,
635        });
636
637        let mut boxed = BoxProcessor::new(svc);
638        let result = boxed.ready().await.unwrap().call(Exchange::default()).await;
639        assert!(result.is_err());
640        assert!(matches!(result.unwrap_err(), CamelError::CircuitOpen(_)));
641    }
642
643    #[tokio::test]
644    async fn catch_branch_throws_new_error_wins() {
645        let try_step = always_fail(CamelError::ProcessorError("orig".into()));
646        let catch_step = always_fail(CamelError::Io("catch-fail".into()));
647        let mut svc = DoTryService::new(vec![try_step]);
648        svc.catch_clauses.push(CatchClause {
649            matcher: CatchMatcher::ByVariant(vec!["ProcessorError".into()]),
650            on_when: None,
651            steps: vec![catch_step],
652            disposition: ExceptionDisposition::Handled,
653        });
654
655        let mut boxed = BoxProcessor::new(svc);
656        let result = boxed.ready().await.unwrap().call(Exchange::default()).await;
657        assert!(result.is_err());
658        assert!(matches!(result.unwrap_err(), CamelError::Io(_)));
659    }
660
661    #[tokio::test]
662    async fn finally_throws_with_no_previous_error_propagates_finally_error() {
663        let finally_step = always_fail(CamelError::Config("fin".into()));
664        let mut svc = DoTryService::new(vec![passthrough()]);
665        svc.finally_steps = vec![finally_step];
666
667        let mut boxed = BoxProcessor::new(svc);
668        let result = boxed.ready().await.unwrap().call(Exchange::default()).await;
669        assert!(result.is_err());
670        assert!(matches!(result.unwrap_err(), CamelError::Config(_)));
671    }
672
673    #[tokio::test]
674    async fn finally_throws_with_previous_error_restores_previous() {
675        let try_step = always_fail(CamelError::ProcessorError("orig".into()));
676        let finally_step = always_fail(CamelError::Config("fin".into()));
677        let mut svc = DoTryService::new(vec![try_step]);
678        svc.finally_steps = vec![finally_step];
679
680        let mut boxed = BoxProcessor::new(svc);
681        let result = boxed.ready().await.unwrap().call(Exchange::default()).await;
682        assert!(result.is_err());
683        assert!(
684            matches!(result.unwrap_err(), CamelError::ProcessorError(_)),
685            "previous error must be restored when finally throws (Camel parity)"
686        );
687    }
688
689    #[tokio::test]
690    async fn finally_on_when_false_skips_finally() {
691        let finally_call = Arc::new(AtomicU32::new(0));
692        let mut svc = DoTryService::new(vec![passthrough()]);
693        svc.finally_steps = vec![record_call(finally_call.clone())];
694        svc.finally_on_when = Some(PredicateSource::Sync(FilterPredicate::new(|_ex| false)));
695
696        let mut boxed = BoxProcessor::new(svc);
697        let _ = boxed.ready().await.unwrap().call(Exchange::default()).await;
698        assert_eq!(finally_call.load(Ordering::SeqCst), 0);
699    }
700
701    #[tokio::test]
702    async fn stale_handled_marker_cleared_on_entry() {
703        let mut ex = Exchange::default();
704        ex.set_property(PROPERTY_EXCEPTION_HANDLED, camel_api::Value::Bool(true));
705        let svc = DoTryService::new(vec![passthrough()]);
706        let mut boxed = BoxProcessor::new(svc);
707        let result = boxed.ready().await.unwrap().call(ex).await;
708        let ex = result.unwrap();
709        assert!(
710            !ex.properties.contains_key(PROPERTY_EXCEPTION_HANDLED),
711            "stale CamelExceptionHandled must be cleared on entry"
712        );
713    }
714
715    #[tokio::test]
716    async fn nested_do_try_inner_catch_does_not_leak_to_outer() {
717        let inner = {
718            let try_step = always_fail(CamelError::Io("inner".into()));
719            let mut d = DoTryService::new(vec![try_step]);
720            d.catch_clauses.push(CatchClause {
721                matcher: CatchMatcher::ByVariant(vec!["Io".into()]),
722                on_when: None,
723                steps: vec![passthrough()],
724                disposition: ExceptionDisposition::Handled,
725            });
726            BoxProcessor::new(d)
727        };
728        let mut outer = DoTryService::new(vec![inner]);
729        outer.catch_clauses.push(CatchClause {
730            matcher: CatchMatcher::ByVariant(vec!["Io".into()]),
731            on_when: None,
732            steps: vec![passthrough()],
733            disposition: ExceptionDisposition::Handled,
734        });
735
736        let mut boxed = BoxProcessor::new(outer);
737        let result = boxed.ready().await.unwrap().call(Exchange::default()).await;
738        assert!(
739            result.is_ok(),
740            "outer must see Ok because inner handled its own error"
741        );
742    }
743
744    #[tokio::test]
745    async fn catch_all_only_fires_when_no_specific_clause_matches() {
746        let try_step = always_fail(CamelError::Io("err".into()));
747        let processor_call = Arc::new(AtomicU32::new(0));
748        let catch_all_call = Arc::new(AtomicU32::new(0));
749
750        let mut svc = DoTryService::new(vec![try_step]);
751        // First clause (specific) targets ProcessorError — won't match Io error.
752        svc.catch_clauses.push(CatchClause {
753            matcher: CatchMatcher::ByVariant(vec!["ProcessorError".into()]),
754            on_when: None,
755            steps: vec![record_call(processor_call.clone())],
756            disposition: ExceptionDisposition::Handled,
757        });
758        // Second clause is the catch-all — should fire.
759        svc.catch_clauses.push(CatchClause {
760            matcher: CatchMatcher::ByVariant(vec!["*".into()]),
761            on_when: None,
762            steps: vec![record_call(catch_all_call.clone())],
763            disposition: ExceptionDisposition::Handled,
764        });
765
766        let mut boxed = BoxProcessor::new(svc);
767        let _ = boxed.ready().await.unwrap().call(Exchange::default()).await;
768        assert_eq!(
769            processor_call.load(Ordering::SeqCst),
770            0,
771            "specific ProcessorError clause must not fire on Io error"
772        );
773        assert_eq!(
774            catch_all_call.load(Ordering::SeqCst),
775            1,
776            "catch-all clause must fire when no specific clause matches"
777        );
778    }
779
780    #[tokio::test]
781    async fn catch_throws_with_finally_runs_finally_and_propagates_catch_err() {
782        let try_step = always_fail(CamelError::ProcessorError("orig".into()));
783        let catch_step = always_fail(CamelError::Io("catch-fail".into()));
784        let finally_flag = Arc::new(AtomicU32::new(0));
785        let finally_step = record_call(finally_flag.clone());
786
787        let mut svc = DoTryService::new(vec![try_step]);
788        svc.catch_clauses.push(CatchClause {
789            matcher: CatchMatcher::ByVariant(vec!["ProcessorError".into()]),
790            on_when: None,
791            steps: vec![catch_step],
792            disposition: ExceptionDisposition::Handled,
793        });
794        svc.finally_steps = vec![finally_step];
795
796        let mut boxed = BoxProcessor::new(svc);
797        let result = boxed.ready().await.unwrap().call(Exchange::default()).await;
798
799        assert!(result.is_err());
800        assert!(
801            matches!(result.unwrap_err(), CamelError::Io(_)),
802            "catch_err must propagate (not original ProcessorError)"
803        );
804        assert_eq!(
805            finally_flag.load(Ordering::SeqCst),
806            1,
807            "doFinally must run even when catch throws"
808        );
809    }
810
811    #[tokio::test]
812    async fn catch_throws_and_finally_throws_restores_catch_err() {
813        let try_step = always_fail(CamelError::ProcessorError("orig".into()));
814        let catch_step = always_fail(CamelError::Io("catch-fail".into()));
815        let finally_step = always_fail(CamelError::Config("fin-fail".into()));
816
817        let mut svc = DoTryService::new(vec![try_step]);
818        svc.catch_clauses.push(CatchClause {
819            matcher: CatchMatcher::ByVariant(vec!["ProcessorError".into()]),
820            on_when: None,
821            steps: vec![catch_step],
822            disposition: ExceptionDisposition::Handled,
823        });
824        svc.finally_steps = vec![finally_step];
825
826        let mut boxed = BoxProcessor::new(svc);
827        let result = boxed.ready().await.unwrap().call(Exchange::default()).await;
828
829        assert!(result.is_err());
830        assert!(
831            matches!(result.unwrap_err(), CamelError::Io(_)),
832            "catch_err (Io) must be restored over finally_err (Config) per Camel parity"
833        );
834    }
835
836    #[tokio::test]
837    async fn finally_on_when_false_with_previous_error_still_propagates_original() {
838        let try_step = always_fail(CamelError::ProcessorError("orig".into()));
839        let finally_flag = Arc::new(AtomicU32::new(0));
840        let finally_step = record_call(finally_flag.clone());
841
842        let mut svc = DoTryService::new(vec![try_step]);
843        // No catch clauses → original error stays.
844        svc.finally_steps = vec![finally_step];
845        svc.finally_on_when = Some(PredicateSource::Sync(FilterPredicate::new(|_ex| false)));
846
847        let mut boxed = BoxProcessor::new(svc);
848        let result = boxed.ready().await.unwrap().call(Exchange::default()).await;
849
850        assert!(result.is_err());
851        assert!(
852            matches!(result.unwrap_err(), CamelError::ProcessorError(_)),
853            "original error must propagate even when finally_on_when skips finally"
854        );
855        assert_eq!(
856            finally_flag.load(Ordering::SeqCst),
857            0,
858            "doFinally must NOT run when on_when returns false"
859        );
860    }
861
862    fn catch_fails_service() -> BoxProcessor {
863        let try_step = always_fail(CamelError::ProcessorError("orig-lost".into()));
864        let catch_step = always_fail(CamelError::Io("catch-fail".into()));
865        let mut svc = DoTryService::new(vec![try_step]);
866        svc.catch_clauses.push(CatchClause {
867            matcher: CatchMatcher::ByVariant(vec!["ProcessorError".into()]),
868            on_when: None,
869            steps: vec![catch_step],
870            disposition: ExceptionDisposition::Handled,
871        });
872        BoxProcessor::new(svc)
873    }
874
875    #[tokio::test]
876    async fn catch_throws_under_propagate_disposition_returns_catch_err() {
877        let try_step = always_fail(CamelError::ProcessorError("orig".into()));
878        let catch_step = always_fail(CamelError::Io("catch-fail".into()));
879        let mut svc = DoTryService::new(vec![try_step]);
880        svc.catch_clauses.push(CatchClause {
881            matcher: CatchMatcher::ByVariant(vec!["ProcessorError".into()]),
882            on_when: None,
883            steps: vec![catch_step],
884            disposition: ExceptionDisposition::Propagate,
885        });
886
887        let mut boxed = BoxProcessor::new(svc);
888        let result = boxed.ready().await.unwrap().call(Exchange::default()).await;
889        assert!(
890            matches!(result, Err(CamelError::Io(_))),
891            "Propagate-disposition catch failure must return the catch error, \
892             same envelope as Handled, got: {result:?}"
893        );
894    }
895
896    #[test]
897    fn catch_throws_logs_original_and_catch_error() {
898        let mut boxed = catch_fails_service();
899        let (result, captured, _span_records) = capture_debugs_with_span_records(|| {
900            tokio::runtime::Builder::new_current_thread()
901                .enable_all()
902                .build()
903                .expect("current-thread runtime")
904                .block_on(async { boxed.ready().await.unwrap().call(Exchange::default()).await })
905        });
906
907        assert!(
908            matches!(result, Err(CamelError::Io(_))),
909            "catch failure must return Err(Io), got: {result:?}"
910        );
911        // Structured field lookup on the envelope record, not message
912        // formatting.
913        let original = record_field(&captured, "do_try catch block failed", "original_error")
914            .unwrap_or_else(|| {
915                panic!("envelope record missing original_error field, captured: {captured:?}")
916            });
917        let catch = record_field(&captured, "do_try catch block failed", "catch_error")
918            .unwrap_or_else(|| {
919                panic!("envelope record missing catch_error field, captured: {captured:?}")
920            });
921        assert!(
922            original.contains("orig-lost"),
923            "original_error must carry the original error, got: {original}"
924        );
925        assert!(
926            catch.contains("catch-fail"),
927            "catch_error must carry the catch error, got: {catch}"
928        );
929    }
930
931    #[test]
932    fn catch_throws_marks_span_error_and_event() {
933        let mut boxed = catch_fails_service();
934        let (result, captured, span_records) = capture_debugs_with_span_records(|| {
935            // Declared `error` field so record_span_error's record lands.
936            let span = tracing::info_span!("dotry_test", error = tracing::field::Empty);
937            let _guard = span.enter();
938            tokio::runtime::Builder::new_current_thread()
939                .enable_all()
940                .build()
941                .expect("current-thread runtime")
942                .block_on(async { boxed.ready().await.unwrap().call(Exchange::default()).await })
943        });
944
945        assert!(
946            matches!(result, Err(CamelError::Io(_))),
947            "catch failure must return Err(Io), got: {result:?}"
948        );
949        // The span recorded the `error` field with the CATCH error.
950        // Structured field lookup: substring `error=` would also match
951        // `original_error=` suffixes.
952        assert!(
953            span_records.iter().any(|line| {
954                captured_field(line, "error").is_some_and(|v| v.contains("catch-fail"))
955            }),
956            "expected span error record carrying the catch error, span records: {span_records:?}"
957        );
958        // The WARN event still carries original_error while a span is active.
959        let original = record_field(
960            &captured,
961            "do_try catch block failed",
962            "original_error",
963        )
964        .unwrap_or_else(|| {
965            panic!(
966                "envelope record missing original_error field under active span, captured: {captured:?}"
967            )
968        });
969        assert!(
970            original.contains("orig-lost"),
971            "original_error must carry the original error under an active span, got: {original}"
972        );
973    }
974
975    #[test]
976    fn catch_and_finally_throw_logs_finally_error() {
977        let try_step = always_fail(CamelError::ProcessorError("orig".into()));
978        let catch_step = always_fail(CamelError::Io("catch-fail".into()));
979        let finally_step = always_fail(CamelError::Config("fin-fail".into()));
980        let mut svc = DoTryService::new(vec![try_step]);
981        svc.catch_clauses.push(CatchClause {
982            matcher: CatchMatcher::ByVariant(vec!["ProcessorError".into()]),
983            on_when: None,
984            steps: vec![catch_step],
985            disposition: ExceptionDisposition::Handled,
986        });
987        svc.finally_steps = vec![finally_step];
988
989        let mut boxed = BoxProcessor::new(svc);
990        let (result, captured, _span_records) = capture_debugs_with_span_records(|| {
991            tokio::runtime::Builder::new_current_thread()
992                .enable_all()
993                .build()
994                .expect("current-thread runtime")
995                .block_on(async { boxed.ready().await.unwrap().call(Exchange::default()).await })
996        });
997
998        assert!(
999            matches!(result, Err(CamelError::Io(_))),
1000            "catch error must be restored over finally error, got: {result:?}"
1001        );
1002        let catch = record_field(
1003            &captured,
1004            "doFinally threw after failed catch; restoring catch error",
1005            "catch_error",
1006        )
1007        .unwrap_or_else(|| {
1008            panic!("catch-failed restore record missing catch_error field, captured: {captured:?}")
1009        });
1010        let finally = record_field(
1011            &captured,
1012            "doFinally threw after failed catch; restoring catch error",
1013            "finally_error",
1014        )
1015        .unwrap_or_else(|| {
1016            panic!(
1017                "catch-failed restore record missing finally_error field, captured: {captured:?}"
1018            )
1019        });
1020        assert!(
1021            catch.contains("catch-fail"),
1022            "catch_error must carry the catch error, got: {catch}"
1023        );
1024        assert!(
1025            finally.contains("fin-fail"),
1026            "finally_error must carry the finally error, got: {finally}"
1027        );
1028    }
1029
1030    // ── Fallible predicate path (language-value-boundary task 1.4) ──
1031
1032    use camel_api::{ExpressionErrorClass, FilterPredicate, PredicateSource};
1033
1034    fn expression_failed() -> CamelError {
1035        CamelError::ExpressionFailed {
1036            language: "rhai".to_string(),
1037            route_id: "r1".to_string(),
1038            step_id: "step#0".to_string(),
1039            verb: "when".to_string(),
1040            class: ExpressionErrorClass::Runtime,
1041            position: None,
1042            conversion: None,
1043            cause: None,
1044        }
1045    }
1046
1047    fn async_err_predicate(err: CamelError) -> PredicateSource {
1048        PredicateSource::Async(Arc::new(move |_: &Exchange| {
1049            let err = err.clone();
1050            Box::pin(async move { Err(err) }) as camel_api::BoxBoolFuture
1051        }))
1052    }
1053
1054    #[tokio::test]
1055    async fn catch_when_predicate_error_chains_original() {
1056        let try_step = always_fail(CamelError::ProcessorError("boom".into()));
1057        let mut svc = DoTryService::new(vec![try_step]);
1058        svc.catch_clauses.push(CatchClause {
1059            matcher: CatchMatcher::Predicate(async_err_predicate(expression_failed())),
1060            on_when: None,
1061            steps: vec![passthrough()],
1062            disposition: ExceptionDisposition::Handled,
1063        });
1064
1065        let mut boxed = BoxProcessor::new(svc);
1066        let result = boxed.ready().await.unwrap().call(Exchange::default()).await;
1067        match result {
1068            Err(CamelError::ExpressionFailed {
1069                cause: Some(cause), ..
1070            }) => {
1071                assert!(
1072                    cause.to_string().contains("boom"),
1073                    "cause must carry the original error, got: {cause}"
1074                );
1075            }
1076            other => panic!("expected Err(ExpressionFailed) with chained cause, got {other:?}"),
1077        }
1078    }
1079
1080    #[tokio::test]
1081    async fn catch_on_when_predicate_error_chains_original() {
1082        let try_step = always_fail(CamelError::ProcessorError("boom".into()));
1083        let mut svc = DoTryService::new(vec![try_step]);
1084        svc.catch_clauses.push(CatchClause {
1085            matcher: CatchMatcher::ByVariant(vec!["*".into()]),
1086            on_when: Some(async_err_predicate(expression_failed())),
1087            steps: vec![passthrough()],
1088            disposition: ExceptionDisposition::Handled,
1089        });
1090
1091        let mut boxed = BoxProcessor::new(svc);
1092        let result = boxed.ready().await.unwrap().call(Exchange::default()).await;
1093        match result {
1094            Err(CamelError::ExpressionFailed {
1095                cause: Some(cause), ..
1096            }) => {
1097                assert!(
1098                    cause.to_string().contains("boom"),
1099                    "cause must carry the original error, got: {cause}"
1100                );
1101            }
1102            other => panic!("expected Err(ExpressionFailed) with chained cause, got {other:?}"),
1103        }
1104    }
1105
1106    #[tokio::test]
1107    async fn finally_on_when_predicate_error_fails() {
1108        let mut svc = DoTryService::new(vec![passthrough()]);
1109        svc.finally_steps = vec![passthrough()];
1110        svc.finally_on_when = Some(async_err_predicate(expression_failed()));
1111
1112        let mut boxed = BoxProcessor::new(svc);
1113        let result = boxed.ready().await.unwrap().call(Exchange::default()).await;
1114        match result {
1115            Err(CamelError::ExpressionFailed { .. }) => {}
1116            other => panic!("expected Err(ExpressionFailed), got {other:?}"),
1117        }
1118    }
1119}