1use 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#[derive(Clone)]
18pub enum CatchMatcher {
19 ByVariant(Vec<String>),
22 Predicate(PredicateSource),
24}
25
26impl CatchMatcher {
27 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
44pub(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#[derive(Clone)]
65pub struct CatchClause {
66 pub matcher: CatchMatcher,
68 pub on_when: Option<PredicateSource>,
70 pub steps: Vec<BoxProcessor>,
72 pub disposition: ExceptionDisposition,
75}
76
77#[derive(Clone)]
79pub struct DoTryService {
80 pub try_steps: Vec<BoxProcessor>,
82 pub catch_clauses: Vec<CatchClause>,
84 pub finally_steps: Vec<BoxProcessor>,
86 pub finally_on_when: Option<PredicateSource>,
88}
89
90impl DoTryService {
91 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 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
118async 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
139async 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#[allow(clippy::large_enum_variant)]
185enum FinallyOutcome {
186 Completed(Exchange),
188 NoPreviousFail(CamelError),
190 Restore {
192 previous: CamelError,
193 finally_err: CamelError,
194 },
195}
196
197async 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 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 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 Ok(FinallyOutcome::Restore { previous, .. }) => Err(previous),
262 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 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 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 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 _ => 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 Err(pred_err) => return Err(pred_err),
370 };
371 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 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 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 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#[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 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 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 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 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 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 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 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 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}