1use std::collections::BTreeMap;
7
8use camel_api::DelayConfig;
9use camel_api::aggregator::{
10 AggregationStrategy, AggregatorConfig, CompletionCondition, CompletionMode, CorrelationStrategy,
11};
12use camel_api::body::Body;
13use camel_api::body_converter::BodyType;
14use camel_api::circuit_breaker::CircuitBreakerConfig;
15use camel_api::dynamic_router::{DynamicRouterConfig, RouterExpression};
16use camel_api::error_handler::{ErrorHandlerConfig, RedeliveryPolicy};
17use camel_api::load_balancer::LoadBalancerConfig;
18use camel_api::loop_eip::{LoopConfig, LoopMode};
19use camel_api::multicast::{MulticastConfig, MulticastStrategy};
20use camel_api::recipient_list::{RecipientListConfig, RecipientListExpression};
21use camel_api::routing_slip::{RoutingSlipConfig, RoutingSlipExpression};
22use camel_api::splitter::SplitterConfig;
23use camel_api::throttler::{ThrottleStrategy, ThrottlerConfig};
24use camel_api::{
25 BoxProcessor, CamelError, CanonicalRouteSpec, EndpointUri, Exchange, FilterPredicate,
26 IdentityProcessor, LanguageExpressionDef, OpaqueProcessor, ProcessorFn, Value,
27 runtime::{
28 CanonicalAggregateSpec, CanonicalAggregateStrategySpec, CanonicalCircuitBreakerSpec,
29 CanonicalSplitAggregationSpec, CanonicalSplitExpressionSpec, CanonicalStepSpec,
30 CanonicalWhenSpec,
31 },
32};
33use camel_component_api::ConcurrencyModel;
34use camel_core::route::{BuilderStep, DeclarativeWhenStep, RouteDefinition, WhenStep};
35use camel_processor::{
36 ConvertBodyTo, DynamicSetHeader, LogLevel, MapBody, MarshalService, SetBody, SetHeader,
37 StreamCacheService, UnmarshalService, builtin_data_format,
38};
39
40pub mod do_try;
42pub use do_try::{DoCatchBuilder, DoFinallyBuilder, DoTryBuilder};
43
44pub trait StepAccumulator: Sized {
50 fn steps_mut(&mut self) -> &mut Vec<BuilderStep>;
51
52 fn to(mut self, endpoint: impl Into<String>) -> Self {
53 self.steps_mut().push(BuilderStep::To(endpoint.into()));
54 self
55 }
56
57 fn process<F, Fut>(mut self, f: F) -> Self
58 where
59 F: Fn(Exchange) -> Fut + Send + Sync + 'static,
60 Fut: std::future::Future<Output = Result<Exchange, CamelError>> + Send + 'static,
61 {
62 let svc = ProcessorFn::new(f);
63 self.steps_mut()
64 .push(BuilderStep::Processor(OpaqueProcessor(BoxProcessor::new(
65 svc,
66 ))));
67 self
68 }
69
70 fn process_fn(mut self, processor: BoxProcessor) -> Self {
71 self.steps_mut()
72 .push(BuilderStep::Processor(OpaqueProcessor(processor)));
73 self
74 }
75
76 fn set_header(mut self, key: impl Into<String>, value: impl Into<Value>) -> Self {
77 let svc = SetHeader::new(IdentityProcessor, key, value);
78 self.steps_mut()
79 .push(BuilderStep::Processor(OpaqueProcessor(BoxProcessor::new(
80 svc,
81 ))));
82 self
83 }
84
85 fn map_body<F>(mut self, mapper: F) -> Self
86 where
87 F: Fn(Body) -> Body + Clone + Send + Sync + 'static,
88 {
89 let svc = MapBody::new(IdentityProcessor, mapper);
90 self.steps_mut()
91 .push(BuilderStep::Processor(OpaqueProcessor(BoxProcessor::new(
92 svc,
93 ))));
94 self
95 }
96
97 fn set_body<B>(mut self, body: B) -> Self
98 where
99 B: Into<Body> + Clone + Send + Sync + 'static,
100 {
101 let body: Body = body.into();
102 let svc = SetBody::new(IdentityProcessor, move |_ex: &Exchange| body.clone());
103 self.steps_mut()
104 .push(BuilderStep::Processor(OpaqueProcessor(BoxProcessor::new(
105 svc,
106 ))));
107 self
108 }
109
110 fn transform<B>(self, body: B) -> Self
115 where
116 B: Into<Body> + Clone + Send + Sync + 'static,
117 {
118 self.set_body(body)
119 }
120
121 fn set_body_fn<F>(mut self, expr: F) -> Self
122 where
123 F: Fn(&Exchange) -> Body + Clone + Send + Sync + 'static,
124 {
125 let svc = SetBody::new(IdentityProcessor, expr);
126 self.steps_mut()
127 .push(BuilderStep::Processor(OpaqueProcessor(BoxProcessor::new(
128 svc,
129 ))));
130 self
131 }
132
133 fn set_header_fn<F>(mut self, key: impl Into<String>, expr: F) -> Self
134 where
135 F: Fn(&Exchange) -> Value + Clone + Send + Sync + 'static,
136 {
137 let svc = DynamicSetHeader::new(IdentityProcessor, key, expr);
138 self.steps_mut()
139 .push(BuilderStep::Processor(OpaqueProcessor(BoxProcessor::new(
140 svc,
141 ))));
142 self
143 }
144
145 fn aggregate(mut self, config: AggregatorConfig) -> Self {
146 self.steps_mut().push(BuilderStep::Aggregate { config });
147 self
148 }
149
150 fn stop(mut self) -> Self {
156 self.steps_mut().push(BuilderStep::Stop);
157 self
158 }
159
160 fn delay(mut self, duration: std::time::Duration) -> Self {
161 self.steps_mut().push(BuilderStep::Delay {
162 config: DelayConfig::from_duration(duration),
163 });
164 self
165 }
166
167 fn delay_with_header(
168 mut self,
169 duration: std::time::Duration,
170 header: impl Into<String>,
171 ) -> Self {
172 self.steps_mut().push(BuilderStep::Delay {
173 config: DelayConfig::from_duration_with_header(duration, header),
174 });
175 self
176 }
177
178 fn log(mut self, message: impl Into<String>, level: LogLevel) -> Self {
182 self.steps_mut().push(BuilderStep::Log {
183 level,
184 message: message.into(),
185 });
186 self
187 }
188
189 fn convert_body_to(mut self, target: BodyType) -> Self {
201 let svc = ConvertBodyTo::new(IdentityProcessor, target);
202 self.steps_mut()
203 .push(BuilderStep::Processor(OpaqueProcessor(BoxProcessor::new(
204 svc,
205 ))));
206 self
207 }
208
209 fn stream_cache(mut self, threshold: usize) -> Self {
210 let config = camel_api::stream_cache::StreamCacheConfig::new(threshold);
211 let svc = StreamCacheService::new(IdentityProcessor, config);
212 self.steps_mut()
213 .push(BuilderStep::Processor(OpaqueProcessor(BoxProcessor::new(
214 svc,
215 ))));
216 self
217 }
218
219 fn stream_cache_default(self) -> Self {
223 self.stream_cache(camel_api::stream_cache::DEFAULT_STREAM_CACHE_THRESHOLD)
224 }
225
226 fn marshal(mut self, format: impl Into<String>) -> Result<Self, CamelError> {
237 let name = format.into();
238 let df = builtin_data_format(&name)
239 .ok_or_else(|| CamelError::Config(format!("unknown data format: '{name}'")))?;
240 let svc = MarshalService::new(IdentityProcessor, df);
241 self.steps_mut()
242 .push(BuilderStep::Processor(OpaqueProcessor(BoxProcessor::new(
243 svc,
244 ))));
245 Ok(self)
246 }
247
248 fn unmarshal(mut self, format: impl Into<String>) -> Result<Self, CamelError> {
259 let name = format.into();
260 let df = builtin_data_format(&name)
261 .ok_or_else(|| CamelError::Config(format!("unknown data format: '{name}'")))?;
262 let svc = UnmarshalService::new(IdentityProcessor, df);
263 self.steps_mut()
264 .push(BuilderStep::Processor(OpaqueProcessor(BoxProcessor::new(
265 svc,
266 ))));
267 Ok(self)
268 }
269
270 fn validate(mut self, expression: impl Into<String>) -> Self {
280 let source = expression.into();
281 let expression = LanguageExpressionDef {
282 language: "simple".into(),
283 source,
284 };
285 self.steps_mut().push(BuilderStep::Validate {
286 predicate: expression,
287 });
288 self
289 }
290
291 fn script(mut self, language: impl Into<String>, script: impl Into<String>) -> Self {
302 self.steps_mut().push(BuilderStep::Script {
303 language: language.into(),
304 script: script.into(),
305 });
306 self
307 }
308
309 fn enrich(mut self, uri: impl Into<String>) -> Self {
315 self.steps_mut().push(BuilderStep::Enrich {
316 uri: uri.into(),
317 strategy: None,
318 timeout_ms: None,
319 });
320 self
321 }
322
323 fn poll_enrich(mut self, uri: impl Into<String>, timeout_ms: u64) -> Self {
329 self.steps_mut().push(BuilderStep::PollEnrich {
330 uri: uri.into(),
331 strategy: None,
332 timeout_ms: Some(timeout_ms),
333 });
334 self
335 }
336
337 fn bean(mut self, name: impl Into<String>, method: impl Into<String>) -> Self {
338 self.steps_mut().push(BuilderStep::Bean {
339 name: name.into(),
340 method: method.into(),
341 });
342 self
343 }
344}
345
346#[derive(Debug, Clone, PartialEq, Eq)]
348enum EndpointSlot {
349 From,
351 Step(usize),
353}
354
355#[derive(Clone)]
371pub struct RouteBuilder {
372 from_uri: String,
373 steps: Vec<BuilderStep>,
374 parameter_assignments: Vec<(EndpointSlot, BTreeMap<String, String>)>,
376 parameter_misuse: Option<String>,
378 error_handler: Option<ErrorHandlerConfig>,
379 error_handler_mode: ErrorHandlerMode,
380 circuit_breaker_config: Option<CircuitBreakerConfig>,
381 security_policy_config: Option<camel_api::security_policy::SecurityPolicyConfig>,
382 security_authenticator: Option<std::sync::Arc<dyn camel_auth::TokenAuthenticator>>,
383 provider_registry: Option<std::sync::Arc<camel_auth::ProviderRegistry>>,
384 concurrency: Option<ConcurrencyModel>,
385 route_id: Option<String>,
386 auto_startup: Option<bool>,
387 startup_order: Option<i32>,
388}
389
390#[derive(Default, Clone)]
391enum ErrorHandlerMode {
392 #[default]
393 None,
394 ExplicitConfig,
395 Shorthand {
396 dlc_uri: Option<String>,
397 specs: Vec<OnExceptionSpec>,
398 },
399 Mixed,
400}
401
402#[derive(Clone)]
403struct OnExceptionSpec {
404 matches: std::sync::Arc<dyn Fn(&CamelError) -> bool + Send + Sync>,
405 retry: Option<RedeliveryPolicy>,
406 handled_by: Option<String>,
407}
408
409impl RouteBuilder {
410 pub fn from(endpoint: &str) -> Self {
412 Self {
413 from_uri: endpoint.to_string(),
414 steps: Vec::new(),
415 parameter_assignments: Vec::new(),
416 parameter_misuse: None,
417 error_handler: None,
418 error_handler_mode: ErrorHandlerMode::None,
419 circuit_breaker_config: None,
420 security_policy_config: None,
421 security_authenticator: None,
422 provider_registry: None,
423 concurrency: None,
424 route_id: None,
425 auto_startup: None,
426 startup_order: None,
427 }
428 }
429
430 pub fn filter<F>(self, predicate: F) -> FilterBuilder
434 where
435 F: Fn(&Exchange) -> bool + Send + Sync + 'static,
436 {
437 FilterBuilder {
438 parent: self,
439 predicate: camel_api::FilterPredicate::new(predicate),
440 steps: vec![],
441 }
442 }
443
444 pub fn choice(self) -> ChoiceBuilder {
450 ChoiceBuilder {
451 parent: self,
452 whens: vec![],
453 _otherwise: None,
454 }
455 }
456
457 pub fn wire_tap(mut self, endpoint: &str) -> Self {
461 self.steps.push(BuilderStep::WireTap {
462 uri: endpoint.to_string(),
463 });
464 self
465 }
466
467 pub fn parameters(mut self, params: BTreeMap<String, String>) -> Self {
475 let slot = if self.steps.is_empty() {
476 EndpointSlot::From
477 } else {
478 EndpointSlot::Step(self.steps.len() - 1)
479 };
480
481 if self
482 .parameter_assignments
483 .iter()
484 .any(|(existing, _)| *existing == slot)
485 {
486 self.parameter_misuse =
487 Some("multiple .parameters() calls attached to the same endpoint".to_string());
488 } else if !endpoint_slot_bears_uri(&slot, &self.steps) {
489 self.parameter_misuse =
490 Some(".parameters() called with no pending endpoint step".to_string());
491 }
492
493 self.parameter_assignments.push((slot, params));
494 self
495 }
496
497 pub fn error_handler(mut self, config: ErrorHandlerConfig) -> Self {
499 self.error_handler_mode = match self.error_handler_mode {
500 ErrorHandlerMode::None | ErrorHandlerMode::ExplicitConfig => {
501 ErrorHandlerMode::ExplicitConfig
502 }
503 ErrorHandlerMode::Shorthand { .. } | ErrorHandlerMode::Mixed => ErrorHandlerMode::Mixed,
504 };
505 self.error_handler = Some(config);
506 self
507 }
508
509 pub fn dead_letter_channel(mut self, uri: impl Into<String>) -> Self {
511 let uri = uri.into();
512 self.error_handler_mode = match self.error_handler_mode {
513 ErrorHandlerMode::None => ErrorHandlerMode::Shorthand {
514 dlc_uri: Some(uri),
515 specs: Vec::new(),
516 },
517 ErrorHandlerMode::Shorthand { specs, .. } => ErrorHandlerMode::Shorthand {
518 dlc_uri: Some(uri),
519 specs,
520 },
521 ErrorHandlerMode::ExplicitConfig | ErrorHandlerMode::Mixed => ErrorHandlerMode::Mixed,
522 };
523 self
524 }
525
526 pub fn on_exception<F>(mut self, matches: F) -> OnExceptionBuilder
528 where
529 F: Fn(&CamelError) -> bool + Send + Sync + 'static,
530 {
531 self.error_handler_mode = match self.error_handler_mode {
532 ErrorHandlerMode::None => ErrorHandlerMode::Shorthand {
533 dlc_uri: None,
534 specs: Vec::new(),
535 },
536 ErrorHandlerMode::ExplicitConfig | ErrorHandlerMode::Mixed => ErrorHandlerMode::Mixed,
537 shorthand @ ErrorHandlerMode::Shorthand { .. } => shorthand,
538 };
539
540 OnExceptionBuilder {
541 parent: self,
542 policy: OnExceptionSpec {
543 matches: std::sync::Arc::new(matches),
544 retry: None,
545 handled_by: None,
546 },
547 }
548 }
549
550 pub fn circuit_breaker(mut self, config: CircuitBreakerConfig) -> Self {
552 self.circuit_breaker_config = Some(config);
553 self
554 }
555
556 pub fn security_policy(
557 mut self,
558 config: camel_api::security_policy::SecurityPolicyConfig,
559 ) -> Self {
560 self.security_policy_config = Some(config);
561 self
562 }
563
564 pub fn security_authenticator(
565 mut self,
566 auth: std::sync::Arc<dyn camel_auth::TokenAuthenticator>,
567 ) -> Self {
568 self.security_authenticator = Some(auth);
569 self
570 }
571
572 pub fn provider_registry(
576 mut self,
577 registry: std::sync::Arc<camel_auth::ProviderRegistry>,
578 ) -> Self {
579 self.provider_registry = Some(registry);
580 self
581 }
582
583 pub fn concurrent(mut self, max: usize) -> Self {
597 let max = if max == 0 { None } else { Some(max) };
598 self.concurrency = Some(ConcurrencyModel::Concurrent { max });
599 self
600 }
601
602 pub fn sequential(mut self) -> Self {
607 self.concurrency = Some(ConcurrencyModel::Sequential);
608 self
609 }
610
611 pub fn route_id(mut self, id: impl Into<String>) -> Self {
615 self.route_id = Some(id.into());
616 self
617 }
618
619 pub fn auto_startup(mut self, auto: bool) -> Self {
623 self.auto_startup = Some(auto);
624 self
625 }
626
627 pub fn startup_order(mut self, order: i32) -> Self {
631 self.startup_order = Some(order);
632 self
633 }
634
635 pub fn split(self, config: SplitterConfig) -> SplitBuilder {
641 SplitBuilder {
642 parent: self,
643 config,
644 steps: Vec::new(),
645 }
646 }
647
648 pub fn multicast(self) -> MulticastBuilder {
654 MulticastBuilder {
655 parent: self,
656 steps: Vec::new(),
657 config: MulticastConfig::new(),
658 }
659 }
660
661 pub fn throttle(self, max_requests: usize, period: std::time::Duration) -> ThrottleBuilder {
668 ThrottleBuilder {
669 parent: self,
670 config: ThrottlerConfig::new(max_requests, period),
671 steps: Vec::new(),
672 }
673 }
674
675 pub fn loop_count(self, count: usize) -> LoopBuilder {
677 LoopBuilder {
678 parent: self,
679 config: LoopConfig::new(LoopMode::Count(count)),
680 steps: vec![],
681 }
682 }
683
684 pub fn loop_while<F>(self, predicate: F) -> LoopBuilder
686 where
687 F: Fn(&Exchange) -> bool + Send + Sync + 'static,
688 {
689 LoopBuilder {
690 parent: self,
691 config: LoopConfig::new(LoopMode::While(camel_api::FilterPredicate::new(predicate))),
692 steps: vec![],
693 }
694 }
695
696 pub fn load_balance(self) -> LoadBalancerBuilder {
702 LoadBalancerBuilder {
703 parent: self,
704 config: LoadBalancerConfig::round_robin(),
705 steps: Vec::new(),
706 }
707 }
708
709 pub fn dynamic_router(self, expression: RouterExpression) -> Self {
725 self.dynamic_router_with_config(DynamicRouterConfig::new(expression))
726 }
727
728 pub fn dynamic_router_with_config(mut self, config: DynamicRouterConfig) -> Self {
732 self.steps.push(BuilderStep::DynamicRouter { config });
733 self
734 }
735
736 pub fn routing_slip(self, expression: RoutingSlipExpression) -> Self {
737 self.routing_slip_with_config(RoutingSlipConfig::new(expression))
738 }
739
740 pub fn routing_slip_with_config(mut self, config: RoutingSlipConfig) -> Self {
741 self.steps.push(BuilderStep::RoutingSlip { config });
742 self
743 }
744
745 pub fn recipient_list(self, expression: RecipientListExpression) -> Self {
746 self.recipient_list_with_config(RecipientListConfig::new(expression))
747 }
748
749 pub fn recipient_list_with_config(mut self, config: RecipientListConfig) -> Self {
750 self.steps.push(BuilderStep::RecipientList { config });
751 self
752 }
753
754 pub fn build(mut self) -> Result<RouteDefinition, CamelError> {
758 validate_uri(&self.from_uri)?;
759 let route_id = self
760 .route_id
761 .filter(|s| !s.trim().is_empty())
762 .ok_or_else(|| {
763 CamelError::RouteError(
764 "route must have a non-empty 'route_id' — call .route_id(\"name\") on the builder"
765 .to_string(),
766 )
767 })?;
768 let resolved_error_handler = match self.error_handler_mode {
769 ErrorHandlerMode::None => self.error_handler,
770 ErrorHandlerMode::ExplicitConfig => self.error_handler,
771 ErrorHandlerMode::Mixed => {
772 return Err(CamelError::RouteError(
773 "mixed error handler modes: cannot combine .error_handler(config) with shorthand methods".into(),
774 ));
775 }
776 ErrorHandlerMode::Shorthand { dlc_uri, specs } => {
777 let mut config = if let Some(uri) = dlc_uri {
778 ErrorHandlerConfig::dead_letter_channel(uri)
779 } else {
780 ErrorHandlerConfig::log_only()
781 };
782
783 for spec in specs {
784 let matcher = spec.matches.clone();
785 let mut builder = config.on_exception(move |e| matcher(e));
786
787 if let Some(retry) = spec.retry {
788 builder = builder.retry(retry.max_attempts).with_backoff(
789 retry.initial_delay,
790 retry.multiplier,
791 retry.max_delay,
792 );
793 if retry.jitter_factor > 0.0 {
794 builder = builder.with_jitter(retry.jitter_factor);
795 }
796 }
797
798 if let Some(uri) = spec.handled_by {
799 builder = builder.handled_by(uri);
800 }
801
802 config = builder.build();
803 }
804
805 Some(config)
806 }
807 };
808
809 apply_parameter_assignments(
812 &mut self.from_uri,
813 &mut self.steps,
814 std::mem::take(&mut self.parameter_assignments),
815 std::mem::take(&mut self.parameter_misuse),
816 )?;
817
818 let definition = RouteDefinition::new(self.from_uri, self.steps);
819 let definition = if let Some(eh) = resolved_error_handler {
820 definition.with_error_handler(eh)
821 } else {
822 definition
823 };
824 let definition = if let Some(cb) = self.circuit_breaker_config {
825 definition.with_circuit_breaker(cb)
826 } else {
827 definition
828 };
829 let definition = if let Some(sp) = self.security_policy_config {
830 definition.with_security_policy(sp)
831 } else {
832 definition
833 };
834 let definition = if let Some(auth) = self.security_authenticator {
835 definition.with_security_authenticator(auth)
836 } else {
837 definition
838 };
839 let definition = if let Some(registry) = self.provider_registry {
840 definition.with_provider_registry(registry)
841 } else {
842 definition
843 };
844 let definition = if let Some(concurrency) = self.concurrency {
845 definition.with_concurrency(concurrency)
846 } else {
847 definition
848 };
849 let definition = definition.with_route_id(route_id);
850 let definition = if let Some(auto) = self.auto_startup {
851 definition.with_auto_startup(auto)
852 } else {
853 definition
854 };
855 let definition = if let Some(order) = self.startup_order {
856 definition.with_startup_order(order)
857 } else {
858 definition
859 };
860 Ok(definition)
861 }
862
863 pub fn build_canonical(mut self) -> Result<CanonicalRouteSpec, CamelError> {
865 validate_uri(&self.from_uri)?;
866 let route_id = self
867 .route_id
868 .filter(|s| !s.trim().is_empty())
869 .ok_or_else(|| {
870 CamelError::RouteError(
871 "route must have a non-empty 'route_id' — call .route_id(\"name\") on the builder"
872 .to_string(),
873 )
874 })?;
875
876 apply_parameter_assignments(
878 &mut self.from_uri,
879 &mut self.steps,
880 std::mem::take(&mut self.parameter_assignments),
881 std::mem::take(&mut self.parameter_misuse),
882 )?;
883
884 let steps = canonicalize_steps(self.steps)?;
885 let circuit_breaker = self
886 .circuit_breaker_config
887 .map(canonicalize_circuit_breaker)
888 .transpose()?;
889
890 if self.security_policy_config.is_some() {
891 return Err(CamelError::RouteError(
892 "routes with security_policy cannot use the canonical/hot-reload path (not yet supported)"
893 .into(),
894 ));
895 }
896
897 let spec = CanonicalRouteSpec {
898 route_id,
899 from: self.from_uri,
900 steps,
901 circuit_breaker,
902 auto_startup: None,
903 startup_order: None,
904 concurrency: None,
905 version: camel_api::CANONICAL_CONTRACT_VERSION,
906 };
907 spec.validate_contract()?;
908 Ok(spec)
909 }
910}
911
912pub struct OnExceptionBuilder {
913 parent: RouteBuilder,
914 policy: OnExceptionSpec,
915}
916
917impl OnExceptionBuilder {
918 pub fn retry(mut self, max_attempts: u32) -> Self {
919 self.policy.retry = Some(RedeliveryPolicy::new(max_attempts));
920 self
921 }
922
923 pub fn with_backoff(
924 mut self,
925 initial: std::time::Duration,
926 multiplier: f64,
927 max: std::time::Duration,
928 ) -> Self {
929 if let Some(ref mut retry) = self.policy.retry {
930 retry.initial_delay = initial;
931 retry.multiplier = multiplier;
932 retry.max_delay = max;
933 } else {
934 tracing::warn!("backoff/jitter configuration has no effect when retry_count is 0");
935 }
936 self
937 }
938
939 pub fn with_jitter(mut self, jitter_factor: f64) -> Self {
940 if let Some(ref mut retry) = self.policy.retry {
941 retry.jitter_factor = jitter_factor.clamp(0.0, 1.0);
942 } else {
943 tracing::warn!("backoff/jitter configuration has no effect when retry_count is 0");
944 }
945 self
946 }
947
948 pub fn handled_by(mut self, uri: impl Into<String>) -> Self {
949 self.policy.handled_by = Some(uri.into());
950 self
951 }
952
953 pub fn end_on_exception(mut self) -> RouteBuilder {
954 if let ErrorHandlerMode::Shorthand { ref mut specs, .. } = self.parent.error_handler_mode {
955 specs.push(self.policy);
956 }
957 self.parent
958 }
959}
960
961fn endpoint_slot_bears_uri(slot: &EndpointSlot, steps: &[BuilderStep]) -> bool {
964 match slot {
965 EndpointSlot::From => true,
966 EndpointSlot::Step(i) => steps.get(*i).is_some_and(|step| {
967 matches!(
968 step,
969 BuilderStep::To(_)
970 | BuilderStep::WireTap { .. }
971 | BuilderStep::Enrich { .. }
972 | BuilderStep::PollEnrich { .. }
973 )
974 }),
975 }
976}
977
978fn apply_parameter_assignments(
987 from_uri: &mut String,
988 steps: &mut [BuilderStep],
989 assignments: Vec<(EndpointSlot, BTreeMap<String, String>)>,
990 misuse: Option<String>,
991) -> Result<(), CamelError> {
992 if let Some(reason) = misuse {
993 return Err(CamelError::RouteError(reason));
994 }
995
996 for (slot, params) in assignments {
997 if params.is_empty() {
1001 continue;
1002 }
1003
1004 let uri = match slot {
1005 EndpointSlot::From => &mut *from_uri,
1006 EndpointSlot::Step(i) => match steps.get_mut(i) {
1007 Some(BuilderStep::To(uri)) => uri,
1008 Some(BuilderStep::WireTap { uri, .. }) => uri,
1009 Some(BuilderStep::Enrich { uri, .. }) => uri,
1010 Some(BuilderStep::PollEnrich { uri, .. }) => uri,
1011 Some(other) => {
1012 return Err(CamelError::RouteError(format!(
1013 ".parameters() attached to a non-endpoint step (`{}`)",
1014 canonical_step_name(other)
1015 )));
1016 }
1017 None => {
1018 return Err(CamelError::RouteError(
1019 ".parameters() attached to an out-of-range step index".to_string(),
1020 ));
1021 }
1022 },
1023 };
1024 let merged = EndpointUri::try_from_uri_and_params(uri, params)?.to_canonical_string();
1025 *uri = merged;
1026 }
1027
1028 Ok(())
1029}
1030
1031fn validate_uri(uri: &str) -> Result<(), CamelError> {
1033 let trimmed = uri.trim();
1034 if trimmed.is_empty() {
1035 return Err(CamelError::RouteError(
1036 "route must have a 'from' URI".to_string(),
1037 ));
1038 }
1039 if !trimmed.contains(':') {
1040 return Err(CamelError::RouteError(
1041 "URI must have a scheme (e.g. 'timer:tick')".to_string(),
1042 ));
1043 }
1044 let scheme = trimmed.split(':').next().unwrap_or("");
1045 if scheme.trim().is_empty() {
1046 return Err(CamelError::RouteError(
1047 "URI scheme must not be empty".to_string(),
1048 ));
1049 }
1050 Ok(())
1051}
1052
1053fn canonicalize_steps(steps: Vec<BuilderStep>) -> Result<Vec<CanonicalStepSpec>, CamelError> {
1054 let mut canonical = Vec::with_capacity(steps.len());
1055 for step in steps {
1056 canonical.push(canonicalize_step(step)?);
1057 }
1058 Ok(canonical)
1059}
1060
1061fn canonicalize_step(step: BuilderStep) -> Result<CanonicalStepSpec, CamelError> {
1062 match step {
1063 BuilderStep::To(uri) => Ok(CanonicalStepSpec::To { uri }),
1064 BuilderStep::Log { message, .. } => Ok(CanonicalStepSpec::Log { message }),
1065 BuilderStep::Stop => Ok(CanonicalStepSpec::Stop),
1066 BuilderStep::WireTap { uri } => Ok(CanonicalStepSpec::WireTap { uri }),
1067 BuilderStep::Delay { config } => Ok(CanonicalStepSpec::Delay {
1068 delay_ms: config.delay_ms,
1069 dynamic_header: config.dynamic_header,
1070 }),
1071 BuilderStep::DeclarativeScript { expression } => {
1072 Ok(CanonicalStepSpec::Script { expression })
1073 }
1074 BuilderStep::DeclarativeFilter { predicate, steps } => Ok(CanonicalStepSpec::Filter {
1075 predicate,
1076 steps: canonicalize_steps(steps)?,
1077 }),
1078 BuilderStep::DeclarativeChoice { whens, otherwise } => {
1079 let mut canonical_whens = Vec::with_capacity(whens.len());
1080 for DeclarativeWhenStep { predicate, steps } in whens {
1081 canonical_whens.push(CanonicalWhenSpec {
1082 predicate,
1083 steps: canonicalize_steps(steps)?,
1084 });
1085 }
1086 let otherwise = match otherwise {
1087 Some(steps) => Some(canonicalize_steps(steps)?),
1088 None => None,
1089 };
1090 Ok(CanonicalStepSpec::Choice {
1091 whens: canonical_whens,
1092 otherwise,
1093 })
1094 }
1095 BuilderStep::DeclarativeSplit {
1096 expression,
1097 aggregation,
1098 parallel,
1099 parallel_limit,
1100 trace_item_threshold: _,
1103 stop_on_exception,
1104 steps,
1105 } => Ok(CanonicalStepSpec::Split {
1106 expression: CanonicalSplitExpressionSpec::Language(expression),
1107 aggregation: canonicalize_split_aggregation(aggregation)?,
1108 parallel,
1109 parallel_limit,
1110 trace_item_threshold: None,
1111 stop_on_exception,
1112 steps: canonicalize_steps(steps)?,
1113 }),
1114 BuilderStep::Aggregate { config } => Ok(CanonicalStepSpec::Aggregate(
1115 canonicalize_aggregate(config)?,
1116 )),
1117 other => {
1118 let step_name = canonical_step_name(&other);
1119 let detail = camel_api::canonical_contract_rejection_reason(step_name)
1120 .unwrap_or("not included in canonical v2");
1121 Err(CamelError::RouteError(format!(
1122 "canonical v2 does not support step `{step_name}`: {detail}"
1123 )))
1124 }
1125 }
1126}
1127
1128fn canonicalize_split_aggregation(
1129 strategy: camel_api::splitter::AggregationStrategy,
1130) -> Result<CanonicalSplitAggregationSpec, CamelError> {
1131 match strategy {
1132 camel_api::splitter::AggregationStrategy::LastWins => {
1133 Ok(CanonicalSplitAggregationSpec::LastWins)
1134 }
1135 camel_api::splitter::AggregationStrategy::CollectAll => {
1136 Ok(CanonicalSplitAggregationSpec::CollectAll)
1137 }
1138 camel_api::splitter::AggregationStrategy::Custom(_) => Err(CamelError::RouteError(
1139 "canonical v2 does not support custom split aggregation".to_string(),
1140 )),
1141 camel_api::splitter::AggregationStrategy::Original => {
1142 Ok(CanonicalSplitAggregationSpec::Original)
1143 }
1144 _ => Err(CamelError::RouteError(
1145 "canonical v2 does not support this split aggregation strategy".to_string(),
1146 )),
1147 }
1148}
1149
1150fn extract_completion_fields(
1151 mode: &CompletionMode,
1152) -> Result<(Option<usize>, Option<u64>), CamelError> {
1153 match mode {
1154 CompletionMode::Single(cond) => match cond {
1155 CompletionCondition::Size(n) => Ok((Some(*n), None)),
1156 CompletionCondition::Timeout(d) => Ok((None, Some(d.as_millis() as u64))),
1157 CompletionCondition::Predicate(_) | CompletionCondition::PredicateExpr { .. } => {
1158 Err(CamelError::RouteError(
1159 "aggregate PredicateExpr/Predicate completion cannot reverse-map to canonical \
1160 (forward-only in rc-zit); build the canonical spec directly"
1161 .to_string(),
1162 ))
1163 }
1164 _ => Err(CamelError::RouteError(
1165 "unsupported completion condition".to_string(),
1166 )),
1167 },
1168 CompletionMode::Any(conds) => {
1169 let mut size = None;
1170 let mut timeout_ms = None;
1171 for cond in conds {
1172 match cond {
1173 CompletionCondition::Size(n) => size = Some(*n),
1174 CompletionCondition::Timeout(d) => timeout_ms = Some(d.as_millis() as u64),
1175 CompletionCondition::Predicate(_)
1176 | CompletionCondition::PredicateExpr { .. } => {
1177 return Err(CamelError::RouteError(
1178 "aggregate PredicateExpr/Predicate completion cannot reverse-map to \
1179 canonical (forward-only in rc-zit); build the canonical spec directly"
1180 .to_string(),
1181 ));
1182 }
1183 _ => {
1184 return Err(CamelError::RouteError(
1185 "unsupported completion condition".to_string(),
1186 ));
1187 }
1188 }
1189 }
1190 Ok((size, timeout_ms))
1191 }
1192 _ => Err(CamelError::RouteError(
1193 "unsupported completion mode".to_string(),
1194 )),
1195 }
1196}
1197
1198fn canonicalize_aggregate(config: AggregatorConfig) -> Result<CanonicalAggregateSpec, CamelError> {
1199 let (completion_size, completion_timeout_ms) = extract_completion_fields(&config.completion)?;
1200
1201 let header = match &config.correlation {
1202 CorrelationStrategy::HeaderName(h) => h.clone(),
1203 CorrelationStrategy::Expression { expr, .. } => expr.clone(),
1204 CorrelationStrategy::Fn(_) => {
1205 return Err(CamelError::RouteError(
1206 "canonical v2 does not support Fn correlation strategy".to_string(),
1207 ));
1208 }
1209 _ => {
1210 return Err(CamelError::RouteError(
1211 "canonical v2 does not support this correlation strategy".to_string(),
1212 ));
1213 }
1214 };
1215
1216 let correlation_key = match &config.correlation {
1217 CorrelationStrategy::HeaderName(_) => None,
1218 CorrelationStrategy::Expression { expr, .. } => Some(expr.clone()),
1219 _ => unreachable!(),
1223 };
1224
1225 let strategy = match config.strategy {
1226 AggregationStrategy::CollectAll => CanonicalAggregateStrategySpec::CollectAll,
1227 AggregationStrategy::Custom(_) => {
1228 return Err(CamelError::RouteError(
1229 "canonical v2 does not support custom aggregate strategy".to_string(),
1230 ));
1231 }
1232 _ => {
1233 return Err(CamelError::RouteError(
1234 "canonical v2 does not support this aggregate strategy".to_string(),
1235 ));
1236 }
1237 };
1238 let bucket_ttl_ms = config
1239 .bucket_ttl
1240 .map(|ttl| u64::try_from(ttl.as_millis()).unwrap_or(u64::MAX));
1241
1242 Ok(CanonicalAggregateSpec {
1243 header,
1244 completion_size,
1245 completion_timeout_ms,
1246 correlation_key,
1247 force_completion_on_stop: if config.force_completion_on_stop {
1248 Some(true)
1249 } else {
1250 None
1251 },
1252 discard_on_timeout: if config.discard_on_timeout {
1253 Some(true)
1254 } else {
1255 None
1256 },
1257 strategy,
1258 max_buckets: config.max_buckets,
1259 max_bucket_size: config.max_bucket_size,
1260 bucket_ttl_ms,
1261 completion_predicate: None,
1262 })
1263}
1264
1265fn canonicalize_circuit_breaker(
1266 config: CircuitBreakerConfig,
1267) -> Result<CanonicalCircuitBreakerSpec, CamelError> {
1268 if config.fallback.is_some() {
1269 return Err(CamelError::RouteError(
1270 "canonical v2 does not support circuit breaker `fallback` (opaque BoxProcessor \
1271 cannot reverse-map to canonical steps); build the canonical spec directly"
1272 .to_string(),
1273 ));
1274 }
1275 Ok(CanonicalCircuitBreakerSpec {
1276 failure_threshold: config.failure_threshold,
1277 open_duration_ms: u64::try_from(config.open_duration.as_millis()).unwrap_or(u64::MAX),
1278 fallback: Vec::new(),
1279 })
1280}
1281
1282fn canonical_step_name(step: &BuilderStep) -> &'static str {
1283 match step {
1284 BuilderStep::Processor(_) => "processor",
1285 BuilderStep::To(_) => "to",
1286 BuilderStep::Stop => "stop",
1287 BuilderStep::Log { .. } => "log",
1288 BuilderStep::DeclarativeSetHeader { .. } => "set_header",
1289 BuilderStep::DeclarativeSetHeaderIfAbsent { .. } => "set_header_if_absent",
1290 BuilderStep::DeclarativeRemoveHeader { .. } => "remove_header",
1291 BuilderStep::DeclarativeSetBody { .. } => "set_body",
1292 BuilderStep::DeclarativeFilter { .. } => "filter",
1293 BuilderStep::DeclarativeChoice { .. } => "choice",
1294 BuilderStep::DeclarativeScript { .. } => "script",
1295 BuilderStep::DeclarativeFunction { .. } => "function",
1296 BuilderStep::DeclarativeSplit { .. } => "split",
1297 BuilderStep::Split { .. } => "split",
1298 BuilderStep::Loop { .. } | BuilderStep::DeclarativeLoop { .. } => "loop",
1299 BuilderStep::Aggregate { .. } => "aggregate",
1300 BuilderStep::Filter { .. } => "filter",
1301 BuilderStep::Choice { .. } => "choice",
1302 BuilderStep::WireTap { .. } => "wire_tap",
1303 BuilderStep::Delay { .. } => "delay",
1304 BuilderStep::Multicast { .. } => "multicast",
1305 BuilderStep::DeclarativeLog { .. } => "log",
1306 BuilderStep::Bean { .. } => "bean",
1307 BuilderStep::Script { .. } => "script",
1308 BuilderStep::Throttle { .. } => "throttle",
1309 BuilderStep::LoadBalance { .. } => "load_balancer",
1310 BuilderStep::DynamicRouter { .. } => "dynamic_router",
1311 BuilderStep::RoutingSlip { .. } => "routing_slip",
1312 BuilderStep::DeclarativeDynamicRouter { .. } => "declarative_dynamic_router",
1313 BuilderStep::DeclarativeRoutingSlip { .. } => "declarative_routing_slip",
1314 BuilderStep::RecipientList { .. } => "recipient_list",
1315 BuilderStep::DeclarativeRecipientList { .. } => "declarative_recipient_list",
1316 BuilderStep::DeclarativeSetProperty { .. } => "set_property",
1317 BuilderStep::DeclarativeStreamSplit { .. } => "stream_split",
1318 BuilderStep::Enrich { .. } => "enrich",
1319 BuilderStep::PollEnrich { .. } => "poll_enrich",
1320 BuilderStep::Validate { .. } => "validate",
1321 BuilderStep::IdempotentConsumer { .. } => "idempotent_consumer",
1322 BuilderStep::ClaimCheck { .. } => "claim_check",
1323 BuilderStep::Cache { .. } => "cache",
1324 BuilderStep::CacheInvalidate { .. } => "cache_invalidate",
1325 BuilderStep::CacheClear { .. } => "cache_clear",
1326 BuilderStep::CacheStats { .. } => "cache_stats",
1327 BuilderStep::CachePeekStale { .. } => "cache_peek_stale",
1328 BuilderStep::Sampling { .. } => "sampling",
1329 BuilderStep::Sort { .. } => "sort",
1330 BuilderStep::DeclarativeDoTry { .. } => "do_try",
1331 BuilderStep::Resequence { .. } => "resequence",
1332 }
1333}
1334
1335impl StepAccumulator for RouteBuilder {
1336 fn steps_mut(&mut self) -> &mut Vec<BuilderStep> {
1337 &mut self.steps
1338 }
1339}
1340
1341pub struct SplitBuilder {
1349 parent: RouteBuilder,
1350 config: SplitterConfig,
1351 steps: Vec<BuilderStep>,
1352}
1353
1354impl SplitBuilder {
1355 pub fn filter<F>(self, predicate: F) -> FilterInSplitBuilder
1357 where
1358 F: Fn(&Exchange) -> bool + Send + Sync + 'static,
1359 {
1360 FilterInSplitBuilder {
1361 parent: self,
1362 predicate: camel_api::FilterPredicate::new(predicate),
1363 steps: vec![],
1364 }
1365 }
1366
1367 pub fn end_split(mut self) -> RouteBuilder {
1370 let split_step = BuilderStep::Split {
1371 config: self.config,
1372 steps: self.steps,
1373 };
1374 self.parent.steps.push(split_step);
1375 self.parent
1376 }
1377}
1378
1379impl StepAccumulator for SplitBuilder {
1380 fn steps_mut(&mut self) -> &mut Vec<BuilderStep> {
1381 &mut self.steps
1382 }
1383}
1384
1385pub struct FilterBuilder {
1387 parent: RouteBuilder,
1388 predicate: FilterPredicate,
1389 steps: Vec<BuilderStep>,
1390}
1391
1392impl FilterBuilder {
1393 pub fn end_filter(mut self) -> RouteBuilder {
1396 let step = BuilderStep::Filter {
1397 predicate: self.predicate,
1398 steps: self.steps,
1399 };
1400 self.parent.steps.push(step);
1401 self.parent
1402 }
1403}
1404
1405impl StepAccumulator for FilterBuilder {
1406 fn steps_mut(&mut self) -> &mut Vec<BuilderStep> {
1407 &mut self.steps
1408 }
1409}
1410
1411pub struct FilterInSplitBuilder {
1413 parent: SplitBuilder,
1414 predicate: FilterPredicate,
1415 steps: Vec<BuilderStep>,
1416}
1417
1418impl FilterInSplitBuilder {
1419 pub fn end_filter(mut self) -> SplitBuilder {
1421 let step = BuilderStep::Filter {
1422 predicate: self.predicate,
1423 steps: self.steps,
1424 };
1425 self.parent.steps.push(step);
1426 self.parent
1427 }
1428}
1429
1430impl StepAccumulator for FilterInSplitBuilder {
1431 fn steps_mut(&mut self) -> &mut Vec<BuilderStep> {
1432 &mut self.steps
1433 }
1434}
1435
1436pub struct ChoiceBuilder {
1443 parent: RouteBuilder,
1444 whens: Vec<WhenStep>,
1445 _otherwise: Option<Vec<BuilderStep>>,
1446}
1447
1448impl ChoiceBuilder {
1449 pub fn when<F>(self, predicate: F) -> WhenBuilder
1452 where
1453 F: Fn(&Exchange) -> bool + Send + Sync + 'static,
1454 {
1455 WhenBuilder {
1456 parent: self,
1457 predicate: camel_api::FilterPredicate::new(predicate),
1458 steps: vec![],
1459 }
1460 }
1461
1462 pub fn otherwise(self) -> OtherwiseBuilder {
1466 OtherwiseBuilder {
1467 parent: self,
1468 steps: vec![],
1469 }
1470 }
1471
1472 pub fn end_choice(mut self) -> RouteBuilder {
1476 let step = BuilderStep::Choice {
1477 whens: self.whens,
1478 otherwise: self._otherwise,
1479 };
1480 self.parent.steps.push(step);
1481 self.parent
1482 }
1483}
1484
1485pub struct WhenBuilder {
1487 parent: ChoiceBuilder,
1488 predicate: camel_api::FilterPredicate,
1489 steps: Vec<BuilderStep>,
1490}
1491
1492impl WhenBuilder {
1493 pub fn end_when(mut self) -> ChoiceBuilder {
1496 self.parent.whens.push(WhenStep {
1497 predicate: self.predicate,
1498 steps: self.steps,
1499 });
1500 self.parent
1501 }
1502}
1503
1504impl StepAccumulator for WhenBuilder {
1505 fn steps_mut(&mut self) -> &mut Vec<BuilderStep> {
1506 &mut self.steps
1507 }
1508}
1509
1510pub struct OtherwiseBuilder {
1512 parent: ChoiceBuilder,
1513 steps: Vec<BuilderStep>,
1514}
1515
1516impl OtherwiseBuilder {
1517 pub fn end_otherwise(self) -> ChoiceBuilder {
1519 let OtherwiseBuilder { mut parent, steps } = self;
1520 parent._otherwise = Some(steps);
1521 parent
1522 }
1523}
1524
1525impl StepAccumulator for OtherwiseBuilder {
1526 fn steps_mut(&mut self) -> &mut Vec<BuilderStep> {
1527 &mut self.steps
1528 }
1529}
1530
1531pub struct MulticastBuilder {
1539 parent: RouteBuilder,
1540 steps: Vec<BuilderStep>,
1541 config: MulticastConfig,
1542}
1543
1544impl MulticastBuilder {
1545 pub fn parallel(mut self, parallel: bool) -> Self {
1546 self.config = self.config.parallel(parallel);
1547 self
1548 }
1549
1550 pub fn parallel_limit(mut self, limit: usize) -> Self {
1551 self.config = self.config.parallel_limit(limit);
1552 self
1553 }
1554
1555 pub fn stop_on_exception(mut self, stop: bool) -> Self {
1556 self.config = self.config.stop_on_exception(stop);
1557 self
1558 }
1559
1560 pub fn timeout(mut self, duration: std::time::Duration) -> Self {
1561 self.config = self.config.timeout(duration);
1562 self
1563 }
1564
1565 pub fn aggregation(mut self, strategy: MulticastStrategy) -> Self {
1566 self.config = self.config.aggregation(strategy);
1567 self
1568 }
1569
1570 pub fn end_multicast(mut self) -> RouteBuilder {
1571 let step = BuilderStep::Multicast {
1572 steps: self.steps,
1573 config: self.config,
1574 };
1575 self.parent.steps.push(step);
1576 self.parent
1577 }
1578}
1579
1580impl StepAccumulator for MulticastBuilder {
1581 fn steps_mut(&mut self) -> &mut Vec<BuilderStep> {
1582 &mut self.steps
1583 }
1584}
1585
1586pub struct ThrottleBuilder {
1594 parent: RouteBuilder,
1595 config: ThrottlerConfig,
1596 steps: Vec<BuilderStep>,
1597}
1598
1599impl ThrottleBuilder {
1600 pub fn strategy(mut self, strategy: ThrottleStrategy) -> Self {
1606 self.config = self.config.strategy(strategy);
1607 self
1608 }
1609
1610 pub fn end_throttle(mut self) -> RouteBuilder {
1613 let step = BuilderStep::Throttle {
1614 config: self.config,
1615 steps: self.steps,
1616 };
1617 self.parent.steps.push(step);
1618 self.parent
1619 }
1620}
1621
1622impl StepAccumulator for ThrottleBuilder {
1623 fn steps_mut(&mut self) -> &mut Vec<BuilderStep> {
1624 &mut self.steps
1625 }
1626}
1627
1628pub struct LoopBuilder {
1630 parent: RouteBuilder,
1631 config: LoopConfig,
1632 steps: Vec<BuilderStep>,
1633}
1634
1635impl LoopBuilder {
1636 pub fn loop_count(self, count: usize) -> LoopInLoopBuilder {
1637 LoopInLoopBuilder {
1638 parent: self,
1639 config: LoopConfig::new(LoopMode::Count(count)),
1640 steps: vec![],
1641 }
1642 }
1643
1644 pub fn loop_while<F>(self, predicate: F) -> LoopInLoopBuilder
1645 where
1646 F: Fn(&Exchange) -> bool + Send + Sync + 'static,
1647 {
1648 LoopInLoopBuilder {
1649 parent: self,
1650 config: LoopConfig::new(LoopMode::While(camel_api::FilterPredicate::new(predicate))),
1651 steps: vec![],
1652 }
1653 }
1654
1655 pub fn end_loop(mut self) -> RouteBuilder {
1656 let step = BuilderStep::Loop {
1657 config: self.config,
1658 steps: self.steps,
1659 };
1660 self.parent.steps.push(step);
1661 self.parent
1662 }
1663}
1664
1665impl StepAccumulator for LoopBuilder {
1666 fn steps_mut(&mut self) -> &mut Vec<BuilderStep> {
1667 &mut self.steps
1668 }
1669}
1670
1671pub struct LoopInLoopBuilder {
1672 parent: LoopBuilder,
1673 config: LoopConfig,
1674 steps: Vec<BuilderStep>,
1675}
1676
1677impl LoopInLoopBuilder {
1678 pub fn end_loop(mut self) -> LoopBuilder {
1679 let step = BuilderStep::Loop {
1680 config: self.config,
1681 steps: self.steps,
1682 };
1683 self.parent.steps.push(step);
1684 self.parent
1685 }
1686}
1687
1688impl StepAccumulator for LoopInLoopBuilder {
1689 fn steps_mut(&mut self) -> &mut Vec<BuilderStep> {
1690 &mut self.steps
1691 }
1692}
1693
1694pub struct LoadBalancerBuilder {
1702 parent: RouteBuilder,
1703 config: LoadBalancerConfig,
1704 steps: Vec<BuilderStep>,
1705}
1706
1707impl LoadBalancerBuilder {
1708 pub fn round_robin(mut self) -> Self {
1710 self.config = LoadBalancerConfig::round_robin();
1711 self
1712 }
1713
1714 pub fn random(mut self) -> Self {
1716 self.config = LoadBalancerConfig::random();
1717 self
1718 }
1719
1720 pub fn weighted(mut self, weights: Vec<(String, u32)>) -> Self {
1725 self.config = LoadBalancerConfig::weighted(weights);
1726 self
1727 }
1728
1729 pub fn failover(mut self) -> Self {
1734 self.config = LoadBalancerConfig::failover();
1735 self
1736 }
1737
1738 pub fn end_load_balance(mut self) -> RouteBuilder {
1741 let step = BuilderStep::LoadBalance {
1742 config: self.config,
1743 steps: self.steps,
1744 };
1745 self.parent.steps.push(step);
1746 self.parent
1747 }
1748}
1749
1750impl StepAccumulator for LoadBalancerBuilder {
1751 fn steps_mut(&mut self) -> &mut Vec<BuilderStep> {
1752 &mut self.steps
1753 }
1754}
1755
1756#[cfg(test)]
1761mod tests {
1762 use super::*;
1763 use camel_api::SpanKindHint;
1764 use camel_api::error_handler::ErrorHandlerConfig;
1765 use camel_api::load_balancer::LoadBalanceStrategy;
1766 use camel_api::{Exchange, Message};
1767 use camel_core::route::BuilderStep;
1768 use std::sync::Arc;
1769 use std::time::Duration;
1770 use tower::{Service, ServiceExt};
1771
1772 #[test]
1773 fn test_builder_from_creates_definition() {
1774 let definition = RouteBuilder::from("timer:tick")
1775 .route_id("test-route")
1776 .build()
1777 .unwrap();
1778 assert_eq!(definition.from_uri(), "timer:tick");
1779 }
1780
1781 #[test]
1782 fn test_builder_empty_from_uri_errors() {
1783 let result = RouteBuilder::from("").route_id("test-route").build();
1784 assert!(result.is_err());
1785 }
1786
1787 #[test]
1788 fn test_build_rejects_schemeless_uri() {
1789 let result = RouteBuilder::from("no-scheme-here")
1790 .route_id("test-route")
1791 .build();
1792 match result {
1793 Err(err) => {
1794 let err_msg = format!("{err}");
1795 assert!(
1796 err_msg.contains("scheme"),
1797 "expected scheme-related error, got: {err_msg}"
1798 );
1799 }
1800 Ok(_) => panic!("schemeless URI should fail"),
1801 }
1802 }
1803
1804 #[test]
1805 fn test_build_rejects_empty_scheme_uri() {
1806 let result = RouteBuilder::from(":missing-scheme")
1807 .route_id("test-route")
1808 .build();
1809 match result {
1810 Err(err) => {
1811 let err_msg = format!("{err}");
1812 assert!(
1813 err_msg.contains("scheme"),
1814 "expected scheme-related error, got: {err_msg}"
1815 );
1816 }
1817 Ok(_) => panic!("empty-scheme URI should fail"),
1818 }
1819 }
1820
1821 #[test]
1822 fn test_build_accepts_valid_uri() {
1823 let result = RouteBuilder::from("timer:tick")
1824 .route_id("test-route")
1825 .build();
1826 assert!(result.is_ok());
1827 }
1828
1829 #[test]
1830 fn test_build_canonical_rejects_schemeless_uri() {
1831 let result = RouteBuilder::from("no-scheme-here")
1832 .route_id("test-route")
1833 .build_canonical();
1834 assert!(result.is_err());
1835 }
1836
1837 #[test]
1838 fn test_builder_to_adds_step() {
1839 let definition = RouteBuilder::from("timer:tick")
1840 .route_id("test-route")
1841 .to("log:info")
1842 .build()
1843 .unwrap();
1844
1845 assert_eq!(definition.from_uri(), "timer:tick");
1846 assert!(matches!(&definition.steps()[0], BuilderStep::To(uri) if uri == "log:info"));
1848 }
1849
1850 #[test]
1851 fn test_builder_filter_adds_filter_step() {
1852 let definition = RouteBuilder::from("timer:tick")
1853 .route_id("test-route")
1854 .filter(|_ex| true)
1855 .to("mock:result")
1856 .end_filter()
1857 .build()
1858 .unwrap();
1859
1860 assert!(matches!(&definition.steps()[0], BuilderStep::Filter { .. }));
1861 }
1862
1863 #[test]
1864 fn test_builder_set_header_adds_processor_step() {
1865 let definition = RouteBuilder::from("timer:tick")
1866 .route_id("test-route")
1867 .set_header("key", Value::String("value".into()))
1868 .build()
1869 .unwrap();
1870
1871 assert!(matches!(&definition.steps()[0], BuilderStep::Processor(_)));
1872 }
1873
1874 #[test]
1875 fn test_builder_map_body_adds_processor_step() {
1876 let definition = RouteBuilder::from("timer:tick")
1877 .route_id("test-route")
1878 .map_body(|body| body)
1879 .build()
1880 .unwrap();
1881
1882 assert!(matches!(&definition.steps()[0], BuilderStep::Processor(_)));
1883 }
1884
1885 #[test]
1886 fn test_builder_process_adds_processor_step() {
1887 let definition = RouteBuilder::from("timer:tick")
1888 .route_id("test-route")
1889 .process(|ex| async move { Ok(ex) })
1890 .build()
1891 .unwrap();
1892
1893 assert!(matches!(&definition.steps()[0], BuilderStep::Processor(_)));
1894 }
1895
1896 #[test]
1897 fn test_builder_chain_multiple_steps() {
1898 let definition = RouteBuilder::from("timer:tick")
1899 .route_id("test-route")
1900 .set_header("source", Value::String("timer".into()))
1901 .filter(|ex| ex.input.header("source").is_some())
1902 .to("log:info")
1903 .end_filter()
1904 .to("mock:result")
1905 .build()
1906 .unwrap();
1907
1908 assert_eq!(definition.steps().len(), 3); assert!(matches!(&definition.steps()[0], BuilderStep::Processor(_))); assert!(matches!(&definition.steps()[1], BuilderStep::Filter { .. })); assert!(matches!(&definition.steps()[2], BuilderStep::To(uri) if uri == "mock:result"));
1912 }
1913
1914 #[test]
1915 fn test_loop_count_builder() {
1916 use camel_api::loop_eip::LoopMode;
1917
1918 let def = RouteBuilder::from("direct:start")
1919 .route_id("loop-test")
1920 .loop_count(3)
1921 .to("mock:inside")
1922 .end_loop()
1923 .to("mock:after")
1924 .build()
1925 .unwrap();
1926
1927 assert_eq!(def.steps().len(), 2);
1928 match &def.steps()[0] {
1929 BuilderStep::Loop { config, steps } => {
1930 assert!(matches!(config.mode, LoopMode::Count(3)));
1931 assert_eq!(steps.len(), 1);
1932 }
1933 other => panic!("Expected Loop, got {:?}", other),
1934 }
1935 assert!(matches!(def.steps()[1], BuilderStep::To(_)));
1936 }
1937
1938 #[test]
1939 fn test_loop_while_builder() {
1940 use camel_api::loop_eip::LoopMode;
1941
1942 let def = RouteBuilder::from("direct:start")
1943 .route_id("loop-while-test")
1944 .loop_while(|_ex| true)
1945 .to("mock:retry")
1946 .end_loop()
1947 .build()
1948 .unwrap();
1949
1950 assert_eq!(def.steps().len(), 1);
1951 match &def.steps()[0] {
1952 BuilderStep::Loop { config, steps } => {
1953 assert!(matches!(config.mode, LoopMode::While(_)));
1954 assert_eq!(steps.len(), 1);
1955 }
1956 other => panic!("Expected Loop, got {:?}", other),
1957 }
1958 }
1959
1960 #[test]
1961 fn test_nested_loop_builder() {
1962 use camel_api::loop_eip::LoopMode;
1963
1964 let def = RouteBuilder::from("direct:start")
1965 .route_id("nested-loop-test")
1966 .loop_count(2)
1967 .to("mock:outer")
1968 .loop_count(3)
1969 .to("mock:inner")
1970 .end_loop()
1971 .end_loop()
1972 .to("mock:after")
1973 .build()
1974 .unwrap();
1975
1976 assert_eq!(def.steps().len(), 2);
1977 match &def.steps()[0] {
1978 BuilderStep::Loop { steps, .. } => {
1979 assert_eq!(steps.len(), 2);
1980 match &steps[1] {
1981 BuilderStep::Loop {
1982 config,
1983 steps: inner_steps,
1984 } => {
1985 assert!(matches!(config.mode, LoopMode::Count(3)));
1986 assert_eq!(inner_steps.len(), 1);
1987 }
1988 other => panic!("Expected nested Loop, got {:?}", other),
1989 }
1990 }
1991 other => panic!("Expected outer Loop, got {:?}", other),
1992 }
1993 }
1994
1995 #[tokio::test]
2000 async fn test_set_header_processor_works() {
2001 let mut svc = SetHeader::new(IdentityProcessor, "greeting", Value::String("hello".into()));
2002 let exchange = Exchange::new(Message::new("test"));
2003 let result = svc.call(exchange).await.unwrap();
2004 assert_eq!(
2005 result.input.header("greeting"),
2006 Some(&Value::String("hello".into()))
2007 );
2008 }
2009
2010 #[tokio::test]
2011 async fn test_filter_processor_passes() {
2012 use camel_api::BoxProcessorExt;
2013 use camel_processor::FilterService;
2014
2015 let sub = BoxProcessor::from_fn(|ex| Box::pin(async move { Ok(ex) }));
2016 let mut svc =
2017 FilterService::new(|ex: &Exchange| ex.input.body.as_text() == Some("pass"), sub);
2018 let exchange = Exchange::new(Message::new("pass"));
2019 let result = svc.ready().await.unwrap().call(exchange).await.unwrap();
2020 assert_eq!(result.input.body.as_text(), Some("pass"));
2021 }
2022
2023 #[tokio::test]
2024 async fn test_filter_processor_blocks() {
2025 use camel_api::BoxProcessorExt;
2026 use camel_processor::FilterService;
2027
2028 let sub = BoxProcessor::from_fn(|_ex| {
2029 Box::pin(async move { Err(CamelError::ProcessorError("should not reach".into())) })
2030 });
2031 let mut svc =
2032 FilterService::new(|ex: &Exchange| ex.input.body.as_text() == Some("pass"), sub);
2033 let exchange = Exchange::new(Message::new("reject"));
2034 let result = svc.ready().await.unwrap().call(exchange).await.unwrap();
2035 assert_eq!(result.input.body.as_text(), Some("reject"));
2036 }
2037
2038 #[tokio::test]
2039 async fn test_map_body_processor_works() {
2040 let mapper = MapBody::new(IdentityProcessor, |body: Body| {
2041 if let Some(text) = body.as_text() {
2042 Body::Text(text.to_uppercase())
2043 } else {
2044 body
2045 }
2046 });
2047 let exchange = Exchange::new(Message::new("hello"));
2048 let result = mapper.oneshot(exchange).await.unwrap();
2049 assert_eq!(result.input.body.as_text(), Some("HELLO"));
2050 }
2051
2052 #[tokio::test]
2053 async fn test_process_custom_processor_works() {
2054 let processor = ProcessorFn::new(|mut ex: Exchange| async move {
2055 ex.set_property("custom", Value::Bool(true));
2056 Ok(ex)
2057 });
2058 let exchange = Exchange::new(Message::default());
2059 let result = processor.oneshot(exchange).await.unwrap();
2060 assert_eq!(result.property("custom"), Some(&Value::Bool(true)));
2061 }
2062
2063 #[tokio::test]
2068 async fn test_compose_pipeline_runs_steps_in_order() {
2069 use camel_core::route::{CompiledStep, PipelineRuntimeCtx, compose_pipeline};
2070
2071 let processors = vec![
2072 CompiledStep::Process {
2073 kind_hint: SpanKindHint::Internal,
2074 processor: BoxProcessor::new(SetHeader::new(
2075 IdentityProcessor,
2076 "step",
2077 Value::String("one".into()),
2078 )),
2079 body_contract: None,
2080 lifecycle: None,
2081 label: None,
2082 to_uri: None,
2083 },
2084 CompiledStep::Process {
2085 kind_hint: SpanKindHint::Internal,
2086 processor: BoxProcessor::new(MapBody::new(IdentityProcessor, |body: Body| {
2087 if let Some(text) = body.as_text() {
2088 Body::Text(format!("{}-processed", text))
2089 } else {
2090 body
2091 }
2092 })),
2093 body_contract: None,
2094 lifecycle: None,
2095 label: None,
2096 to_uri: None,
2097 },
2098 ];
2099
2100 let pipeline = compose_pipeline(processors, PipelineRuntimeCtx::compile_time());
2101 let exchange = Exchange::new(Message::new("hello"));
2102 let result = pipeline.oneshot(exchange).await.unwrap();
2103
2104 assert_eq!(
2105 result.input.header("step"),
2106 Some(&Value::String("one".into()))
2107 );
2108 assert_eq!(result.input.body.as_text(), Some("hello-processed"));
2109 }
2110
2111 #[tokio::test]
2112 async fn test_compose_pipeline_empty_is_identity() {
2113 use camel_core::route::{PipelineRuntimeCtx, compose_pipeline};
2114
2115 let pipeline = compose_pipeline(vec![], PipelineRuntimeCtx::compile_time());
2116 let exchange = Exchange::new(Message::new("unchanged"));
2117 let result = pipeline.oneshot(exchange).await.unwrap();
2118 assert_eq!(result.input.body.as_text(), Some("unchanged"));
2119 }
2120
2121 #[test]
2126 fn test_builder_circuit_breaker_sets_config() {
2127 use camel_api::circuit_breaker::CircuitBreakerConfig;
2128
2129 let config = CircuitBreakerConfig::new().failure_threshold(5);
2130 let definition = RouteBuilder::from("timer:tick")
2131 .route_id("test-route")
2132 .circuit_breaker(config)
2133 .build()
2134 .unwrap();
2135
2136 let cb = definition
2137 .circuit_breaker_config()
2138 .expect("circuit breaker should be set");
2139 assert_eq!(cb.failure_threshold, 5);
2140 }
2141
2142 #[test]
2143 fn test_builder_circuit_breaker_with_error_handler() {
2144 use camel_api::circuit_breaker::CircuitBreakerConfig;
2145 use camel_api::error_handler::ErrorHandlerConfig;
2146
2147 let cb_config = CircuitBreakerConfig::new().failure_threshold(3);
2148 let eh_config = ErrorHandlerConfig::log_only();
2149
2150 let definition = RouteBuilder::from("timer:tick")
2151 .route_id("test-route")
2152 .to("log:info")
2153 .circuit_breaker(cb_config)
2154 .error_handler(eh_config)
2155 .build()
2156 .unwrap();
2157
2158 assert!(
2159 definition.circuit_breaker_config().is_some(),
2160 "circuit breaker config should be set"
2161 );
2162 }
2164
2165 #[test]
2166 fn test_builder_on_exception_shorthand_multiple_clauses_preserve_order() {
2167 let definition = RouteBuilder::from("direct:start")
2168 .route_id("test-route")
2169 .dead_letter_channel("log:dlc")
2170 .on_exception(|e| matches!(e, CamelError::Io(_)))
2171 .retry(3)
2172 .handled_by("log:io")
2173 .end_on_exception()
2174 .on_exception(|e| matches!(e, CamelError::ProcessorError(_)))
2175 .retry(1)
2176 .end_on_exception()
2177 .to("mock:out")
2178 .build()
2179 .expect("route should build");
2180
2181 let cfg = definition
2182 .error_handler_config()
2183 .expect("error handler should be set");
2184 assert_eq!(cfg.policies.len(), 2);
2185 assert_eq!(cfg.dlc_uri.as_deref(), Some("log:dlc"));
2186 assert_eq!(
2187 cfg.policies[0].retry.as_ref().map(|p| p.max_attempts),
2188 Some(3)
2189 );
2190 assert_eq!(cfg.policies[0].handled_by.as_deref(), Some("log:io"));
2191 assert_eq!(
2192 cfg.policies[1].retry.as_ref().map(|p| p.max_attempts),
2193 Some(1)
2194 );
2195 }
2196
2197 #[test]
2198 fn test_builder_on_exception_mixed_mode_rejected() {
2199 let result = RouteBuilder::from("direct:start")
2200 .route_id("test-route")
2201 .error_handler(ErrorHandlerConfig::log_only())
2202 .on_exception(|_e| true)
2203 .end_on_exception()
2204 .to("mock:out")
2205 .build();
2206
2207 let err = result.err().expect("mixed mode should fail with an error");
2208
2209 assert!(
2210 format!("{err}").contains("mixed error handler modes"),
2211 "unexpected error: {err}"
2212 );
2213 }
2214
2215 #[test]
2216 fn test_builder_on_exception_backoff_and_jitter_without_retry_noop() {
2217 let definition = RouteBuilder::from("direct:start")
2218 .route_id("test-route")
2219 .on_exception(|_e| true)
2220 .with_backoff(Duration::from_millis(5), 3.0, Duration::from_millis(100))
2221 .with_jitter(0.5)
2222 .end_on_exception()
2223 .to("mock:out")
2224 .build()
2225 .expect("route should build");
2226
2227 let cfg = definition
2228 .error_handler_config()
2229 .expect("error handler should be set");
2230 assert_eq!(cfg.policies.len(), 1);
2231 assert!(cfg.policies[0].retry.is_none());
2232 }
2233
2234 #[test]
2235 fn test_builder_dead_letter_channel_without_on_exception_sets_dlc() {
2236 let definition = RouteBuilder::from("direct:start")
2237 .route_id("test-route")
2238 .dead_letter_channel("log:dlc")
2239 .to("mock:out")
2240 .build()
2241 .expect("route should build");
2242
2243 let cfg = definition
2244 .error_handler_config()
2245 .expect("error handler should be set");
2246 assert_eq!(cfg.dlc_uri.as_deref(), Some("log:dlc"));
2247 assert!(cfg.policies.is_empty());
2248 }
2249
2250 #[test]
2251 fn test_builder_dead_letter_channel_called_twice_uses_latest_and_keeps_policies() {
2252 let definition = RouteBuilder::from("direct:start")
2253 .route_id("test-route")
2254 .dead_letter_channel("log:first")
2255 .on_exception(|e| matches!(e, CamelError::Io(_)))
2256 .retry(2)
2257 .end_on_exception()
2258 .dead_letter_channel("log:second")
2259 .to("mock:out")
2260 .build()
2261 .expect("route should build");
2262
2263 let cfg = definition
2264 .error_handler_config()
2265 .expect("error handler should be set");
2266 assert_eq!(cfg.dlc_uri.as_deref(), Some("log:second"));
2267 assert_eq!(cfg.policies.len(), 1);
2268 assert_eq!(
2269 cfg.policies[0].retry.as_ref().map(|p| p.max_attempts),
2270 Some(2)
2271 );
2272 }
2273
2274 #[test]
2275 fn test_builder_on_exception_without_dlc_defaults_to_log_only() {
2276 let definition = RouteBuilder::from("direct:start")
2277 .route_id("test-route")
2278 .on_exception(|e| matches!(e, CamelError::ProcessorError(_)))
2279 .retry(1)
2280 .end_on_exception()
2281 .to("mock:out")
2282 .build()
2283 .expect("route should build");
2284
2285 let cfg = definition
2286 .error_handler_config()
2287 .expect("error handler should be set");
2288 assert!(cfg.dlc_uri.is_none());
2289 assert_eq!(cfg.policies.len(), 1);
2290 }
2291
2292 #[test]
2293 fn test_builder_error_handler_explicit_overwrite_stays_explicit_mode() {
2294 let first = ErrorHandlerConfig::dead_letter_channel("log:first");
2295 let second = ErrorHandlerConfig::dead_letter_channel("log:second");
2296
2297 let definition = RouteBuilder::from("direct:start")
2298 .route_id("test-route")
2299 .error_handler(first)
2300 .error_handler(second)
2301 .to("mock:out")
2302 .build()
2303 .expect("route should build");
2304
2305 let cfg = definition
2306 .error_handler_config()
2307 .expect("error handler should be set");
2308 assert_eq!(cfg.dlc_uri.as_deref(), Some("log:second"));
2309 }
2310
2311 #[test]
2314 fn test_split_builder_typestate() {
2315 use camel_api::splitter::{SplitterConfig, split_body_lines};
2316
2317 let definition = RouteBuilder::from("timer:test?period=1000")
2319 .route_id("test-route")
2320 .split(SplitterConfig::new(split_body_lines()))
2321 .to("mock:per-fragment")
2322 .end_split()
2323 .to("mock:final")
2324 .build()
2325 .unwrap();
2326
2327 assert_eq!(definition.steps().len(), 2);
2329 }
2330
2331 #[test]
2332 fn test_split_builder_steps_collected() {
2333 use camel_api::splitter::{SplitterConfig, split_body_lines};
2334
2335 let definition = RouteBuilder::from("timer:test?period=1000")
2336 .route_id("test-route")
2337 .split(SplitterConfig::new(split_body_lines()))
2338 .set_header("fragment", Value::String("yes".into()))
2339 .to("mock:per-fragment")
2340 .end_split()
2341 .build()
2342 .unwrap();
2343
2344 assert_eq!(definition.steps().len(), 1);
2346 match &definition.steps()[0] {
2347 BuilderStep::Split { steps, .. } => {
2348 assert_eq!(steps.len(), 2); }
2350 other => panic!("Expected Split, got {:?}", other),
2351 }
2352 }
2353
2354 #[test]
2355 fn test_split_builder_config_propagated() {
2356 use camel_api::splitter::{AggregationStrategy, SplitterConfig, split_body_lines};
2357
2358 let definition = RouteBuilder::from("timer:test?period=1000")
2359 .route_id("test-route")
2360 .split(
2361 SplitterConfig::new(split_body_lines())
2362 .parallel(true)
2363 .parallel_limit(4)
2364 .aggregation(AggregationStrategy::CollectAll),
2365 )
2366 .to("mock:per-fragment")
2367 .end_split()
2368 .build()
2369 .unwrap();
2370
2371 match &definition.steps()[0] {
2372 BuilderStep::Split { config, .. } => {
2373 assert!(config.parallel);
2374 assert_eq!(config.parallel_limit, Some(4));
2375 assert!(matches!(
2376 config.aggregation,
2377 AggregationStrategy::CollectAll
2378 ));
2379 }
2380 other => panic!("Expected Split, got {:?}", other),
2381 }
2382 }
2383
2384 #[test]
2385 fn test_aggregate_builder_adds_step() {
2386 use camel_api::aggregator::AggregatorConfig;
2387 use camel_core::route::BuilderStep;
2388
2389 let definition = RouteBuilder::from("timer:tick")
2390 .route_id("test-route")
2391 .aggregate(
2392 AggregatorConfig::correlate_by("key")
2393 .complete_when_size(2)
2394 .build()
2395 .unwrap(),
2396 )
2397 .build()
2398 .unwrap();
2399
2400 assert_eq!(definition.steps().len(), 1);
2401 assert!(matches!(
2402 definition.steps()[0],
2403 BuilderStep::Aggregate { .. }
2404 ));
2405 }
2406
2407 #[test]
2408 fn test_aggregate_in_split_builder() {
2409 use camel_api::aggregator::AggregatorConfig;
2410 use camel_api::splitter::{SplitterConfig, split_body_lines};
2411 use camel_core::route::BuilderStep;
2412
2413 let definition = RouteBuilder::from("timer:tick")
2414 .route_id("test-route")
2415 .split(SplitterConfig::new(split_body_lines()))
2416 .aggregate(
2417 AggregatorConfig::correlate_by("key")
2418 .complete_when_size(1)
2419 .build()
2420 .unwrap(),
2421 )
2422 .end_split()
2423 .build()
2424 .unwrap();
2425
2426 assert_eq!(definition.steps().len(), 1);
2427 if let BuilderStep::Split { steps, .. } = &definition.steps()[0] {
2428 assert!(matches!(steps[0], BuilderStep::Aggregate { .. }));
2429 } else {
2430 panic!("expected Split step");
2431 }
2432 }
2433
2434 #[test]
2437 fn test_builder_set_body_static_adds_processor() {
2438 let definition = RouteBuilder::from("timer:tick")
2439 .route_id("test-route")
2440 .set_body("fixed")
2441 .build()
2442 .unwrap();
2443 assert!(matches!(&definition.steps()[0], BuilderStep::Processor(_)));
2444 }
2445
2446 #[test]
2447 fn test_builder_set_body_fn_adds_processor() {
2448 let definition = RouteBuilder::from("timer:tick")
2449 .route_id("test-route")
2450 .set_body_fn(|_ex: &Exchange| Body::Text("dynamic".into()))
2451 .build()
2452 .unwrap();
2453 assert!(matches!(&definition.steps()[0], BuilderStep::Processor(_)));
2454 }
2455
2456 #[test]
2457 fn transform_alias_produces_same_as_set_body() {
2458 let route_transform = RouteBuilder::from("timer:tick")
2459 .route_id("test-route")
2460 .transform("hello")
2461 .build()
2462 .unwrap();
2463
2464 let route_set_body = RouteBuilder::from("timer:tick")
2465 .route_id("test-route")
2466 .set_body("hello")
2467 .build()
2468 .unwrap();
2469
2470 assert_eq!(route_transform.steps().len(), route_set_body.steps().len());
2471 }
2472
2473 #[test]
2474 fn test_builder_set_header_fn_adds_processor() {
2475 let definition = RouteBuilder::from("timer:tick")
2476 .route_id("test-route")
2477 .set_header_fn("k", |_ex: &Exchange| Value::String("v".into()))
2478 .build()
2479 .unwrap();
2480 assert!(matches!(&definition.steps()[0], BuilderStep::Processor(_)));
2481 }
2482
2483 #[tokio::test]
2484 async fn test_set_body_static_processor_works() {
2485 use camel_core::route::{CompiledStep, PipelineRuntimeCtx, compose_pipeline};
2486 let def = RouteBuilder::from("t:t")
2487 .route_id("test-route")
2488 .set_body("replaced")
2489 .build()
2490 .unwrap();
2491 let pipeline = compose_pipeline(
2492 def.steps()
2493 .iter()
2494 .filter_map(|s| {
2495 if let BuilderStep::Processor(op) = s {
2496 Some(op.0.clone())
2497 } else {
2498 None
2499 }
2500 })
2501 .map(|p| CompiledStep::Process {
2502 kind_hint: SpanKindHint::Internal,
2503 processor: p,
2504 body_contract: None,
2505 lifecycle: None,
2506 label: None,
2507 to_uri: None,
2508 })
2509 .collect(),
2510 PipelineRuntimeCtx::compile_time(),
2511 );
2512 let exchange = Exchange::new(Message::new("original"));
2513 let result = pipeline.oneshot(exchange).await.unwrap();
2514 assert_eq!(result.input.body.as_text(), Some("replaced"));
2515 }
2516
2517 #[tokio::test]
2518 async fn test_set_body_fn_processor_works() {
2519 use camel_core::route::{CompiledStep, PipelineRuntimeCtx, compose_pipeline};
2520 let def = RouteBuilder::from("t:t")
2521 .route_id("test-route")
2522 .set_body_fn(|ex: &Exchange| {
2523 Body::Text(ex.input.body.as_text().unwrap_or("").to_uppercase())
2524 })
2525 .build()
2526 .unwrap();
2527 let pipeline = compose_pipeline(
2528 def.steps()
2529 .iter()
2530 .filter_map(|s| {
2531 if let BuilderStep::Processor(op) = s {
2532 Some(op.0.clone())
2533 } else {
2534 None
2535 }
2536 })
2537 .map(|p| CompiledStep::Process {
2538 kind_hint: SpanKindHint::Internal,
2539 processor: p,
2540 body_contract: None,
2541 lifecycle: None,
2542 label: None,
2543 to_uri: None,
2544 })
2545 .collect(),
2546 PipelineRuntimeCtx::compile_time(),
2547 );
2548 let exchange = Exchange::new(Message::new("hello"));
2549 let result = pipeline.oneshot(exchange).await.unwrap();
2550 assert_eq!(result.input.body.as_text(), Some("HELLO"));
2551 }
2552
2553 #[tokio::test]
2554 async fn test_set_header_fn_processor_works() {
2555 use camel_core::route::{CompiledStep, PipelineRuntimeCtx, compose_pipeline};
2556 let def = RouteBuilder::from("t:t")
2557 .route_id("test-route")
2558 .set_header_fn("echo", |ex: &Exchange| {
2559 ex.input
2560 .body
2561 .as_text()
2562 .map(|t| Value::String(t.into()))
2563 .unwrap_or(Value::Null)
2564 })
2565 .build()
2566 .unwrap();
2567 let pipeline = compose_pipeline(
2568 def.steps()
2569 .iter()
2570 .filter_map(|s| {
2571 if let BuilderStep::Processor(op) = s {
2572 Some(op.0.clone())
2573 } else {
2574 None
2575 }
2576 })
2577 .map(|p| CompiledStep::Process {
2578 kind_hint: SpanKindHint::Internal,
2579 processor: p,
2580 body_contract: None,
2581 lifecycle: None,
2582 label: None,
2583 to_uri: None,
2584 })
2585 .collect(),
2586 PipelineRuntimeCtx::compile_time(),
2587 );
2588 let exchange = Exchange::new(Message::new("ping"));
2589 let result = pipeline.oneshot(exchange).await.unwrap();
2590 assert_eq!(
2591 result.input.header("echo"),
2592 Some(&Value::String("ping".into()))
2593 );
2594 }
2595
2596 #[test]
2599 fn test_filter_builder_typestate() {
2600 let result = RouteBuilder::from("timer:tick?period=50&repeatCount=1")
2601 .route_id("test-route")
2602 .filter(|_ex| true)
2603 .to("mock:inner")
2604 .end_filter()
2605 .to("mock:outer")
2606 .build();
2607 assert!(result.is_ok());
2608 }
2609
2610 #[test]
2611 fn test_filter_builder_steps_collected() {
2612 let definition = RouteBuilder::from("timer:tick?period=50&repeatCount=1")
2613 .route_id("test-route")
2614 .filter(|_ex| true)
2615 .to("mock:inner")
2616 .end_filter()
2617 .build()
2618 .unwrap();
2619
2620 assert_eq!(definition.steps().len(), 1);
2621 assert!(matches!(&definition.steps()[0], BuilderStep::Filter { .. }));
2622 }
2623
2624 #[test]
2625 fn test_wire_tap_builder_adds_step() {
2626 let definition = RouteBuilder::from("timer:tick")
2627 .route_id("test-route")
2628 .wire_tap("mock:tap")
2629 .to("mock:result")
2630 .build()
2631 .unwrap();
2632
2633 assert_eq!(definition.steps().len(), 2);
2634 assert!(
2635 matches!(&definition.steps()[0], BuilderStep::WireTap { uri } if uri == "mock:tap")
2636 );
2637 assert!(matches!(&definition.steps()[1], BuilderStep::To(uri) if uri == "mock:result"));
2638 }
2639
2640 #[test]
2643 fn test_multicast_builder_typestate() {
2644 let definition = RouteBuilder::from("timer:tick")
2645 .route_id("test-route")
2646 .multicast()
2647 .to("direct:a")
2648 .to("direct:b")
2649 .end_multicast()
2650 .to("mock:result")
2651 .build()
2652 .unwrap();
2653
2654 assert_eq!(definition.steps().len(), 2); }
2656
2657 #[test]
2658 fn test_multicast_builder_steps_collected() {
2659 let definition = RouteBuilder::from("timer:tick")
2660 .route_id("test-route")
2661 .multicast()
2662 .to("direct:a")
2663 .to("direct:b")
2664 .end_multicast()
2665 .build()
2666 .unwrap();
2667
2668 match &definition.steps()[0] {
2669 BuilderStep::Multicast { steps, .. } => {
2670 assert_eq!(steps.len(), 2);
2671 }
2672 other => panic!("Expected Multicast, got {:?}", other),
2673 }
2674 }
2675
2676 #[test]
2679 fn test_builder_concurrent_sets_concurrency() {
2680 use camel_component_api::ConcurrencyModel;
2681
2682 let definition = RouteBuilder::from("http://0.0.0.0:8080/test")
2683 .route_id("test-route")
2684 .concurrent(16)
2685 .to("log:info")
2686 .build()
2687 .unwrap();
2688
2689 assert_eq!(
2690 definition.concurrency_override(),
2691 Some(&ConcurrencyModel::Concurrent { max: Some(16) })
2692 );
2693 }
2694
2695 #[test]
2696 fn test_builder_concurrent_zero_means_unbounded() {
2697 use camel_component_api::ConcurrencyModel;
2698
2699 let definition = RouteBuilder::from("http://0.0.0.0:8080/test")
2700 .route_id("test-route")
2701 .concurrent(0)
2702 .to("log:info")
2703 .build()
2704 .unwrap();
2705
2706 assert_eq!(
2707 definition.concurrency_override(),
2708 Some(&ConcurrencyModel::Concurrent { max: None })
2709 );
2710 }
2711
2712 #[test]
2713 fn test_builder_sequential_sets_concurrency() {
2714 use camel_component_api::ConcurrencyModel;
2715
2716 let definition = RouteBuilder::from("http://0.0.0.0:8080/test")
2717 .route_id("test-route")
2718 .sequential()
2719 .to("log:info")
2720 .build()
2721 .unwrap();
2722
2723 assert_eq!(
2724 definition.concurrency_override(),
2725 Some(&ConcurrencyModel::Sequential)
2726 );
2727 }
2728
2729 #[test]
2730 fn test_builder_default_concurrency_is_none() {
2731 let definition = RouteBuilder::from("timer:tick")
2732 .route_id("test-route")
2733 .to("log:info")
2734 .build()
2735 .unwrap();
2736
2737 assert_eq!(definition.concurrency_override(), None);
2738 }
2739
2740 #[test]
2743 fn test_builder_route_id_sets_id() {
2744 let definition = RouteBuilder::from("timer:tick")
2745 .route_id("my-route")
2746 .build()
2747 .unwrap();
2748
2749 assert_eq!(definition.route_id(), "my-route");
2750 }
2751
2752 #[test]
2753 fn test_build_without_route_id_fails() {
2754 let result = RouteBuilder::from("timer:tick?period=1000")
2755 .to("log:info")
2756 .build();
2757 let err = match result {
2758 Err(e) => e.to_string(),
2759 Ok(_) => panic!("build() should fail without route_id"),
2760 };
2761 assert!(
2762 err.contains("route_id"),
2763 "error should mention route_id, got: {}",
2764 err
2765 );
2766 }
2767
2768 #[test]
2769 fn test_builder_empty_route_id_rejected() {
2770 let result = RouteBuilder::from("timer:tick").route_id("").build();
2771 let err = result.err().expect("empty route_id should be rejected");
2772 assert!(matches!(err, CamelError::RouteError(_)));
2773 }
2774
2775 #[test]
2776 fn test_builder_whitespace_route_id_rejected() {
2777 let result = RouteBuilder::from("timer:tick").route_id(" ").build();
2778 assert!(result.is_err());
2779 }
2780
2781 #[test]
2782 fn test_builder_auto_startup_false() {
2783 let definition = RouteBuilder::from("timer:tick")
2784 .route_id("test-route")
2785 .auto_startup(false)
2786 .build()
2787 .unwrap();
2788
2789 assert!(!definition.auto_startup());
2790 }
2791
2792 #[test]
2793 fn test_builder_startup_order_custom() {
2794 let definition = RouteBuilder::from("timer:tick")
2795 .route_id("test-route")
2796 .startup_order(50)
2797 .build()
2798 .unwrap();
2799
2800 assert_eq!(definition.startup_order(), 50);
2801 }
2802
2803 #[test]
2804 fn test_builder_defaults() {
2805 let definition = RouteBuilder::from("timer:tick")
2806 .route_id("test-route")
2807 .build()
2808 .unwrap();
2809
2810 assert_eq!(definition.route_id(), "test-route");
2811 assert!(definition.auto_startup());
2812 assert_eq!(definition.startup_order(), 1000);
2813 }
2814
2815 #[test]
2818 fn test_choice_builder_single_when() {
2819 let definition = RouteBuilder::from("timer:tick")
2820 .route_id("test-route")
2821 .choice()
2822 .when(|ex: &Exchange| ex.input.header("type").is_some())
2823 .to("mock:typed")
2824 .end_when()
2825 .end_choice()
2826 .build()
2827 .unwrap();
2828 assert_eq!(definition.steps().len(), 1);
2829 assert!(
2830 matches!(&definition.steps()[0], BuilderStep::Choice { whens, otherwise }
2831 if whens.len() == 1 && otherwise.is_none())
2832 );
2833 }
2834
2835 #[test]
2836 fn test_choice_builder_when_otherwise() {
2837 let definition = RouteBuilder::from("timer:tick")
2838 .route_id("test-route")
2839 .choice()
2840 .when(|ex: &Exchange| ex.input.header("a").is_some())
2841 .to("mock:a")
2842 .end_when()
2843 .otherwise()
2844 .to("mock:fallback")
2845 .end_otherwise()
2846 .end_choice()
2847 .build()
2848 .unwrap();
2849 assert!(
2850 matches!(&definition.steps()[0], BuilderStep::Choice { whens, otherwise }
2851 if whens.len() == 1 && otherwise.is_some())
2852 );
2853 }
2854
2855 #[test]
2856 fn test_choice_builder_multiple_whens() {
2857 let definition = RouteBuilder::from("timer:tick")
2858 .route_id("test-route")
2859 .choice()
2860 .when(|ex: &Exchange| ex.input.header("a").is_some())
2861 .to("mock:a")
2862 .end_when()
2863 .when(|ex: &Exchange| ex.input.header("b").is_some())
2864 .to("mock:b")
2865 .end_when()
2866 .end_choice()
2867 .build()
2868 .unwrap();
2869 assert!(
2870 matches!(&definition.steps()[0], BuilderStep::Choice { whens, .. }
2871 if whens.len() == 2)
2872 );
2873 }
2874
2875 #[test]
2876 fn test_choice_step_after_choice() {
2877 let definition = RouteBuilder::from("timer:tick")
2879 .route_id("test-route")
2880 .choice()
2881 .when(|_ex: &Exchange| true)
2882 .to("mock:inner")
2883 .end_when()
2884 .end_choice()
2885 .to("mock:outer") .build()
2887 .unwrap();
2888 assert_eq!(definition.steps().len(), 2);
2889 assert!(matches!(&definition.steps()[1], BuilderStep::To(uri) if uri == "mock:outer"));
2890 }
2891
2892 #[test]
2895 fn test_throttle_builder_typestate() {
2896 let definition = RouteBuilder::from("timer:tick")
2897 .route_id("test-route")
2898 .throttle(10, std::time::Duration::from_secs(1))
2899 .to("mock:result")
2900 .end_throttle()
2901 .build()
2902 .unwrap();
2903
2904 assert_eq!(definition.steps().len(), 1);
2905 assert!(matches!(
2906 &definition.steps()[0],
2907 BuilderStep::Throttle { .. }
2908 ));
2909 }
2910
2911 #[test]
2912 fn test_throttle_builder_with_strategy() {
2913 let definition = RouteBuilder::from("timer:tick")
2914 .route_id("test-route")
2915 .throttle(10, std::time::Duration::from_secs(1))
2916 .strategy(ThrottleStrategy::Reject)
2917 .to("mock:result")
2918 .end_throttle()
2919 .build()
2920 .unwrap();
2921
2922 if let BuilderStep::Throttle { config, .. } = &definition.steps()[0] {
2923 assert_eq!(config.strategy, ThrottleStrategy::Reject);
2924 } else {
2925 panic!("Expected Throttle step");
2926 }
2927 }
2928
2929 #[test]
2930 fn test_throttle_builder_steps_collected() {
2931 let definition = RouteBuilder::from("timer:tick")
2932 .route_id("test-route")
2933 .throttle(5, std::time::Duration::from_secs(1))
2934 .set_header("throttled", Value::Bool(true))
2935 .to("mock:throttled")
2936 .end_throttle()
2937 .build()
2938 .unwrap();
2939
2940 match &definition.steps()[0] {
2941 BuilderStep::Throttle { steps, .. } => {
2942 assert_eq!(steps.len(), 2); }
2944 other => panic!("Expected Throttle, got {:?}", other),
2945 }
2946 }
2947
2948 #[test]
2949 fn test_throttle_step_after_throttle() {
2950 let definition = RouteBuilder::from("timer:tick")
2952 .route_id("test-route")
2953 .throttle(10, std::time::Duration::from_secs(1))
2954 .to("mock:inner")
2955 .end_throttle()
2956 .to("mock:outer")
2957 .build()
2958 .unwrap();
2959
2960 assert_eq!(definition.steps().len(), 2);
2961 assert!(matches!(&definition.steps()[1], BuilderStep::To(uri) if uri == "mock:outer"));
2962 }
2963
2964 #[test]
2967 fn test_load_balance_builder_typestate() {
2968 let definition = RouteBuilder::from("timer:tick")
2969 .route_id("test-route")
2970 .load_balance()
2971 .round_robin()
2972 .to("mock:a")
2973 .to("mock:b")
2974 .end_load_balance()
2975 .build()
2976 .unwrap();
2977
2978 assert_eq!(definition.steps().len(), 1);
2979 assert!(matches!(
2980 &definition.steps()[0],
2981 BuilderStep::LoadBalance { .. }
2982 ));
2983 }
2984
2985 #[test]
2986 fn test_load_balance_builder_with_strategy() {
2987 let definition = RouteBuilder::from("timer:tick")
2988 .route_id("test-route")
2989 .load_balance()
2990 .random()
2991 .to("mock:result")
2992 .end_load_balance()
2993 .build()
2994 .unwrap();
2995
2996 if let BuilderStep::LoadBalance { config, .. } = &definition.steps()[0] {
2997 assert_eq!(config.strategy, LoadBalanceStrategy::Random);
2998 } else {
2999 panic!("Expected LoadBalance step");
3000 }
3001 }
3002
3003 #[test]
3004 fn test_load_balance_builder_steps_collected() {
3005 let definition = RouteBuilder::from("timer:tick")
3006 .route_id("test-route")
3007 .load_balance()
3008 .set_header("lb", Value::Bool(true))
3009 .to("mock:a")
3010 .end_load_balance()
3011 .build()
3012 .unwrap();
3013
3014 match &definition.steps()[0] {
3015 BuilderStep::LoadBalance { steps, .. } => {
3016 assert_eq!(steps.len(), 2); }
3018 other => panic!("Expected LoadBalance, got {:?}", other),
3019 }
3020 }
3021
3022 #[test]
3023 fn test_load_balance_step_after_load_balance() {
3024 let definition = RouteBuilder::from("timer:tick")
3026 .route_id("test-route")
3027 .load_balance()
3028 .to("mock:inner")
3029 .end_load_balance()
3030 .to("mock:outer")
3031 .build()
3032 .unwrap();
3033
3034 assert_eq!(definition.steps().len(), 2);
3035 assert!(matches!(&definition.steps()[1], BuilderStep::To(uri) if uri == "mock:outer"));
3036 }
3037
3038 #[test]
3041 fn test_dynamic_router_builder() {
3042 let definition = RouteBuilder::from("timer:tick")
3043 .route_id("test-route")
3044 .dynamic_router(Arc::new(|_| Some("mock:result".to_string())))
3045 .build()
3046 .unwrap();
3047
3048 assert_eq!(definition.steps().len(), 1);
3049 assert!(matches!(
3050 &definition.steps()[0],
3051 BuilderStep::DynamicRouter { .. }
3052 ));
3053 }
3054
3055 #[test]
3056 fn test_dynamic_router_builder_with_config() {
3057 let config = DynamicRouterConfig::new(Arc::new(|_| Some("mock:a".to_string())))
3058 .max_iterations(100)
3059 .cache_size(500);
3060
3061 let definition = RouteBuilder::from("timer:tick")
3062 .route_id("test-route")
3063 .dynamic_router_with_config(config)
3064 .build()
3065 .unwrap();
3066
3067 assert_eq!(definition.steps().len(), 1);
3068 if let BuilderStep::DynamicRouter { config } = &definition.steps()[0] {
3069 assert_eq!(config.max_iterations, 100);
3070 assert_eq!(config.cache_size, 500);
3071 } else {
3072 panic!("Expected DynamicRouter step");
3073 }
3074 }
3075
3076 #[test]
3077 fn test_dynamic_router_step_after_router() {
3078 let definition = RouteBuilder::from("timer:tick")
3080 .route_id("test-route")
3081 .dynamic_router(Arc::new(|_| Some("mock:inner".to_string())))
3082 .to("mock:outer")
3083 .build()
3084 .unwrap();
3085
3086 assert_eq!(definition.steps().len(), 2);
3087 assert!(matches!(
3088 &definition.steps()[0],
3089 BuilderStep::DynamicRouter { .. }
3090 ));
3091 assert!(matches!(&definition.steps()[1], BuilderStep::To(uri) if uri == "mock:outer"));
3092 }
3093
3094 #[test]
3095 fn routing_slip_builder_creates_step() {
3096 use camel_api::RoutingSlipExpression;
3097
3098 let expression: RoutingSlipExpression = Arc::new(|_| Some("direct:a,direct:b".to_string()));
3099
3100 let route = RouteBuilder::from("direct:start")
3101 .route_id("routing-slip-test")
3102 .routing_slip(expression)
3103 .build()
3104 .unwrap();
3105
3106 assert!(
3107 matches!(route.steps()[0], BuilderStep::RoutingSlip { .. }),
3108 "Expected RoutingSlip step"
3109 );
3110 }
3111
3112 #[test]
3113 fn routing_slip_with_config_builder_creates_step() {
3114 use camel_api::RoutingSlipConfig;
3115
3116 let config = RoutingSlipConfig::new(Arc::new(|_| Some("mock:a".to_string())))
3117 .uri_delimiter("|")
3118 .cache_size(50)
3119 .ignore_invalid_endpoints(true);
3120
3121 let route = RouteBuilder::from("direct:start")
3122 .route_id("routing-slip-config-test")
3123 .routing_slip_with_config(config)
3124 .build()
3125 .unwrap();
3126
3127 if let BuilderStep::RoutingSlip { config } = &route.steps()[0] {
3128 assert_eq!(config.uri_delimiter, "|");
3129 assert_eq!(config.cache_size, 50);
3130 assert!(config.ignore_invalid_endpoints);
3131 } else {
3132 panic!("Expected RoutingSlip step");
3133 }
3134 }
3135
3136 #[test]
3137 fn test_builder_marshal_adds_processor_step() {
3138 let definition = RouteBuilder::from("timer:tick")
3139 .route_id("test-route")
3140 .marshal("json")
3141 .unwrap()
3142 .build()
3143 .unwrap();
3144 assert!(matches!(&definition.steps()[0], BuilderStep::Processor(_)));
3145 }
3146
3147 #[test]
3148 fn test_builder_unmarshal_adds_processor_step() {
3149 let definition = RouteBuilder::from("timer:tick")
3150 .route_id("test-route")
3151 .unmarshal("json")
3152 .unwrap()
3153 .build()
3154 .unwrap();
3155 assert!(matches!(&definition.steps()[0], BuilderStep::Processor(_)));
3156 }
3157
3158 #[test]
3159 fn test_builder_stream_cache_adds_processor_step() {
3160 let definition = RouteBuilder::from("timer:tick")
3161 .route_id("test-route")
3162 .stream_cache(1024)
3163 .build()
3164 .unwrap();
3165 assert!(matches!(&definition.steps()[0], BuilderStep::Processor(_)));
3166 }
3167
3168 #[test]
3169 fn validate_adds_validate_step() {
3170 let def = RouteBuilder::from("direct:in")
3171 .route_id("test")
3172 .validate("schemas/order.xsd")
3173 .build()
3174 .unwrap();
3175 let steps = def.steps();
3176 assert_eq!(steps.len(), 1);
3177 assert!(
3178 matches!(&steps[0], BuilderStep::Validate { predicate } if predicate.language == "simple" && predicate.source == "schemas/order.xsd"),
3179 "got: {:?}",
3180 steps[0]
3181 );
3182 }
3183
3184 #[test]
3185 fn test_builder_marshal_returns_err_for_unknown_format() {
3186 let result = RouteBuilder::from("timer:tick")
3187 .route_id("test-route")
3188 .marshal("protobuf");
3189 let err = match result {
3190 Err(e) => e,
3191 Ok(_) => panic!("marshal with unknown format should return Err"),
3192 };
3193 let msg = err.to_string();
3194 assert!(
3195 msg.contains("unknown data format"),
3196 "error should mention unknown format, got: {msg}"
3197 );
3198 assert!(
3199 msg.contains("protobuf"),
3200 "error should mention format name, got: {msg}"
3201 );
3202 }
3203
3204 #[test]
3205 fn test_builder_unmarshal_returns_err_for_unknown_format() {
3206 let result = RouteBuilder::from("timer:tick")
3207 .route_id("test-route")
3208 .unmarshal("protobuf");
3209 let err = match result {
3210 Err(e) => e,
3211 Ok(_) => panic!("unmarshal with unknown format should return Err"),
3212 };
3213 let msg = err.to_string();
3214 assert!(
3215 msg.contains("unknown data format"),
3216 "error should mention unknown format, got: {msg}"
3217 );
3218 assert!(
3219 msg.contains("protobuf"),
3220 "error should mention format name, got: {msg}"
3221 );
3222 }
3223
3224 #[test]
3225 fn test_builder_recipient_list_creates_step() {
3226 let route = RouteBuilder::from("direct:start")
3227 .route_id("recipient-list-test")
3228 .recipient_list(Arc::new(|_| "direct:a,direct:b".to_string()))
3229 .build()
3230 .unwrap();
3231
3232 assert!(matches!(
3233 &route.steps()[0],
3234 BuilderStep::RecipientList { .. }
3235 ));
3236 }
3237
3238 #[test]
3239 fn test_builder_recipient_list_with_config_creates_step() {
3240 let config = RecipientListConfig::new(Arc::new(|_| "mock:a".to_string()));
3241
3242 let route = RouteBuilder::from("direct:start")
3243 .route_id("recipient-list-config-test")
3244 .recipient_list_with_config(config)
3245 .build()
3246 .unwrap();
3247
3248 assert!(matches!(
3249 &route.steps()[0],
3250 BuilderStep::RecipientList { .. }
3251 ));
3252 }
3253
3254 #[test]
3255 fn test_builder_script_adds_script_step() {
3256 let route = RouteBuilder::from("direct:start")
3257 .route_id("script-test")
3258 .script("rhai", "headers[\"x\"] = \"y\"")
3259 .build()
3260 .unwrap();
3261
3262 assert!(matches!(
3263 &route.steps()[0],
3264 BuilderStep::Script { language, script }
3265 if language == "rhai" && script == "headers[\"x\"] = \"y\""
3266 ));
3267 }
3268
3269 #[test]
3270 fn test_builder_delay_and_delay_with_header_add_steps() {
3271 let route = RouteBuilder::from("direct:start")
3272 .route_id("delay-test")
3273 .delay(Duration::from_millis(250))
3274 .delay_with_header(Duration::from_millis(500), "x-delay")
3275 .build()
3276 .unwrap();
3277
3278 assert_eq!(route.steps().len(), 2);
3279 assert!(matches!(&route.steps()[0], BuilderStep::Delay { .. }));
3280 assert!(matches!(&route.steps()[1], BuilderStep::Delay { .. }));
3281 }
3282
3283 #[test]
3284 fn test_builder_log_and_stop_add_steps_in_order() {
3285 let route = RouteBuilder::from("direct:start")
3286 .route_id("log-stop-test")
3287 .log("hello", LogLevel::Info)
3288 .stop()
3289 .to("mock:after")
3290 .build()
3291 .unwrap();
3292
3293 assert_eq!(route.steps().len(), 3);
3294 assert!(matches!(
3295 &route.steps()[0],
3296 BuilderStep::Log { message, .. } if message == "hello"
3297 ));
3298 assert!(matches!(&route.steps()[1], BuilderStep::Stop));
3299 assert!(matches!(&route.steps()[2], BuilderStep::To(uri) if uri == "mock:after"));
3300 }
3301
3302 #[test]
3303 fn test_builder_stream_cache_default_adds_processor_step() {
3304 let route = RouteBuilder::from("direct:start")
3305 .route_id("stream-cache-default-test")
3306 .stream_cache_default()
3307 .build()
3308 .unwrap();
3309
3310 assert!(matches!(&route.steps()[0], BuilderStep::Processor(_)));
3311 }
3312
3313 #[test]
3314 fn test_validate_creates_validate_step_with_expression() {
3315 let route = RouteBuilder::from("direct:in")
3316 .route_id("validate-prefix-test")
3317 .validate("${body.size()} > 0")
3318 .build()
3319 .unwrap();
3320
3321 assert!(matches!(
3322 &route.steps()[0],
3323 BuilderStep::Validate { predicate } if predicate.language == "simple" && predicate.source == "${body.size()} > 0"
3324 ));
3325 }
3326
3327 #[test]
3328 fn test_load_balance_builder_weighted_failover_config() {
3329 let route = RouteBuilder::from("direct:start")
3330 .route_id("lb-weighted-failover")
3331 .load_balance()
3332 .weighted(vec![
3333 ("direct:a".to_string(), 3),
3334 ("direct:b".to_string(), 1),
3335 ])
3336 .failover()
3337 .to("mock:result")
3338 .end_load_balance()
3339 .build()
3340 .unwrap();
3341
3342 if let BuilderStep::LoadBalance { config, .. } = &route.steps()[0] {
3343 assert_eq!(config.strategy, LoadBalanceStrategy::Failover);
3344 } else {
3345 panic!("Expected LoadBalance step");
3346 }
3347 }
3348
3349 #[test]
3350 fn test_multicast_builder_all_config_setters() {
3351 let route = RouteBuilder::from("direct:start")
3352 .route_id("multicast-config-test")
3353 .multicast()
3354 .parallel(true)
3355 .parallel_limit(4)
3356 .stop_on_exception(true)
3357 .timeout(Duration::from_millis(300))
3358 .aggregation(MulticastStrategy::Original)
3359 .to("mock:a")
3360 .end_multicast()
3361 .build()
3362 .unwrap();
3363
3364 if let BuilderStep::Multicast { config, .. } = &route.steps()[0] {
3365 assert!(config.parallel);
3366 assert_eq!(config.parallel_limit, Some(4));
3367 assert!(config.stop_on_exception);
3368 assert_eq!(config.timeout, Some(Duration::from_millis(300)));
3369 assert!(matches!(config.aggregation, MulticastStrategy::Original));
3370 } else {
3371 panic!("Expected Multicast step");
3372 }
3373 }
3374
3375 #[test]
3376 fn test_build_canonical_rejects_unsupported_processor_step() {
3377 let err = RouteBuilder::from("direct:start")
3378 .route_id("canonical-reject")
3379 .set_header("k", Value::String("v".into()))
3380 .build_canonical()
3381 .unwrap_err();
3382
3383 assert!(format!("{err}").contains("does not support step `processor`"));
3384 }
3385
3386 #[test]
3389 fn test_load_balance_builder_weighted_strategy() {
3390 let route = RouteBuilder::from("direct:start")
3391 .route_id("lb-weighted")
3392 .load_balance()
3393 .weighted(vec![
3394 ("direct:a".to_string(), 5),
3395 ("direct:b".to_string(), 2),
3396 ("direct:c".to_string(), 1),
3397 ])
3398 .to("mock:result")
3399 .end_load_balance()
3400 .build()
3401 .unwrap();
3402
3403 if let BuilderStep::LoadBalance { config, .. } = &route.steps()[0] {
3404 assert!(matches!(config.strategy, LoadBalanceStrategy::Weighted(_)));
3405 } else {
3406 panic!("Expected LoadBalance step");
3407 }
3408 }
3409
3410 #[test]
3411 fn test_load_balance_builder_failover_strategy() {
3412 let route = RouteBuilder::from("direct:start")
3413 .route_id("lb-failover")
3414 .load_balance()
3415 .failover()
3416 .to("mock:primary")
3417 .end_load_balance()
3418 .build()
3419 .unwrap();
3420
3421 if let BuilderStep::LoadBalance { config, .. } = &route.steps()[0] {
3422 assert_eq!(config.strategy, LoadBalanceStrategy::Failover);
3423 } else {
3424 panic!("Expected LoadBalance step");
3425 }
3426 }
3427
3428 #[test]
3431 fn test_filter_in_split_builder_typestate() {
3432 use camel_api::splitter::{SplitterConfig, split_body_lines};
3433
3434 let definition = RouteBuilder::from("timer:test")
3435 .route_id("filter-in-split")
3436 .split(SplitterConfig::new(split_body_lines()))
3437 .filter(|_ex| true)
3438 .to("mock:filtered")
3439 .end_filter()
3440 .end_split()
3441 .build()
3442 .unwrap();
3443
3444 assert_eq!(definition.steps().len(), 1);
3445 if let BuilderStep::Split { steps, .. } = &definition.steps()[0] {
3446 assert_eq!(steps.len(), 1);
3447 assert!(matches!(&steps[0], BuilderStep::Filter { .. }));
3448 } else {
3449 panic!("Expected Split step");
3450 }
3451 }
3452
3453 #[test]
3454 fn test_filter_in_split_builder_multiple_steps() {
3455 use camel_api::splitter::{SplitterConfig, split_body_lines};
3456
3457 let definition = RouteBuilder::from("timer:test")
3458 .route_id("filter-in-split-multi")
3459 .split(SplitterConfig::new(split_body_lines()))
3460 .to("mock:before-filter")
3461 .filter(|_ex| true)
3462 .to("mock:inside-filter")
3463 .end_filter()
3464 .to("mock:after-filter")
3465 .end_split()
3466 .build()
3467 .unwrap();
3468
3469 if let BuilderStep::Split { steps, .. } = &definition.steps()[0] {
3470 assert_eq!(steps.len(), 3);
3472 } else {
3473 panic!("Expected Split step");
3474 }
3475 }
3476
3477 #[test]
3480 fn test_build_canonical_with_circuit_breaker() {
3481 use camel_api::circuit_breaker::CircuitBreakerConfig;
3482
3483 let spec = RouteBuilder::from("direct:start")
3484 .route_id("canonical-cb")
3485 .circuit_breaker(CircuitBreakerConfig::new().failure_threshold(10))
3486 .to("mock:result")
3487 .build_canonical()
3488 .unwrap();
3489
3490 let cb = spec.circuit_breaker.expect("circuit breaker should be set");
3491 assert_eq!(cb.failure_threshold, 10);
3492 }
3493
3494 #[test]
3495 fn test_build_canonical_rejects_custom_split_aggregation() {
3496 use camel_api::splitter::{SplitterConfig, split_body_lines};
3497
3498 let err = RouteBuilder::from("direct:start")
3499 .route_id("canonical-custom-split")
3500 .split(SplitterConfig::new(split_body_lines()).aggregation(
3501 camel_api::splitter::AggregationStrategy::Custom(Arc::new(|_, ex| ex)),
3502 ))
3503 .to("mock:frag")
3504 .end_split()
3505 .build_canonical()
3506 .unwrap_err();
3507
3508 assert!(format!("{err}").contains("canonical v2 does not support step `split`"));
3510 }
3511
3512 #[test]
3513 fn test_build_canonical_rejects_custom_aggregate_strategy() {
3514 let err = RouteBuilder::from("direct:start")
3515 .route_id("canonical-custom-agg")
3516 .aggregate(
3517 AggregatorConfig::correlate_by("key")
3518 .complete_when_size(2)
3519 .strategy(AggregationStrategy::Custom(Arc::new(|_, ex| ex)))
3520 .build()
3521 .unwrap(),
3522 )
3523 .build_canonical()
3524 .unwrap_err();
3525
3526 assert!(format!("{err}").contains("custom aggregate strategy"));
3527 }
3528
3529 #[test]
3530 fn test_build_canonical_rejects_fn_correlation_strategy() {
3531 let err = RouteBuilder::from("direct:start")
3532 .route_id("canonical-fn-corr")
3533 .aggregate(AggregatorConfig {
3534 header_name: "key".to_string(),
3535 completion: CompletionMode::Single(CompletionCondition::Size(1)),
3536 correlation: CorrelationStrategy::Fn(Arc::new(|_| Some("key".to_string()))),
3537 strategy: AggregationStrategy::CollectAll,
3538 max_buckets: None,
3539 max_bucket_size: None,
3540 bucket_ttl: None,
3541 force_completion_on_stop: false,
3542 discard_on_timeout: false,
3543 max_timeout_tasks: 1024,
3544 })
3545 .build_canonical()
3546 .unwrap_err();
3547
3548 assert!(format!("{err}").contains("Fn correlation strategy"));
3549 }
3550
3551 #[test]
3552 fn test_build_canonical_rejects_predicate_completion() {
3553 let err = RouteBuilder::from("direct:start")
3554 .route_id("canonical-pred-completion")
3555 .aggregate(AggregatorConfig {
3556 header_name: "key".to_string(),
3557 completion: CompletionMode::Single(CompletionCondition::Predicate(Arc::new(
3558 |_| false,
3559 ))),
3560 correlation: CorrelationStrategy::HeaderName("key".to_string()),
3561 strategy: AggregationStrategy::CollectAll,
3562 max_buckets: None,
3563 max_bucket_size: None,
3564 bucket_ttl: None,
3565 force_completion_on_stop: false,
3566 discard_on_timeout: false,
3567 max_timeout_tasks: 1024,
3568 })
3569 .build_canonical()
3570 .unwrap_err();
3571
3572 assert!(
3573 format!("{err}").contains("cannot reverse-map"),
3574 "reject message must explain forward-only: {}",
3575 err
3576 );
3577 }
3578
3579 #[test]
3580 fn extract_completion_fields_rejects_predicate_expr() {
3581 let mode = CompletionMode::Single(CompletionCondition::PredicateExpr {
3582 expr: "${body} == 'DONE'".to_string(),
3583 language: "simple".to_string(),
3584 });
3585 let result = extract_completion_fields(&mode);
3586 assert!(
3587 result.is_err(),
3588 "PredicateExpr must be rejected (forward-only)"
3589 );
3590 let msg = format!("{}", result.unwrap_err());
3591 assert!(
3592 msg.contains("cannot reverse-map"),
3593 "reject message must explain forward-only: {}",
3594 msg
3595 );
3596 }
3597
3598 #[test]
3599 fn extract_completion_fields_rejects_predicate_expr_any_mode() {
3600 let mode = CompletionMode::Any(vec![
3601 CompletionCondition::Size(5),
3602 CompletionCondition::PredicateExpr {
3603 expr: "${body} == 'DONE'".to_string(),
3604 language: "simple".to_string(),
3605 },
3606 ]);
3607 let result = extract_completion_fields(&mode);
3608 assert!(
3609 result.is_err(),
3610 "PredicateExpr in Any must be rejected (forward-only)"
3611 );
3612 let msg = format!("{}", result.unwrap_err());
3613 assert!(
3614 msg.contains("cannot reverse-map"),
3615 "reject message must explain forward-only: {}",
3616 msg
3617 );
3618 }
3619
3620 #[test]
3621 fn test_build_canonical_with_expression_correlation() {
3622 let spec = RouteBuilder::from("direct:start")
3623 .route_id("canonical-expr-corr")
3624 .aggregate(AggregatorConfig {
3625 header_name: "key".to_string(),
3626 completion: CompletionMode::Single(CompletionCondition::Size(1)),
3627 correlation: CorrelationStrategy::Expression {
3628 expr: "header.key".to_string(),
3629 language: "simple".to_string(),
3630 },
3631 strategy: AggregationStrategy::CollectAll,
3632 max_buckets: None,
3633 max_bucket_size: None,
3634 bucket_ttl: None,
3635 force_completion_on_stop: false,
3636 discard_on_timeout: false,
3637 max_timeout_tasks: 1024,
3638 })
3639 .build_canonical()
3640 .unwrap();
3641
3642 assert!(spec.steps.iter().any(|s| matches!(s, CanonicalStepSpec::Aggregate(a) if a.correlation_key == Some("header.key".to_string()))));
3643 }
3644
3645 #[test]
3646 fn test_build_canonical_split_rejected_with_closure_expression() {
3647 use camel_api::splitter::{AggregationStrategy, SplitterConfig, split_body_lines};
3648
3649 let err = RouteBuilder::from("direct:start")
3651 .route_id("canonical-split-last")
3652 .split(
3653 SplitterConfig::new(split_body_lines()).aggregation(AggregationStrategy::LastWins),
3654 )
3655 .to("mock:frag")
3656 .end_split()
3657 .build_canonical()
3658 .unwrap_err();
3659
3660 assert!(format!("{err}").contains("canonical v2 does not support step `split`"));
3661 }
3662
3663 #[test]
3666 fn test_on_exception_full_chain_retry_backoff_jitter_handled_by() {
3667 let definition = RouteBuilder::from("direct:start")
3668 .route_id("on-exception-full")
3669 .dead_letter_channel("log:dlc")
3670 .on_exception(|e| matches!(e, CamelError::Io(_)))
3671 .retry(5)
3672 .with_backoff(Duration::from_millis(10), 2.0, Duration::from_millis(500))
3673 .with_jitter(0.3)
3674 .handled_by("log:io-handler")
3675 .end_on_exception()
3676 .to("mock:out")
3677 .build()
3678 .unwrap();
3679
3680 let cfg = definition
3681 .error_handler_config()
3682 .expect("error handler should be set");
3683 assert_eq!(cfg.policies.len(), 1);
3684 let policy = &cfg.policies[0];
3685 let retry = policy.retry.as_ref().expect("retry should be set");
3686 assert_eq!(retry.max_attempts, 5);
3687 assert_eq!(retry.initial_delay, Duration::from_millis(10));
3688 assert_eq!(retry.multiplier, 2.0);
3689 assert_eq!(retry.max_delay, Duration::from_millis(500));
3690 assert!((retry.jitter_factor - 0.3).abs() < f64::EPSILON);
3691 assert_eq!(policy.handled_by.as_deref(), Some("log:io-handler"));
3692 }
3693
3694 #[test]
3695 fn test_on_exception_jitter_clamped_to_valid_range() {
3696 let definition = RouteBuilder::from("direct:start")
3697 .route_id("jitter-clamp")
3698 .on_exception(|_e| true)
3699 .retry(1)
3700 .with_jitter(5.0)
3701 .end_on_exception()
3702 .to("mock:out")
3703 .build()
3704 .unwrap();
3705
3706 let cfg = definition.error_handler_config().unwrap();
3707 let retry = cfg.policies[0].retry.as_ref().unwrap();
3708 assert!((retry.jitter_factor - 1.0).abs() < f64::EPSILON);
3709 }
3710
3711 #[test]
3714 fn test_builder_process_fn_adds_processor_step() {
3715 use camel_api::BoxProcessorExt;
3716 let processor = BoxProcessor::from_fn(|ex| Box::pin(async move { Ok(ex) }));
3717 let definition = RouteBuilder::from("timer:tick")
3718 .route_id("process-fn-test")
3719 .process_fn(processor)
3720 .build()
3721 .unwrap();
3722
3723 assert!(matches!(&definition.steps()[0], BuilderStep::Processor(_)));
3724 }
3725
3726 #[test]
3727 fn test_builder_convert_body_to_adds_processor_step() {
3728 let definition = RouteBuilder::from("timer:tick")
3729 .route_id("convert-body-test")
3730 .convert_body_to(BodyType::Json)
3731 .build()
3732 .unwrap();
3733
3734 assert!(matches!(&definition.steps()[0], BuilderStep::Processor(_)));
3735 }
3736
3737 #[test]
3738 fn test_builder_bean_adds_bean_step() {
3739 let definition = RouteBuilder::from("timer:tick")
3740 .route_id("bean-test")
3741 .bean("myBean", "process")
3742 .build()
3743 .unwrap();
3744
3745 assert!(
3746 matches!(&definition.steps()[0], BuilderStep::Bean { name, method }
3747 if name == "myBean" && method == "process")
3748 );
3749 }
3750
3751 #[test]
3754 fn test_throttle_builder_delay_strategy() {
3755 let definition = RouteBuilder::from("timer:tick")
3756 .route_id("throttle-delay")
3757 .throttle(10, Duration::from_secs(1))
3758 .strategy(ThrottleStrategy::Delay)
3759 .to("mock:result")
3760 .end_throttle()
3761 .build()
3762 .unwrap();
3763
3764 if let BuilderStep::Throttle { config, .. } = &definition.steps()[0] {
3765 assert_eq!(config.strategy, ThrottleStrategy::Delay);
3766 } else {
3767 panic!("Expected Throttle step");
3768 }
3769 }
3770
3771 #[test]
3772 fn test_throttle_builder_drop_strategy() {
3773 let definition = RouteBuilder::from("timer:tick")
3774 .route_id("throttle-drop")
3775 .throttle(10, Duration::from_secs(1))
3776 .strategy(ThrottleStrategy::Drop)
3777 .to("mock:result")
3778 .end_throttle()
3779 .build()
3780 .unwrap();
3781
3782 if let BuilderStep::Throttle { config, .. } = &definition.steps()[0] {
3783 assert_eq!(config.strategy, ThrottleStrategy::Drop);
3784 } else {
3785 panic!("Expected Throttle step");
3786 }
3787 }
3788
3789 #[test]
3792 fn test_nested_loop_while_builder() {
3793 use camel_api::loop_eip::LoopMode;
3794
3795 let def = RouteBuilder::from("direct:start")
3796 .route_id("nested-loop-while")
3797 .loop_count(2)
3798 .to("mock:outer")
3799 .loop_while(|_ex| true)
3800 .to("mock:inner")
3801 .end_loop()
3802 .end_loop()
3803 .build()
3804 .unwrap();
3805
3806 assert_eq!(def.steps().len(), 1);
3807 if let BuilderStep::Loop { steps, .. } = &def.steps()[0] {
3808 assert_eq!(steps.len(), 2);
3809 if let BuilderStep::Loop { config, .. } = &steps[1] {
3810 assert!(matches!(config.mode, LoopMode::While(_)));
3811 } else {
3812 panic!("Expected inner Loop step");
3813 }
3814 } else {
3815 panic!("Expected outer Loop step");
3816 }
3817 }
3818
3819 #[test]
3822 fn test_choice_builder_multiple_whens_with_otherwise() {
3823 let definition = RouteBuilder::from("timer:tick")
3824 .route_id("choice-multi-otherwise")
3825 .choice()
3826 .when(|ex: &Exchange| ex.input.header("a").is_some())
3827 .to("mock:a")
3828 .end_when()
3829 .when(|ex: &Exchange| ex.input.header("b").is_some())
3830 .to("mock:b")
3831 .end_when()
3832 .when(|ex: &Exchange| ex.input.header("c").is_some())
3833 .to("mock:c")
3834 .end_when()
3835 .otherwise()
3836 .to("mock:fallback")
3837 .end_otherwise()
3838 .end_choice()
3839 .build()
3840 .unwrap();
3841
3842 if let BuilderStep::Choice { whens, otherwise } = &definition.steps()[0] {
3843 assert_eq!(whens.len(), 3);
3844 assert!(otherwise.is_some());
3845 assert_eq!(otherwise.as_ref().unwrap().len(), 1);
3846 } else {
3847 panic!("Expected Choice step");
3848 }
3849 }
3850
3851 #[test]
3854 fn test_multicast_builder_parallel_only() {
3855 let route = RouteBuilder::from("direct:start")
3856 .route_id("multicast-parallel")
3857 .multicast()
3858 .parallel(true)
3859 .to("mock:a")
3860 .end_multicast()
3861 .build()
3862 .unwrap();
3863
3864 if let BuilderStep::Multicast { config, .. } = &route.steps()[0] {
3865 assert!(config.parallel);
3866 assert_eq!(config.parallel_limit, None);
3867 } else {
3868 panic!("Expected Multicast step");
3869 }
3870 }
3871
3872 #[test]
3873 fn test_multicast_builder_timeout_only() {
3874 let route = RouteBuilder::from("direct:start")
3875 .route_id("multicast-timeout")
3876 .multicast()
3877 .timeout(Duration::from_secs(5))
3878 .to("mock:a")
3879 .end_multicast()
3880 .build()
3881 .unwrap();
3882
3883 if let BuilderStep::Multicast { config, .. } = &route.steps()[0] {
3884 assert_eq!(config.timeout, Some(Duration::from_secs(5)));
3885 } else {
3886 panic!("Expected Multicast step");
3887 }
3888 }
3889
3890 #[test]
3891 fn test_multicast_builder_aggregation_collect_all() {
3892 let route = RouteBuilder::from("direct:start")
3893 .route_id("multicast-collect")
3894 .multicast()
3895 .aggregation(MulticastStrategy::CollectAll)
3896 .to("mock:a")
3897 .end_multicast()
3898 .build()
3899 .unwrap();
3900
3901 if let BuilderStep::Multicast { config, .. } = &route.steps()[0] {
3902 assert!(matches!(config.aggregation, MulticastStrategy::CollectAll));
3903 } else {
3904 panic!("Expected Multicast step");
3905 }
3906 }
3907
3908 #[test]
3911 fn test_build_canonical_aggregate_any_completion_mode() {
3912 let spec = RouteBuilder::from("direct:start")
3913 .route_id("canonical-any-completion")
3914 .aggregate(
3915 AggregatorConfig::correlate_by("key")
3916 .complete_on_size_or_timeout(10, Duration::from_secs(30))
3917 .build()
3918 .unwrap(),
3919 )
3920 .build_canonical()
3921 .unwrap();
3922
3923 if let CanonicalStepSpec::Aggregate(agg) = &spec.steps[0] {
3924 assert_eq!(agg.completion_size, Some(10));
3925 assert_eq!(agg.completion_timeout_ms, Some(30_000));
3926 } else {
3927 panic!("Expected Aggregate step");
3928 }
3929 }
3930
3931 #[test]
3932 fn test_build_canonical_aggregate_timeout_completion() {
3933 let spec = RouteBuilder::from("direct:start")
3934 .route_id("canonical-timeout-completion")
3935 .aggregate(
3936 AggregatorConfig::correlate_by("key")
3937 .complete_on_timeout(Duration::from_millis(500))
3938 .build()
3939 .unwrap(),
3940 )
3941 .build_canonical()
3942 .unwrap();
3943
3944 if let CanonicalStepSpec::Aggregate(agg) = &spec.steps[0] {
3945 assert_eq!(agg.completion_size, None);
3946 assert_eq!(agg.completion_timeout_ms, Some(500));
3947 } else {
3948 panic!("Expected Aggregate step");
3949 }
3950 }
3951
3952 #[test]
3955 fn test_build_canonical_aggregate_discard_on_timeout() {
3956 use camel_api::aggregator::AggregatorConfig;
3957
3958 let spec = RouteBuilder::from("direct:start")
3959 .route_id("canonical-discard-timeout")
3960 .aggregate(
3961 AggregatorConfig::correlate_by("key")
3962 .complete_when_size(1)
3963 .discard_on_timeout(true)
3964 .build()
3965 .unwrap(),
3966 )
3967 .build_canonical()
3968 .unwrap();
3969
3970 if let CanonicalStepSpec::Aggregate(agg) = &spec.steps[0] {
3971 assert_eq!(agg.discard_on_timeout, Some(true));
3972 } else {
3973 panic!("Expected Aggregate step");
3974 }
3975 }
3976
3977 #[test]
3978 fn test_build_canonical_aggregate_force_completion_on_stop() {
3979 use camel_api::aggregator::AggregatorConfig;
3980
3981 let spec = RouteBuilder::from("direct:start")
3982 .route_id("canonical-force-stop")
3983 .aggregate(
3984 AggregatorConfig::correlate_by("key")
3985 .complete_when_size(1)
3986 .force_completion_on_stop(true)
3987 .build()
3988 .unwrap(),
3989 )
3990 .build_canonical()
3991 .unwrap();
3992
3993 if let CanonicalStepSpec::Aggregate(agg) = &spec.steps[0] {
3994 assert_eq!(agg.force_completion_on_stop, Some(true));
3995 } else {
3996 panic!("Expected Aggregate step");
3997 }
3998 }
3999
4000 #[test]
4003 fn test_build_canonical_aggregate_max_buckets_and_ttl() {
4004 use camel_api::aggregator::AggregatorConfig;
4005
4006 let spec = RouteBuilder::from("direct:start")
4007 .route_id("canonical-buckets-ttl")
4008 .aggregate(
4009 AggregatorConfig::correlate_by("key")
4010 .complete_when_size(1)
4011 .max_buckets(100)
4012 .bucket_ttl(Duration::from_secs(60))
4013 .build()
4014 .unwrap(),
4015 )
4016 .build_canonical()
4017 .unwrap();
4018
4019 if let CanonicalStepSpec::Aggregate(agg) = &spec.steps[0] {
4020 assert_eq!(agg.max_buckets, Some(100));
4021 assert_eq!(agg.bucket_ttl_ms, Some(60_000));
4022 } else {
4023 panic!("Expected Aggregate step");
4024 }
4025 }
4026
4027 #[test]
4036 fn canonicalize_aggregate_expression_maps_correlation_key() {
4037 use camel_api::aggregator::AggregatorConfig;
4038
4039 let config = AggregatorConfig::correlate_by("seed")
4040 .correlate_by_expr("${header.orderId}", "simple")
4041 .complete_when_size(2)
4042 .build()
4043 .unwrap();
4044
4045 let spec = canonicalize_aggregate(config).unwrap();
4046
4047 assert_eq!(spec.correlation_key, Some("${header.orderId}".to_string()));
4048 assert_eq!(spec.header, "${header.orderId}");
4049 }
4050
4051 #[test]
4052 fn canonicalize_aggregate_header_only_maps_no_correlation_key() {
4053 use camel_api::aggregator::AggregatorConfig;
4054
4055 let config = AggregatorConfig::correlate_by("orderId")
4056 .complete_when_size(2)
4057 .build()
4058 .unwrap();
4059
4060 let spec = canonicalize_aggregate(config).unwrap();
4061
4062 assert_eq!(spec.correlation_key, None);
4063 assert_eq!(spec.header, "orderId");
4064 }
4065
4066 #[test]
4069 fn test_split_builder_with_filter_inside() {
4070 use camel_api::splitter::{SplitterConfig, split_body_lines};
4071
4072 let definition = RouteBuilder::from("timer:test")
4073 .route_id("split-with-filter")
4074 .split(SplitterConfig::new(split_body_lines()))
4075 .filter(|_ex| true)
4076 .to("mock:filtered-frag")
4077 .end_filter()
4078 .end_split()
4079 .build()
4080 .unwrap();
4081
4082 if let BuilderStep::Split { steps, .. } = &definition.steps()[0] {
4083 assert_eq!(steps.len(), 1);
4084 assert!(matches!(&steps[0], BuilderStep::Filter { .. }));
4085 } else {
4086 panic!("Expected Split step");
4087 }
4088 }
4089
4090 #[test]
4093 fn test_wire_tap_multiple_taps() {
4094 let definition = RouteBuilder::from("timer:tick")
4095 .route_id("multi-wire-tap")
4096 .wire_tap("mock:tap1")
4097 .wire_tap("mock:tap2")
4098 .to("mock:result")
4099 .build()
4100 .unwrap();
4101
4102 assert_eq!(definition.steps().len(), 3);
4103 assert!(
4104 matches!(&definition.steps()[0], BuilderStep::WireTap { uri } if uri == "mock:tap1")
4105 );
4106 assert!(
4107 matches!(&definition.steps()[1], BuilderStep::WireTap { uri } if uri == "mock:tap2")
4108 );
4109 }
4110
4111 #[test]
4114 fn test_builder_shorthand_then_explicit_mixed_mode() {
4115 let result = RouteBuilder::from("direct:start")
4116 .route_id("mixed-mode-2")
4117 .dead_letter_channel("log:dlc")
4118 .error_handler(ErrorHandlerConfig::log_only())
4119 .to("mock:out")
4120 .build();
4121
4122 let err = result.err().expect("mixed mode should fail");
4123 assert!(format!("{err}").contains("mixed error handler modes"));
4124 }
4125
4126 #[test]
4129 fn test_build_canonical_empty_from_uri_errors() {
4130 let result = RouteBuilder::from("").route_id("test").build_canonical();
4131 assert!(result.is_err());
4132 }
4133
4134 #[test]
4135 fn test_build_canonical_missing_route_id_errors() {
4136 let result = RouteBuilder::from("direct:start").build_canonical();
4137 assert!(result.is_err());
4138 let err = result.unwrap_err().to_string();
4139 assert!(err.contains("route_id"));
4140 }
4141
4142 #[test]
4145 fn test_split_builder_with_aggregate_inside() {
4146 use camel_api::aggregator::AggregatorConfig;
4147 use camel_api::splitter::{SplitterConfig, split_body_lines};
4148
4149 let definition = RouteBuilder::from("timer:test")
4150 .route_id("split-agg")
4151 .split(SplitterConfig::new(split_body_lines()))
4152 .aggregate(
4153 AggregatorConfig::correlate_by("frag-key")
4154 .complete_when_size(3)
4155 .build()
4156 .unwrap(),
4157 )
4158 .end_split()
4159 .build()
4160 .unwrap();
4161
4162 if let BuilderStep::Split { steps, .. } = &definition.steps()[0] {
4163 assert_eq!(steps.len(), 1);
4164 assert!(matches!(&steps[0], BuilderStep::Aggregate { .. }));
4165 } else {
4166 panic!("Expected Split step");
4167 }
4168 }
4169
4170 #[test]
4173 fn test_throttle_builder_with_steps_inside() {
4174 let definition = RouteBuilder::from("timer:tick")
4175 .route_id("throttle-steps")
4176 .throttle(10, Duration::from_secs(1))
4177 .set_header("throttled", Value::Bool(true))
4178 .to("mock:throttled")
4179 .end_throttle()
4180 .build()
4181 .unwrap();
4182
4183 if let BuilderStep::Throttle { steps, .. } = &definition.steps()[0] {
4184 assert_eq!(steps.len(), 2);
4185 } else {
4186 panic!("Expected Throttle step");
4187 }
4188 }
4189
4190 #[test]
4193 fn test_load_balance_builder_with_steps_inside() {
4194 let definition = RouteBuilder::from("timer:tick")
4195 .route_id("lb-steps")
4196 .load_balance()
4197 .round_robin()
4198 .set_header("lb", Value::Bool(true))
4199 .to("mock:lb")
4200 .end_load_balance()
4201 .build()
4202 .unwrap();
4203
4204 if let BuilderStep::LoadBalance { steps, .. } = &definition.steps()[0] {
4205 assert_eq!(steps.len(), 2);
4206 } else {
4207 panic!("Expected LoadBalance step");
4208 }
4209 }
4210
4211 #[test]
4214 fn test_multicast_builder_with_steps_inside() {
4215 let definition = RouteBuilder::from("timer:tick")
4216 .route_id("multicast-steps")
4217 .multicast()
4218 .set_header("mc", Value::Bool(true))
4219 .to("mock:multicast")
4220 .end_multicast()
4221 .build()
4222 .unwrap();
4223
4224 if let BuilderStep::Multicast { steps, .. } = &definition.steps()[0] {
4225 assert_eq!(steps.len(), 2);
4226 } else {
4227 panic!("Expected Multicast step");
4228 }
4229 }
4230
4231 #[test]
4234 fn test_loop_builder_with_steps_inside() {
4235 let definition = RouteBuilder::from("timer:tick")
4236 .route_id("loop-steps")
4237 .loop_count(3)
4238 .set_header("loop", Value::Bool(true))
4239 .to("mock:loop")
4240 .end_loop()
4241 .build()
4242 .unwrap();
4243
4244 if let BuilderStep::Loop { steps, .. } = &definition.steps()[0] {
4245 assert_eq!(steps.len(), 2);
4246 } else {
4247 panic!("Expected Loop step");
4248 }
4249 }
4250
4251 #[test]
4254 fn test_build_canonical_rejects_loop_step() {
4255 let err = RouteBuilder::from("direct:start")
4256 .route_id("canonical-loop")
4257 .loop_count(3)
4258 .to("mock:loop")
4259 .end_loop()
4260 .build_canonical()
4261 .unwrap_err();
4262
4263 assert!(format!("{err}").contains("does not support step `loop`"));
4264 }
4265
4266 #[test]
4267 fn test_build_canonical_rejects_multicast_step() {
4268 let err = RouteBuilder::from("direct:start")
4269 .route_id("canonical-multicast")
4270 .multicast()
4271 .to("mock:a")
4272 .end_multicast()
4273 .build_canonical()
4274 .unwrap_err();
4275
4276 assert!(format!("{err}").contains("does not support step `multicast`"));
4277 }
4278
4279 #[test]
4280 fn test_build_canonical_rejects_throttle_step() {
4281 let err = RouteBuilder::from("direct:start")
4282 .route_id("canonical-throttle")
4283 .throttle(10, Duration::from_secs(1))
4284 .to("mock:result")
4285 .end_throttle()
4286 .build_canonical()
4287 .unwrap_err();
4288
4289 assert!(format!("{err}").contains("does not support step `throttle`"));
4290 }
4291
4292 #[test]
4293 fn test_build_canonical_rejects_load_balancer_step() {
4294 let err = RouteBuilder::from("direct:start")
4295 .route_id("canonical-lb")
4296 .load_balance()
4297 .round_robin()
4298 .to("mock:result")
4299 .end_load_balance()
4300 .build_canonical()
4301 .unwrap_err();
4302
4303 assert!(format!("{err}").contains("does not support step `load_balancer`"));
4304 }
4305
4306 #[test]
4307 fn test_build_canonical_rejects_bean_step() {
4308 let err = RouteBuilder::from("direct:start")
4309 .route_id("canonical-bean")
4310 .bean("myBean", "process")
4311 .build_canonical()
4312 .unwrap_err();
4313
4314 assert!(format!("{err}").contains("does not support step `bean`"));
4315 }
4316
4317 #[test]
4318 fn test_build_canonical_rejects_script_step() {
4319 let err = RouteBuilder::from("direct:start")
4320 .route_id("canonical-script")
4321 .script("rhai", "x = 1")
4322 .build_canonical()
4323 .unwrap_err();
4324
4325 assert!(format!("{err}").contains("does not support step `script`"));
4326 }
4327
4328 #[test]
4329 fn test_build_canonical_accepts_delay_step() {
4330 let spec = RouteBuilder::from("direct:start")
4331 .route_id("canonical-delay")
4332 .delay(Duration::from_millis(100))
4333 .build_canonical()
4334 .unwrap();
4335
4336 assert!(
4337 spec.steps.iter().any(
4338 |s| matches!(s, CanonicalStepSpec::Delay { delay_ms, .. } if *delay_ms == 100)
4339 )
4340 );
4341 }
4342
4343 #[test]
4344 fn test_build_canonical_accepts_wire_tap_step() {
4345 let spec = RouteBuilder::from("direct:start")
4346 .route_id("canonical-wiretap")
4347 .wire_tap("mock:tap")
4348 .build_canonical()
4349 .unwrap();
4350
4351 assert!(
4352 spec.steps
4353 .iter()
4354 .any(|s| matches!(s, CanonicalStepSpec::WireTap { uri } if uri == "mock:tap"))
4355 );
4356 }
4357
4358 #[test]
4359 fn test_build_canonical_rejects_dynamic_router_step() {
4360 let err = RouteBuilder::from("direct:start")
4361 .route_id("canonical-dyn-router")
4362 .dynamic_router(Arc::new(|_| Some("mock:a".to_string())))
4363 .build_canonical()
4364 .unwrap_err();
4365
4366 assert!(format!("{err}").contains("does not support step `dynamic_router`"));
4367 }
4368
4369 #[test]
4370 fn test_build_canonical_rejects_routing_slip_step() {
4371 let err = RouteBuilder::from("direct:start")
4372 .route_id("canonical-routing-slip")
4373 .routing_slip(Arc::new(|_| Some("mock:a".to_string())))
4374 .build_canonical()
4375 .unwrap_err();
4376
4377 assert!(format!("{err}").contains("does not support step `routing_slip`"));
4378 }
4379
4380 #[test]
4381 fn test_build_canonical_rejects_recipient_list_step() {
4382 let err = RouteBuilder::from("direct:start")
4383 .route_id("canonical-recipient")
4384 .recipient_list(Arc::new(|_| "mock:a".to_string()))
4385 .build_canonical()
4386 .unwrap_err();
4387
4388 assert!(format!("{err}").contains("does not support step `recipient_list`"));
4389 }
4390
4391 #[test]
4394 fn test_build_canonical_rejects_any_mode_with_predicate() {
4395 let err = RouteBuilder::from("direct:start")
4396 .route_id("canonical-any-pred")
4397 .aggregate(AggregatorConfig {
4398 header_name: "key".to_string(),
4399 completion: CompletionMode::Any(vec![
4400 CompletionCondition::Size(5),
4401 CompletionCondition::Predicate(Arc::new(|_| false)),
4402 ]),
4403 correlation: CorrelationStrategy::HeaderName("key".to_string()),
4404 strategy: AggregationStrategy::CollectAll,
4405 max_buckets: None,
4406 max_bucket_size: None,
4407 bucket_ttl: None,
4408 force_completion_on_stop: false,
4409 discard_on_timeout: false,
4410 max_timeout_tasks: 1024,
4411 })
4412 .build_canonical()
4413 .unwrap_err();
4414
4415 assert!(
4416 format!("{err}").contains("cannot reverse-map"),
4417 "reject message must explain forward-only: {}",
4418 err
4419 );
4420 }
4421
4422 #[test]
4425 fn test_builder_validation_missing_from_uri() {
4426 let result = RouteBuilder::from("")
4427 .route_id("missing-uri-route")
4428 .to("log:info")
4429 .build();
4430 assert!(result.is_err(), "empty from URI should fail validation");
4431 let err = result.err().unwrap().to_string();
4432 assert!(
4433 err.contains("'from'") || err.contains("URI"),
4434 "error should mention from/URI, got: {err}"
4435 );
4436 }
4437
4438 #[test]
4439 fn test_builder_validation_invalid_step_uri_scheme() {
4440 let result = RouteBuilder::from("timer:tick")
4441 .route_id("bad-step-route")
4442 .to("not-a-valid-uri") .build();
4444 assert!(
4447 result.is_ok(),
4448 "builder should accept opaque step URIs; resolution happens later"
4449 );
4450 }
4451
4452 #[test]
4455 fn test_builder_duplicate_route_ids_produce_identical_definitions() {
4456 let route1 = RouteBuilder::from("direct:a")
4459 .route_id("dup-route")
4460 .to("mock:out")
4461 .build();
4462 let route2 = RouteBuilder::from("direct:b")
4463 .route_id("dup-route")
4464 .to("mock:out")
4465 .build();
4466
4467 assert!(route1.is_ok());
4468 assert!(route2.is_ok());
4469 assert_eq!(route1.unwrap().route_id(), route2.unwrap().route_id());
4470 }
4471
4472 #[test]
4473 fn test_builder_clone_reuse_as_template() {
4474 let template = RouteBuilder::from("direct:in")
4477 .set_header("stage", Value::String("shared".into()))
4478 .log("shared prefix", LogLevel::Info);
4479
4480 let route_a = template
4481 .clone()
4482 .route_id("route-a")
4483 .to("mock:a")
4484 .build()
4485 .expect("clone A builds");
4486 let route_b = template
4487 .route_id("route-b")
4488 .to("mock:b")
4489 .build()
4490 .expect("clone B builds");
4491
4492 assert_eq!(route_a.route_id(), "route-a");
4493 assert_eq!(route_b.route_id(), "route-b");
4494 assert_eq!(route_a.steps().len(), route_b.steps().len());
4497 assert_eq!(route_a.from_uri(), route_b.from_uri());
4498 }
4499}