1use crate::do_try::CatchMatcher;
7use crate::do_try::chain_predicate_error;
8use crate::error_handler::record_span_error;
9use camel_api::error_handler::ExceptionDisposition;
10use camel_api::outcome_pipeline::OutcomePipeline;
11use camel_api::pipeline_outcome::PipelineOutcome;
12use camel_api::{CamelError, Exchange, PredicateSource};
13use std::future::Future;
14use std::pin::Pin;
15
16#[derive(Clone)]
20pub struct CatchClauseSegment {
21 pub matcher: CatchMatcher,
22 pub on_when: Option<PredicateSource>,
23 pub body: camel_api::OutcomeSegment,
24 pub disposition: ExceptionDisposition,
25}
26
27#[derive(Clone)]
29pub struct FinallyClauseSegment {
30 pub on_when: Option<PredicateSource>,
31 pub body: camel_api::OutcomeSegment,
32}
33
34pub struct DoTrySegment {
40 pub try_body: camel_api::OutcomeSegment,
41 pub catches: Vec<CatchClauseSegment>,
42 pub finally: Option<FinallyClauseSegment>,
43}
44
45impl Clone for DoTrySegment {
46 fn clone(&self) -> Self {
47 Self {
48 try_body: self.try_body.clone(),
49 catches: self.catches.clone(),
50 finally: self.finally.clone(),
51 }
52 }
53}
54
55enum FinallyOutcome {
60 Stopped(Box<Exchange>),
61 Failed(CamelError),
62 PredicateFailed(CamelError),
65}
66
67#[allow(clippy::result_large_err)]
74async fn run_finally_body(
75 finally: &mut Option<FinallyClauseSegment>,
76 ex: Exchange,
77) -> Result<Exchange, FinallyOutcome> {
78 let Some(f) = finally.as_mut() else {
79 return Ok(ex);
80 };
81 if let Some(on_when) = &f.on_when {
82 match on_when.matches(&ex).await {
83 Ok(true) => {}
84 Ok(false) => return Ok(ex),
85 Err(err) => return Err(FinallyOutcome::PredicateFailed(err)),
86 }
87 }
88 match f.body.run(ex).await {
89 PipelineOutcome::Completed(e) => Ok(e),
90 PipelineOutcome::Stopped(e) => Err(FinallyOutcome::Stopped(Box::new(e))),
91 PipelineOutcome::Failed(e) => Err(FinallyOutcome::Failed(e)),
92 }
93}
94
95async fn fail_after_predicate_error(
101 finally: &mut Option<FinallyClauseSegment>,
102 ex: Exchange,
103 chained: CamelError,
104) -> PipelineOutcome {
105 match run_finally_body(finally, ex).await {
106 Ok(_) => PipelineOutcome::Failed(chained),
107 Err(FinallyOutcome::Stopped(e)) => PipelineOutcome::Stopped(*e),
108 Err(FinallyOutcome::Failed(_finally_err)) => {
109 tracing::warn!(
110 error = %chained,
111 "doFinally threw after catch predicate failure; restoring chained error"
112 );
113 PipelineOutcome::Failed(chained)
114 }
115 Err(FinallyOutcome::PredicateFailed(pred_err)) => PipelineOutcome::Failed(pred_err),
116 }
117}
118
119impl OutcomePipeline for DoTrySegment {
120 fn clone_box(&self) -> Box<dyn OutcomePipeline> {
121 Box::new(self.clone())
122 }
123
124 fn run<'a>(
125 &'a mut self,
126 exchange: Exchange,
127 ) -> Pin<Box<dyn Future<Output = PipelineOutcome> + Send + 'a>> {
128 Box::pin(async move {
129 let exchange_for_unmatched = exchange.clone();
132
133 let try_outcome = self.try_body.run(exchange).await;
135 let returned_ex = match try_outcome {
136 PipelineOutcome::Stopped(ex) => return PipelineOutcome::Stopped(ex),
137 PipelineOutcome::Completed(ex) => ex,
138 PipelineOutcome::Failed(err) => {
139 let mut current_ex = exchange_for_unmatched;
145 current_ex.set_error(err.clone());
146 for catch in self.catches.iter_mut() {
147 let matched = match catch.matcher.matches(&err, ¤t_ex).await {
148 Ok(m) => m,
149 Err(pred_err) => {
150 let chained = chain_predicate_error(pred_err, err);
151 return fail_after_predicate_error(
152 &mut self.finally,
153 current_ex,
154 chained,
155 )
156 .await;
157 }
158 };
159 if !matched {
160 continue;
161 }
162 if let Some(on_when) = &catch.on_when {
163 match on_when.matches(¤t_ex).await {
164 Ok(true) => {}
165 Ok(false) => continue,
166 Err(pred_err) => {
167 let chained = chain_predicate_error(pred_err, err);
168 return fail_after_predicate_error(
169 &mut self.finally,
170 current_ex,
171 chained,
172 )
173 .await;
174 }
175 }
176 }
177 match catch.body.run(current_ex).await {
178 PipelineOutcome::Stopped(stopped_ex) => {
180 return PipelineOutcome::Stopped(stopped_ex);
181 }
182 PipelineOutcome::Completed(next) => {
184 match catch.disposition {
185 ExceptionDisposition::Handled => {
186 current_ex = next;
187 break;
188 }
189 ExceptionDisposition::Continued => {
190 current_ex = next;
191 break;
192 }
193 _ => {
196 match run_finally_body(&mut self.finally, next).await {
197 Ok(_) => {}
198 Err(FinallyOutcome::Stopped(e)) => {
199 return PipelineOutcome::Stopped(*e);
200 }
201 Err(FinallyOutcome::Failed(_finally_err)) => {
202 tracing::warn!(
203 error = %err,
204 "doFinally threw during Propagate; \
205 restoring original"
206 );
207 return PipelineOutcome::Failed(err);
208 }
209 Err(FinallyOutcome::PredicateFailed(pred_err)) => {
212 return PipelineOutcome::Failed(pred_err);
213 }
214 }
215 return PipelineOutcome::Failed(err);
216 }
217 }
218 }
219 PipelineOutcome::Failed(catch_err) => {
226 tracing::warn!(
233 original_error = %err,
234 catch_error = %catch_err,
235 "do_try catch block failed; catch error supersedes original"
236 );
237 record_span_error(&catch_err);
238 return PipelineOutcome::Failed(catch_err);
239 }
240 }
241 }
242 current_ex
243 }
244 };
245 match run_finally_body(&mut self.finally, returned_ex).await {
248 Ok(ex) => PipelineOutcome::Completed(ex),
249 Err(FinallyOutcome::Stopped(e)) => PipelineOutcome::Stopped(*e),
250 Err(FinallyOutcome::Failed(finally_err)) => {
251 tracing::warn!(
252 error = %finally_err,
253 "doFinally threw during/after catch; surfacing finally error"
254 );
255 PipelineOutcome::Failed(finally_err)
256 }
257 Err(FinallyOutcome::PredicateFailed(pred_err)) => PipelineOutcome::Failed(pred_err),
260 }
261 })
262 }
263}
264
265#[cfg(test)]
266mod tests {
267 use super::*;
268 use crate::do_try::CatchMatcher;
269 use crate::test_log_capture::{capture_debugs_with_span_records, captured_field, record_field};
270 use camel_api::pipeline_outcome::PipelineOutcome;
271 use std::sync::Arc;
272 use std::sync::atomic::{AtomicU32, Ordering};
273
274 struct CompleteSegment;
277
278 impl OutcomePipeline for CompleteSegment {
279 fn clone_box(&self) -> Box<dyn OutcomePipeline> {
280 Box::new(CompleteSegment)
281 }
282 fn run<'a>(
283 &'a mut self,
284 exchange: Exchange,
285 ) -> Pin<Box<dyn Future<Output = PipelineOutcome> + Send + 'a>> {
286 Box::pin(async move { PipelineOutcome::Completed(exchange) })
287 }
288 }
289
290 fn seg_complete() -> camel_api::OutcomeSegment {
291 camel_api::OutcomeSegment::new(Box::new(CompleteSegment))
292 }
293
294 struct FailSegment(CamelError);
295
296 impl OutcomePipeline for FailSegment {
297 fn clone_box(&self) -> Box<dyn OutcomePipeline> {
298 Box::new(FailSegment(self.0.clone()))
299 }
300 fn run<'a>(
301 &'a mut self,
302 _exchange: Exchange,
303 ) -> Pin<Box<dyn Future<Output = PipelineOutcome> + Send + 'a>> {
304 let e = self.0.clone();
305 Box::pin(async move { PipelineOutcome::Failed(e) })
306 }
307 }
308
309 fn seg_fail(err: CamelError) -> camel_api::OutcomeSegment {
310 camel_api::OutcomeSegment::new(Box::new(FailSegment(err)))
311 }
312
313 struct MutateThenStop {
314 mutator: Arc<dyn Fn(&mut Exchange) + Send + Sync>,
315 }
316
317 impl OutcomePipeline for MutateThenStop {
318 fn clone_box(&self) -> Box<dyn OutcomePipeline> {
319 Box::new(MutateThenStop {
320 mutator: Arc::clone(&self.mutator),
321 })
322 }
323 fn run<'a>(
324 &'a mut self,
325 mut exchange: Exchange,
326 ) -> Pin<Box<dyn Future<Output = PipelineOutcome> + Send + 'a>> {
327 let m = Arc::clone(&self.mutator);
328 Box::pin(async move {
329 m(&mut exchange);
330 PipelineOutcome::Stopped(exchange)
331 })
332 }
333 }
334
335 fn seg_stop_with(
336 mutator: impl Fn(&mut Exchange) + Send + Sync + 'static,
337 ) -> camel_api::OutcomeSegment {
338 camel_api::OutcomeSegment::new(Box::new(MutateThenStop {
339 mutator: Arc::new(mutator),
340 }))
341 }
342
343 struct RecordCall {
344 counter: Arc<AtomicU32>,
345 }
346
347 impl OutcomePipeline for RecordCall {
348 fn clone_box(&self) -> Box<dyn OutcomePipeline> {
349 Box::new(RecordCall {
350 counter: Arc::clone(&self.counter),
351 })
352 }
353 fn run<'a>(
354 &'a mut self,
355 exchange: Exchange,
356 ) -> Pin<Box<dyn Future<Output = PipelineOutcome> + Send + 'a>> {
357 let c = Arc::clone(&self.counter);
358 Box::pin(async move {
359 c.fetch_add(1, Ordering::SeqCst);
360 PipelineOutcome::Completed(exchange)
361 })
362 }
363 }
364
365 fn seg_record(counter: Arc<AtomicU32>) -> camel_api::OutcomeSegment {
366 camel_api::OutcomeSegment::new(Box::new(RecordCall { counter }))
367 }
368
369 #[tokio::test]
372 async fn stop_inside_try_skips_catch_and_finally() {
373 let catch_call = Arc::new(AtomicU32::new(0));
374 let finally_call = Arc::new(AtomicU32::new(0));
375
376 let mut seg = DoTrySegment {
377 try_body: seg_stop_with(|ex| {
378 ex.set_property("mutated", camel_api::Value::Bool(true));
379 }),
380 catches: vec![CatchClauseSegment {
381 matcher: CatchMatcher::ByVariant(vec!["*".into()]),
382 on_when: None,
383 body: seg_record(catch_call.clone()),
384 disposition: ExceptionDisposition::Handled,
385 }],
386 finally: Some(FinallyClauseSegment {
387 on_when: None,
388 body: seg_record(finally_call.clone()),
389 }),
390 };
391
392 let result = seg.run(Exchange::default()).await;
393 match result {
394 PipelineOutcome::Stopped(ex) => {
395 assert_eq!(
396 ex.properties.get("mutated"),
397 Some(&camel_api::Value::Bool(true)),
398 "try body mutation must be preserved in Stopped exchange"
399 );
400 }
401 other => panic!("expected Stopped, got {:?}", other),
402 }
403 assert_eq!(
404 catch_call.load(Ordering::SeqCst),
405 0,
406 "catch must NOT run when try stops"
407 );
408 assert_eq!(
409 finally_call.load(Ordering::SeqCst),
410 0,
411 "finally must NOT run when try stops"
412 );
413 }
414
415 #[tokio::test]
416 async fn stop_inside_catch_skips_finally() {
417 let finally_call = Arc::new(AtomicU32::new(0));
418
419 let mut seg = DoTrySegment {
420 try_body: seg_fail(CamelError::ProcessorError("boom".into())),
421 catches: vec![CatchClauseSegment {
422 matcher: CatchMatcher::ByVariant(vec!["ProcessorError".into()]),
423 on_when: None,
424 body: seg_stop_with(|ex| {
425 ex.set_property("catch_mutated", camel_api::Value::Bool(true));
426 }),
427 disposition: ExceptionDisposition::Handled,
428 }],
429 finally: Some(FinallyClauseSegment {
430 on_when: None,
431 body: seg_record(finally_call.clone()),
432 }),
433 };
434
435 let result = seg.run(Exchange::default()).await;
436 match result {
437 PipelineOutcome::Stopped(ex) => {
438 assert_eq!(
439 ex.properties.get("catch_mutated"),
440 Some(&camel_api::Value::Bool(true)),
441 "catch body mutation must be preserved in Stopped exchange"
442 );
443 }
444 other => panic!("expected Stopped, got {:?}", other),
445 }
446 assert_eq!(
447 finally_call.load(Ordering::SeqCst),
448 0,
449 "finally must NOT run when catch stops"
450 );
451 }
452
453 #[tokio::test]
454 async fn stop_inside_finally_stops_outer_route() {
455 let mut seg = DoTrySegment {
456 try_body: seg_complete(),
457 catches: vec![],
458 finally: Some(FinallyClauseSegment {
459 on_when: None,
460 body: seg_stop_with(|ex| {
461 ex.set_property("finally_mutated", camel_api::Value::Bool(true));
462 }),
463 }),
464 };
465
466 let result = seg.run(Exchange::default()).await;
467 match result {
468 PipelineOutcome::Stopped(ex) => {
469 assert_eq!(
470 ex.properties.get("finally_mutated"),
471 Some(&camel_api::Value::Bool(true)),
472 "finally body mutation must be preserved in Stopped exchange"
473 );
474 }
475 other => panic!("expected Stopped, got {:?}", other),
476 }
477 }
478
479 #[tokio::test]
480 async fn catch_on_when_false_falls_through_to_next_catch() {
481 let first_call = Arc::new(AtomicU32::new(0));
482 let second_call = Arc::new(AtomicU32::new(0));
483
484 let mut seg = DoTrySegment {
485 try_body: seg_fail(CamelError::Io("disk err".into())),
486 catches: vec![
487 CatchClauseSegment {
488 matcher: CatchMatcher::ByVariant(vec!["Io".into()]),
489 on_when: Some(PredicateSource::Sync(FilterPredicate::new(|_ex| false))),
490 body: seg_record(first_call.clone()),
491 disposition: ExceptionDisposition::Handled,
492 },
493 CatchClauseSegment {
494 matcher: CatchMatcher::ByVariant(vec!["*".into()]),
495 on_when: None,
496 body: seg_record(second_call.clone()),
497 disposition: ExceptionDisposition::Handled,
498 },
499 ],
500 finally: None,
501 };
502
503 let result = seg.run(Exchange::default()).await;
504 assert!(
505 matches!(result, PipelineOutcome::Completed(_)),
506 "expected Completed after second catch"
507 );
508 assert_eq!(
509 first_call.load(Ordering::SeqCst),
510 0,
511 "first catch must NOT fire (on_when=false)"
512 );
513 assert_eq!(
514 second_call.load(Ordering::SeqCst),
515 1,
516 "second catch must fire"
517 );
518 }
519
520 #[tokio::test]
521 async fn finally_on_when_false_skips_finally_entirely() {
522 let finally_call = Arc::new(AtomicU32::new(0));
523 let mut ex = Exchange::default();
524 ex.set_property("try_set", camel_api::Value::Bool(true));
525
526 let mut seg = DoTrySegment {
527 try_body: seg_complete(),
528 catches: vec![],
529 finally: Some(FinallyClauseSegment {
530 on_when: Some(PredicateSource::Sync(FilterPredicate::new(|_ex| false))),
531 body: seg_record(finally_call.clone()),
532 }),
533 };
534
535 let result = seg.run(ex).await;
536 match result {
537 PipelineOutcome::Completed(ex) => {
538 assert_eq!(
539 ex.properties.get("try_set"),
540 Some(&camel_api::Value::Bool(true)),
541 "exchange state from try must be preserved"
542 );
543 }
544 other => panic!("expected Completed, got {:?}", other),
545 }
546 assert_eq!(
547 finally_call.load(Ordering::SeqCst),
548 0,
549 "finally must NOT run when on_when=false"
550 );
551 }
552
553 fn catch_fails_segment() -> DoTrySegment {
556 DoTrySegment {
557 try_body: seg_fail(CamelError::ProcessorError("orig".into())),
558 catches: vec![CatchClauseSegment {
559 matcher: CatchMatcher::ByVariant(vec!["ProcessorError".into()]),
560 on_when: None,
561 body: seg_fail(CamelError::Io("catch-fail".into())),
562 disposition: ExceptionDisposition::Handled,
563 }],
564 finally: None,
565 }
566 }
567
568 #[test]
569 fn catch_body_failure_returns_catch_err_and_marks_span() {
570 let mut seg = catch_fails_segment();
571 let (result, _captured, span_records) = capture_debugs_with_span_records(|| {
572 let span = tracing::info_span!("dotry_seg_test", error = tracing::field::Empty);
574 let _guard = span.enter();
575 tokio::runtime::Builder::new_current_thread()
576 .enable_all()
577 .build()
578 .expect("current-thread runtime")
579 .block_on(async { seg.run(Exchange::default()).await })
580 });
581
582 assert!(
583 matches!(result, PipelineOutcome::Failed(CamelError::Io(_))),
584 "catch failure must surface the catch error as Failed(Io), got: {result:?}"
585 );
586 assert!(
590 span_records.iter().any(|line| {
591 captured_field(line, "error").is_some_and(|v| v.contains("catch-fail"))
592 }),
593 "expected span error record carrying the catch error, span records: {span_records:?}"
594 );
595 }
596
597 #[test]
598 fn catch_body_failure_emits_envelope_log() {
599 let mut seg = catch_fails_segment();
600 let (result, captured, _span_records) = capture_debugs_with_span_records(|| {
601 tokio::runtime::Builder::new_current_thread()
602 .enable_all()
603 .build()
604 .expect("current-thread runtime")
605 .block_on(async { seg.run(Exchange::default()).await })
606 });
607
608 assert!(
609 matches!(result, PipelineOutcome::Failed(CamelError::Io(_))),
610 "catch failure must surface the catch error as Failed(Io), got: {result:?}"
611 );
612 let original = record_field(&captured, "do_try catch block failed", "original_error")
615 .unwrap_or_else(|| {
616 panic!("envelope record missing original_error field, captured: {captured:?}")
617 });
618 let catch = record_field(&captured, "do_try catch block failed", "catch_error")
619 .unwrap_or_else(|| {
620 panic!("envelope record missing catch_error field, captured: {captured:?}")
621 });
622 assert!(
623 original.contains("orig"),
624 "original_error must carry the original error, got: {original}"
625 );
626 assert!(
627 catch.contains("catch-fail"),
628 "catch_error must carry the catch error, got: {catch}"
629 );
630 }
631
632 use camel_api::{ExpressionErrorClass, FilterPredicate, PredicateSource};
635
636 fn expression_failed() -> CamelError {
637 CamelError::ExpressionFailed {
638 language: "rhai".to_string(),
639 route_id: "r1".to_string(),
640 step_id: "step#0".to_string(),
641 verb: "on_when".to_string(),
642 class: ExpressionErrorClass::Runtime,
643 position: None,
644 conversion: None,
645 cause: None,
646 }
647 }
648
649 fn async_err_predicate(err: CamelError) -> PredicateSource {
650 PredicateSource::Async(Arc::new(move |_: &Exchange| {
651 let err = err.clone();
652 Box::pin(async move { Err(err) }) as camel_api::BoxBoolFuture
653 }))
654 }
655
656 #[tokio::test]
657 async fn finally_on_when_predicate_error_fails_segment() {
658 let finally_call = Arc::new(AtomicU32::new(0));
659 let mut seg = DoTrySegment {
660 try_body: seg_complete(),
661 catches: vec![],
662 finally: Some(FinallyClauseSegment {
663 on_when: Some(async_err_predicate(expression_failed())),
664 body: seg_record(finally_call.clone()),
665 }),
666 };
667
668 let result = seg.run(Exchange::default()).await;
669 match result {
670 PipelineOutcome::Failed(CamelError::ExpressionFailed { .. }) => {}
671 other => panic!("expected Failed(ExpressionFailed), got {other:?}"),
672 }
673 assert_eq!(
674 finally_call.load(Ordering::SeqCst),
675 0,
676 "finally body must not run when its on_when predicate errors"
677 );
678 }
679}