1use std::sync::Arc;
5
6use camel_api::UnitOfWorkConfig;
7use camel_api::circuit_breaker::CircuitBreakerConfig;
8use camel_api::error_handler::ErrorHandlerConfig;
9use camel_api::loop_eip::LoopConfig;
10use camel_api::security_policy::SecurityPolicyConfig;
11use camel_api::{
12 AggregatorConfig, FilterPredicate, MulticastConfig, OpaqueProcessor, ResequencePolicyConfig,
13 SpanKindHint, SplitterConfig,
14};
15use camel_auth::TokenAuthenticator;
16use camel_component_api::ConcurrencyModel;
17
18#[derive(Clone)]
20pub struct WhenStep {
21 pub predicate: FilterPredicate,
22 pub steps: Vec<BuilderStep>,
23}
24
25impl std::fmt::Debug for WhenStep {
26 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
27 f.debug_struct("WhenStep")
28 .field("predicate", &self.predicate)
29 .field("steps", &self.steps)
30 .finish()
31 }
32}
33
34pub use camel_api::declarative::{LanguageExpressionDef, ValueSourceDef};
35
36#[derive(Debug, Clone)]
38pub struct DeclarativeWhenStep {
39 pub predicate: LanguageExpressionDef,
40 pub steps: Vec<BuilderStep>,
41}
42
43#[derive(Debug, Clone)]
45pub struct DoTryCatchClauseBuilder {
46 pub exception: Option<Vec<String>>,
47 pub when: Option<LanguageExpressionDef>,
48 pub on_when: Option<LanguageExpressionDef>,
49 pub disposition: camel_api::error_handler::ExceptionDisposition,
50 pub steps: Vec<BuilderStep>,
51}
52
53#[derive(Debug, Clone)]
55pub struct DoTryFinallyBuilder {
56 pub on_when: Option<LanguageExpressionDef>,
57 pub steps: Vec<BuilderStep>,
58}
59
60#[derive(Debug, Clone)]
62pub enum BuilderStep {
63 Processor(OpaqueProcessor),
65 To(String),
67 Stop,
69 Log {
71 level: camel_processor::LogLevel,
72 message: String,
73 },
74 DeclarativeSetHeader {
76 key: String,
77 value: ValueSourceDef,
78 },
79 DeclarativeSetHeaderIfAbsent {
81 key: String,
82 value: ValueSourceDef,
83 },
84 DeclarativeRemoveHeader {
86 key: String,
87 },
88 DeclarativeSetProperty {
89 key: String,
90 value_source: ValueSourceDef,
91 },
92 DeclarativeSetBody {
94 value: ValueSourceDef,
95 },
96 DeclarativeFilter {
98 predicate: LanguageExpressionDef,
99 steps: Vec<BuilderStep>,
100 },
101 DeclarativeChoice {
103 whens: Vec<DeclarativeWhenStep>,
104 otherwise: Option<Vec<BuilderStep>>,
105 },
106 DeclarativeScript {
108 expression: LanguageExpressionDef,
109 },
110 DeclarativeFunction {
111 definition: camel_api::FunctionDefinition,
112 },
113 DeclarativeSplit {
115 expression: LanguageExpressionDef,
116 aggregation: camel_api::splitter::AggregationStrategy,
117 parallel: bool,
118 parallel_limit: Option<usize>,
119 trace_item_threshold: Option<usize>,
122 stop_on_exception: bool,
123 steps: Vec<BuilderStep>,
124 },
125 DeclarativeStreamSplit {
127 stream_config: camel_api::StreamSplitConfig,
128 aggregation: camel_api::splitter::AggregationStrategy,
129 stop_on_exception: bool,
130 steps: Vec<BuilderStep>,
131 },
132 DeclarativeDynamicRouter {
133 expression: LanguageExpressionDef,
134 uri_delimiter: String,
135 cache_size: i32,
136 ignore_invalid_endpoints: bool,
137 max_iterations: usize,
138 },
139 DeclarativeRoutingSlip {
140 expression: LanguageExpressionDef,
141 uri_delimiter: String,
142 cache_size: i32,
143 ignore_invalid_endpoints: bool,
144 },
145 Split {
147 config: SplitterConfig,
148 steps: Vec<BuilderStep>,
149 },
150 Aggregate {
152 config: AggregatorConfig,
153 },
154 Filter {
156 predicate: FilterPredicate,
157 steps: Vec<BuilderStep>,
158 },
159 Choice {
162 whens: Vec<WhenStep>,
163 otherwise: Option<Vec<BuilderStep>>,
164 },
165 WireTap {
167 uri: String,
168 },
169 Multicast {
171 steps: Vec<BuilderStep>,
172 config: MulticastConfig,
173 },
174 DeclarativeLog {
176 level: camel_processor::LogLevel,
177 message: ValueSourceDef,
178 },
179 Bean {
181 name: String,
182 method: String,
183 },
184 Script {
187 language: String,
188 script: String,
189 },
190 Throttle {
192 config: camel_api::ThrottlerConfig,
193 steps: Vec<BuilderStep>,
194 },
195 LoadBalance {
197 config: camel_api::LoadBalancerConfig,
198 steps: Vec<BuilderStep>,
199 },
200 DynamicRouter {
202 config: camel_api::DynamicRouterConfig,
203 },
204 RoutingSlip {
205 config: camel_api::RoutingSlipConfig,
206 },
207 RecipientList {
208 config: camel_api::recipient_list::RecipientListConfig,
209 },
210 DeclarativeRecipientList {
211 expression: LanguageExpressionDef,
212 delimiter: String,
213 parallel: bool,
214 parallel_limit: Option<usize>,
215 stop_on_exception: bool,
216 aggregation: String,
217 },
218 Delay {
219 config: camel_api::DelayConfig,
220 },
221 Loop {
223 config: LoopConfig,
224 steps: Vec<BuilderStep>,
225 },
226 DeclarativeLoop {
228 count: Option<usize>,
229 while_predicate: Option<LanguageExpressionDef>,
230 steps: Vec<BuilderStep>,
231 max_iterations: Option<usize>,
232 },
233 Enrich {
235 uri: String,
236 strategy: Option<String>,
237 timeout_ms: Option<u64>,
238 },
239 PollEnrich {
241 uri: String,
242 strategy: Option<String>,
243 timeout_ms: Option<u64>,
244 },
245 Validate {
248 predicate: LanguageExpressionDef,
249 },
250 ClaimCheck {
254 repository: String,
255 operation: String,
256 key: LanguageExpressionDef,
257 filter: Option<String>,
258 },
259 Sampling {
263 period: usize,
264 },
265 Sort {
268 expression: LanguageExpressionDef,
269 reverse: bool,
270 },
271 IdempotentConsumer {
275 repository: String,
276 expression: LanguageExpressionDef,
277 steps: Vec<BuilderStep>,
278 eager: bool,
279 remove_on_failure: bool,
280 },
281 Cache {
285 repository: Option<String>,
286 key: LanguageExpressionDef,
287 ttl: Option<String>,
288 max_entry_bytes: Option<usize>,
289 coalesce_misses: bool,
291 on_miss: Vec<BuilderStep>,
292 },
293 CacheInvalidate {
298 repository: Option<String>,
299 key: Option<LanguageExpressionDef>,
300 key_prefix: Option<LanguageExpressionDef>,
301 },
302 CacheClear {
305 repository: Option<String>,
306 },
307 CacheStats {
310 repository: Option<String>,
311 },
312 CachePeekStale {
315 repository: Option<String>,
316 key: LanguageExpressionDef,
317 on_miss: camel_processor::PeekStaleMissPolicy,
318 },
319 DeclarativeDoTry {
321 try_steps: Vec<BuilderStep>,
322 catch: Vec<DoTryCatchClauseBuilder>,
323 finally: Option<DoTryFinallyBuilder>,
324 },
325 Resequence {
328 policy_config: ResequencePolicyConfig,
329 },
330}
331
332impl BuilderStep {
333 pub(crate) fn span_label(&self) -> Option<String> {
347 match self {
348 Self::Processor(_) | Self::Stop => None,
350
351 Self::To(uri) => uri.contains(':').then(|| {
352 let scheme = uri.split(':').next().unwrap_or_default();
353 format!("to:{scheme}")
354 }),
355
356 Self::Log { .. } | Self::DeclarativeLog { .. } => Some("log".into()),
357
358 Self::DeclarativeSetHeader { .. } => Some("set-header".into()),
359 Self::DeclarativeSetHeaderIfAbsent { .. } => Some("set-header-if-absent".into()),
360 Self::DeclarativeRemoveHeader { .. } => Some("remove-header".into()),
361 Self::DeclarativeSetProperty { .. } => Some("set-property".into()),
362 Self::DeclarativeSetBody { .. } => Some("set-body".into()),
363
364 Self::DeclarativeFilter { .. } | Self::Filter { .. } => Some("filter".into()),
365 Self::DeclarativeChoice { .. } | Self::Choice { .. } => Some("choice".into()),
366 Self::DeclarativeScript { .. } | Self::Script { .. } => Some("script".into()),
367 Self::DeclarativeFunction { .. } => Some("function".into()),
368
369 Self::DeclarativeSplit { .. }
370 | Self::DeclarativeStreamSplit { .. }
371 | Self::Split { .. } => Some("split".into()),
372
373 Self::DeclarativeDynamicRouter { .. } | Self::DynamicRouter { .. } => {
374 Some("dynamic-router".into())
375 }
376 Self::DeclarativeRoutingSlip { .. } | Self::RoutingSlip { .. } => {
377 Some("routing-slip".into())
378 }
379 Self::DeclarativeRecipientList { .. } | Self::RecipientList { .. } => {
380 Some("recipient-list".into())
381 }
382
383 Self::Aggregate { .. } => Some("aggregate".into()),
384 Self::WireTap { .. } => Some("wire-tap".into()),
385 Self::Multicast { .. } => Some("multicast".into()),
386 Self::Bean { .. } => Some("bean".into()),
387 Self::Throttle { .. } => Some("throttle".into()),
388 Self::LoadBalance { .. } => Some("load-balance".into()),
389 Self::Delay { .. } => Some("delay".into()),
390 Self::Loop { .. } | Self::DeclarativeLoop { .. } => Some("loop".into()),
391 Self::Enrich { .. } => Some("enrich".into()),
392 Self::PollEnrich { .. } => Some("poll-enrich".into()),
393 Self::Validate { .. } => Some("validate".into()),
394 Self::ClaimCheck { .. } => Some("claim-check".into()),
395 Self::Sampling { .. } => Some("sampling".into()),
396 Self::Sort { .. } => Some("sort".into()),
397 Self::IdempotentConsumer { .. } => Some("idempotent-consumer".into()),
398 Self::Cache { .. } => Some("cache".into()),
399 Self::CacheInvalidate { .. } => Some("cache-invalidate".into()),
400 Self::CacheClear { .. } => Some("cache-clear".into()),
401 Self::CacheStats { .. } => Some("cache-stats".into()),
402 Self::CachePeekStale { .. } => Some("cache-peek-stale".into()),
403 Self::DeclarativeDoTry { .. } => Some("do-try".into()),
404 Self::Resequence { .. } => Some("resequence".into()),
405 }
406 }
407
408 pub(crate) fn to_uri_metadata(&self) -> Option<Arc<str>> {
419 match self {
420 Self::To(uri) => Some(Arc::from(uri.as_str())),
421
422 Self::Processor(..)
426 | Self::Stop
427 | Self::Log { .. }
428 | Self::DeclarativeSetHeader { .. }
429 | Self::DeclarativeSetHeaderIfAbsent { .. }
430 | Self::DeclarativeRemoveHeader { .. }
431 | Self::DeclarativeSetProperty { .. }
432 | Self::DeclarativeSetBody { .. }
433 | Self::DeclarativeFilter { .. }
434 | Self::DeclarativeChoice { .. }
435 | Self::DeclarativeScript { .. }
436 | Self::DeclarativeFunction { .. }
437 | Self::DeclarativeSplit { .. }
438 | Self::DeclarativeStreamSplit { .. }
439 | Self::DeclarativeDynamicRouter { .. }
440 | Self::DeclarativeRoutingSlip { .. }
441 | Self::Split { .. }
442 | Self::Aggregate { .. }
443 | Self::Filter { .. }
444 | Self::Choice { .. }
445 | Self::Multicast { .. }
446 | Self::DeclarativeLog { .. }
447 | Self::Bean { .. }
448 | Self::Script { .. }
449 | Self::Throttle { .. }
450 | Self::LoadBalance { .. }
451 | Self::DynamicRouter { .. }
452 | Self::RoutingSlip { .. }
453 | Self::RecipientList { .. }
454 | Self::DeclarativeRecipientList { .. }
455 | Self::Delay { .. }
456 | Self::Loop { .. }
457 | Self::DeclarativeLoop { .. }
458 | Self::Enrich { .. }
459 | Self::PollEnrich { .. }
460 | Self::WireTap { .. }
461 | Self::Validate { .. }
462 | Self::ClaimCheck { .. }
463 | Self::Sampling { .. }
464 | Self::Sort { .. }
465 | Self::IdempotentConsumer { .. }
466 | Self::Cache { .. }
467 | Self::CacheInvalidate { .. }
468 | Self::CacheClear { .. }
469 | Self::CacheStats { .. }
470 | Self::CachePeekStale { .. }
471 | Self::DeclarativeDoTry { .. }
472 | Self::Resequence { .. } => None,
473 }
474 }
475
476 pub(crate) fn span_kind_hint(&self) -> SpanKindHint {
492 match self {
493 Self::To(uri)
495 | Self::Enrich { uri, .. }
496 | Self::PollEnrich { uri, .. }
497 | Self::WireTap { uri, .. } => uri_span_kind(uri),
498
499 Self::Processor(..)
503 | Self::Stop
504 | Self::Log { .. }
505 | Self::DeclarativeSetHeader { .. }
506 | Self::DeclarativeSetHeaderIfAbsent { .. }
507 | Self::DeclarativeRemoveHeader { .. }
508 | Self::DeclarativeSetProperty { .. }
509 | Self::DeclarativeSetBody { .. }
510 | Self::DeclarativeFilter { .. }
511 | Self::DeclarativeChoice { .. }
512 | Self::DeclarativeScript { .. }
513 | Self::DeclarativeFunction { .. }
514 | Self::DeclarativeSplit { .. }
515 | Self::DeclarativeStreamSplit { .. }
516 | Self::DeclarativeDynamicRouter { .. }
517 | Self::DeclarativeRoutingSlip { .. }
518 | Self::Split { .. }
519 | Self::Aggregate { .. }
520 | Self::Filter { .. }
521 | Self::Choice { .. }
522 | Self::Multicast { .. }
523 | Self::DeclarativeLog { .. }
524 | Self::Bean { .. }
525 | Self::Script { .. }
526 | Self::Throttle { .. }
527 | Self::LoadBalance { .. }
528 | Self::DynamicRouter { .. }
529 | Self::RoutingSlip { .. }
530 | Self::RecipientList { .. }
531 | Self::DeclarativeRecipientList { .. }
532 | Self::Delay { .. }
533 | Self::Loop { .. }
534 | Self::DeclarativeLoop { .. }
535 | Self::Validate { .. }
536 | Self::ClaimCheck { .. }
537 | Self::Sampling { .. }
538 | Self::Sort { .. }
539 | Self::IdempotentConsumer { .. }
540 | Self::Cache { .. }
541 | Self::CacheInvalidate { .. }
542 | Self::CacheClear { .. }
543 | Self::CacheStats { .. }
544 | Self::CachePeekStale { .. }
545 | Self::DeclarativeDoTry { .. }
546 | Self::Resequence { .. } => SpanKindHint::Internal,
547 }
548 }
549}
550
551fn uri_span_kind(uri: &str) -> SpanKindHint {
558 const PRODUCER_SCHEMES: [&str; 5] = ["kafka", "jms", "activemq", "artemis", "mqtt"];
560 const CLIENT_SCHEMES: [&str; 12] = [
562 "http",
563 "https",
564 "grpc",
565 "grpcs",
566 "ws",
567 "redis",
568 "opensearch",
569 "sql",
570 "surrealdb",
571 "cxf",
572 "llm",
573 "mcp",
574 ];
575
576 let scheme = uri.split(':').next().unwrap_or_default();
577 if PRODUCER_SCHEMES
578 .iter()
579 .any(|s| s.eq_ignore_ascii_case(scheme))
580 {
581 SpanKindHint::Producer
582 } else if CLIENT_SCHEMES
583 .iter()
584 .any(|s| s.eq_ignore_ascii_case(scheme))
585 {
586 SpanKindHint::Client
587 } else {
588 SpanKindHint::Internal
589 }
590}
591
592pub struct RouteDefinition {
594 pub(crate) from_uri: String,
595 pub(crate) steps: Vec<BuilderStep>,
596 pub(crate) error_handler: Option<ErrorHandlerConfig>,
598 pub(crate) circuit_breaker: Option<CircuitBreakerConfig>,
600 pub(crate) circuit_breaker_fallback: Vec<BuilderStep>,
608 pub(crate) security_policy: Option<SecurityPolicyConfig>,
609 pub(crate) security_authenticator: Option<Arc<dyn TokenAuthenticator>>,
611 pub(crate) provider_registry: Option<Arc<camel_auth::ProviderRegistry>>,
614 pub(crate) security_provider: Option<String>,
618 pub(crate) security_audiences: Option<Vec<String>>,
623 pub(crate) unit_of_work: Option<UnitOfWorkConfig>,
625 pub(crate) concurrency: Option<ConcurrencyModel>,
628 pub(crate) route_id: String,
630 pub(crate) auto_startup: bool,
632 pub(crate) startup_order: i32,
634 pub(crate) source_hash: Option<u64>,
635}
636
637impl RouteDefinition {
638 pub fn new(from_uri: impl Into<String>, steps: Vec<BuilderStep>) -> Self {
640 Self {
641 from_uri: from_uri.into(),
642 steps,
643 error_handler: None,
644 circuit_breaker: None,
645 circuit_breaker_fallback: Vec::new(),
646 security_policy: None,
647 security_authenticator: None,
648 provider_registry: None,
649 security_provider: None,
650 security_audiences: None,
651 unit_of_work: None,
652 concurrency: None,
653 route_id: String::new(), auto_startup: true,
655 startup_order: 1000,
656 source_hash: None,
657 }
658 }
659
660 pub fn from_uri(&self) -> &str {
662 &self.from_uri
663 }
664
665 pub fn steps(&self) -> &[BuilderStep] {
667 &self.steps
668 }
669
670 pub fn circuit_breaker_fallback(&self) -> &[BuilderStep] {
671 &self.circuit_breaker_fallback
672 }
673
674 pub fn map_steps(mut self, f: impl FnOnce(Vec<BuilderStep>) -> Vec<BuilderStep>) -> Self {
679 self.steps = f(self.steps);
680 self
681 }
682
683 pub fn with_error_handler(mut self, config: ErrorHandlerConfig) -> Self {
685 self.error_handler = Some(config);
686 self
687 }
688
689 pub fn error_handler_config(&self) -> Option<&ErrorHandlerConfig> {
691 self.error_handler.as_ref()
692 }
693
694 pub fn with_circuit_breaker(mut self, config: CircuitBreakerConfig) -> Self {
696 self.circuit_breaker = Some(config);
697 self
698 }
699
700 pub fn with_circuit_breaker_fallback(mut self, steps: Vec<BuilderStep>) -> Self {
707 self.circuit_breaker_fallback = steps;
708 self
709 }
710
711 pub fn with_security_policy(mut self, config: SecurityPolicyConfig) -> Self {
713 self.security_policy = Some(config);
714 self
715 }
716
717 pub fn with_security_authenticator(
719 mut self,
720 authenticator: Arc<dyn TokenAuthenticator>,
721 ) -> Self {
722 self.security_authenticator = Some(authenticator);
723 self
724 }
725
726 pub fn with_provider_registry(mut self, registry: Arc<camel_auth::ProviderRegistry>) -> Self {
732 self.provider_registry = Some(registry);
733 self
734 }
735
736 pub fn with_security_provider(mut self, name: impl Into<String>) -> Self {
738 self.security_provider = Some(name.into());
739 self
740 }
741
742 pub fn with_security_audiences(mut self, audiences: Vec<String>) -> Self {
744 self.security_audiences = Some(audiences);
745 self
746 }
747
748 pub fn security_provider(&self) -> Option<&str> {
750 self.security_provider.as_deref()
751 }
752
753 pub fn security_audiences(&self) -> Option<&[String]> {
755 self.security_audiences.as_deref()
756 }
757
758 pub fn with_unit_of_work(mut self, config: UnitOfWorkConfig) -> Self {
760 self.unit_of_work = Some(config);
761 self
762 }
763
764 pub fn unit_of_work_config(&self) -> Option<&UnitOfWorkConfig> {
766 self.unit_of_work.as_ref()
767 }
768
769 pub fn circuit_breaker_config(&self) -> Option<&CircuitBreakerConfig> {
771 self.circuit_breaker.as_ref()
772 }
773
774 pub fn security_policy_config(&self) -> Option<&SecurityPolicyConfig> {
775 self.security_policy.as_ref()
776 }
777
778 pub fn security_authenticator(&self) -> Option<&Arc<dyn TokenAuthenticator>> {
779 self.security_authenticator.as_ref()
780 }
781
782 pub fn concurrency_override(&self) -> Option<&ConcurrencyModel> {
784 self.concurrency.as_ref()
785 }
786
787 pub fn with_concurrency(mut self, model: ConcurrencyModel) -> Self {
789 self.concurrency = Some(model);
790 self
791 }
792
793 pub fn route_id(&self) -> &str {
795 &self.route_id
796 }
797
798 pub fn auto_startup(&self) -> bool {
800 self.auto_startup
801 }
802
803 pub fn startup_order(&self) -> i32 {
805 self.startup_order
806 }
807
808 pub fn with_route_id(mut self, id: impl Into<String>) -> Self {
810 self.route_id = id.into();
811 self
812 }
813
814 pub fn with_auto_startup(mut self, auto: bool) -> Self {
816 self.auto_startup = auto;
817 self
818 }
819
820 pub fn with_startup_order(mut self, order: i32) -> Self {
822 self.startup_order = order;
823 self
824 }
825
826 pub fn with_source_hash(mut self, hash: u64) -> Self {
827 self.source_hash = Some(hash);
828 self
829 }
830
831 pub fn source_hash(&self) -> Option<u64> {
832 self.source_hash
833 }
834
835 pub fn to_info(&self) -> RouteDefinitionInfo {
838 RouteDefinitionInfo {
839 route_id: self.route_id.clone(),
840 auto_startup: self.auto_startup,
841 startup_order: self.startup_order,
842 source_hash: self.source_hash,
843 }
844 }
845}
846
847#[derive(Clone)]
853pub struct RouteDefinitionInfo {
854 route_id: String,
855 auto_startup: bool,
856 startup_order: i32,
857 pub(crate) source_hash: Option<u64>,
858}
859
860impl RouteDefinitionInfo {
861 pub fn route_id(&self) -> &str {
863 &self.route_id
864 }
865
866 pub fn auto_startup(&self) -> bool {
868 self.auto_startup
869 }
870
871 pub fn startup_order(&self) -> i32 {
873 self.startup_order
874 }
875
876 pub fn source_hash(&self) -> Option<u64> {
877 self.source_hash
878 }
879}
880
881#[cfg(test)]
882mod tests {
883 use super::*;
884
885 #[test]
892 fn builder_step_span_label_mapping() {
893 use camel_api::declarative::LanguageExpressionDef;
894 use camel_api::splitter::AggregationStrategy;
895 use camel_api::{BoxProcessor, IdentityProcessor, OpaqueProcessor};
896
897 let expr = LanguageExpressionDef {
898 language: "simple".into(),
899 source: "${body}".into(),
900 };
901
902 assert_eq!(
904 BuilderStep::To("direct:tree-sub".into())
905 .span_label()
906 .as_deref(),
907 Some("to:direct")
908 );
909 assert_eq!(
910 BuilderStep::To("http://api.example/x".into())
911 .span_label()
912 .as_deref(),
913 Some("to:http")
914 );
915 assert_eq!(BuilderStep::To("garbage".into()).span_label(), None);
917
918 assert_eq!(
920 BuilderStep::Log {
921 level: camel_processor::LogLevel::Info,
922 message: "m".into(),
923 }
924 .span_label()
925 .as_deref(),
926 Some("log")
927 );
928 assert_eq!(
929 BuilderStep::Split {
930 config: camel_api::splitter::SplitterConfig::new(
931 camel_api::splitter::split_body_lines()
932 ),
933 steps: vec![BuilderStep::Stop],
934 }
935 .span_label()
936 .as_deref(),
937 Some("split")
938 );
939 assert_eq!(
940 BuilderStep::DeclarativeSplit {
941 expression: expr.clone(),
942 aggregation: AggregationStrategy::Original,
943 parallel: false,
944 parallel_limit: None,
945 trace_item_threshold: None,
946 stop_on_exception: true,
947 steps: vec![BuilderStep::Stop],
948 }
949 .span_label()
950 .as_deref(),
951 Some("split")
952 );
953 assert_eq!(
954 BuilderStep::DeclarativeStreamSplit {
955 stream_config: camel_api::StreamSplitConfig::default(),
956 aggregation: AggregationStrategy::Original,
957 stop_on_exception: true,
958 steps: vec![BuilderStep::Stop],
959 }
960 .span_label()
961 .as_deref(),
962 Some("split")
963 );
964 assert_eq!(BuilderStep::Stop.span_label(), None);
965
966 assert_eq!(
968 BuilderStep::Processor(OpaqueProcessor(BoxProcessor::new(IdentityProcessor)))
969 .span_label(),
970 None
971 );
972 }
973
974 #[test]
986 fn builder_step_span_kind_hint_mapping() {
987 use camel_api::splitter::split_body_lines;
988 use camel_api::{Exchange, FilterPredicate, SpanKindHint};
989
990 let kind = |uri: &str| BuilderStep::To(uri.into()).span_kind_hint();
991
992 assert_eq!(kind("kafka:orders"), SpanKindHint::Producer);
994 assert_eq!(kind("jms:q"), SpanKindHint::Producer);
995 assert_eq!(kind("activemq:q"), SpanKindHint::Producer);
996 assert_eq!(kind("artemis:q"), SpanKindHint::Producer);
997 assert_eq!(kind("mqtt:t"), SpanKindHint::Producer);
998 assert_eq!(kind("KAFKA:orders"), SpanKindHint::Producer);
1000
1001 assert_eq!(kind("http://x"), SpanKindHint::Client);
1003 assert_eq!(kind("https://x"), SpanKindHint::Client);
1004 assert_eq!(kind("grpc://x"), SpanKindHint::Client);
1005 assert_eq!(kind("grpcs://x"), SpanKindHint::Client);
1006 assert_eq!(kind("ws://x"), SpanKindHint::Client);
1007 assert_eq!(kind("redis://x"), SpanKindHint::Client);
1008 assert_eq!(kind("opensearch://x"), SpanKindHint::Client);
1009 assert_eq!(kind("sql:db"), SpanKindHint::Client);
1010 assert_eq!(kind("surrealdb://x"), SpanKindHint::Client);
1011 assert_eq!(kind("cxf://x"), SpanKindHint::Client);
1012 assert_eq!(kind("llm://x"), SpanKindHint::Client);
1013 assert_eq!(kind("mcp://x"), SpanKindHint::Client);
1014
1015 assert_eq!(kind("direct:y"), SpanKindHint::Internal);
1017 assert_eq!(kind("timer:z"), SpanKindHint::Internal);
1018 assert_eq!(kind("garbage"), SpanKindHint::Internal);
1020
1021 assert_eq!(
1023 BuilderStep::Log {
1024 level: camel_processor::LogLevel::Info,
1025 message: "m".into(),
1026 }
1027 .span_kind_hint(),
1028 SpanKindHint::Internal
1029 );
1030 assert_eq!(
1031 BuilderStep::Filter {
1032 predicate: FilterPredicate::new(|_: &Exchange| true),
1033 steps: vec![BuilderStep::Stop],
1034 }
1035 .span_kind_hint(),
1036 SpanKindHint::Internal
1037 );
1038 assert_eq!(
1039 BuilderStep::Split {
1040 config: camel_api::splitter::SplitterConfig::new(split_body_lines()),
1041 steps: vec![BuilderStep::Stop],
1042 }
1043 .span_kind_hint(),
1044 SpanKindHint::Internal
1045 );
1046
1047 let enrich = |uri: &str| BuilderStep::Enrich {
1051 uri: uri.into(),
1052 strategy: None,
1053 timeout_ms: None,
1054 };
1055 let poll_enrich = |uri: &str| BuilderStep::PollEnrich {
1056 uri: uri.into(),
1057 strategy: None,
1058 timeout_ms: None,
1059 };
1060 let wire_tap = |uri: &str| BuilderStep::WireTap { uri: uri.into() };
1061
1062 assert_eq!(enrich("http://x").span_kind_hint(), SpanKindHint::Client);
1064 assert_eq!(
1065 poll_enrich("http://x").span_kind_hint(),
1066 SpanKindHint::Client
1067 );
1068 assert_eq!(wire_tap("http://x").span_kind_hint(), SpanKindHint::Client);
1069 assert_eq!(
1071 enrich("kafka:orders").span_kind_hint(),
1072 SpanKindHint::Producer
1073 );
1074 assert_eq!(
1075 poll_enrich("kafka:orders").span_kind_hint(),
1076 SpanKindHint::Producer
1077 );
1078 assert_eq!(
1079 wire_tap("kafka:orders").span_kind_hint(),
1080 SpanKindHint::Producer
1081 );
1082 assert_eq!(enrich("direct:y").span_kind_hint(), SpanKindHint::Internal);
1084 assert_eq!(
1085 poll_enrich("seda:q").span_kind_hint(),
1086 SpanKindHint::Internal
1087 );
1088 assert_eq!(
1089 wire_tap("direct:y").span_kind_hint(),
1090 SpanKindHint::Internal
1091 );
1092 }
1093
1094 #[test]
1099 fn golden_debug_output_all_variants() {
1100 use camel_api::declarative::LanguageExpressionDef;
1101 use camel_api::loop_eip::LoopMode;
1102 use camel_api::recipient_list::RecipientListConfig;
1103 use camel_api::splitter::{AggregationStrategy, StreamSplitConfig, StreamSplitFormat};
1104 use camel_api::{
1105 BoxProcessor, DynamicRouterConfig, Exchange, FilterPredicate, FunctionDefinition,
1106 FunctionId, IdentityProcessor, MulticastConfig, OpaqueProcessor, RoutingSlipConfig,
1107 Value,
1108 };
1109 use std::sync::Arc;
1110
1111 let expr = LanguageExpressionDef {
1112 language: "simple".into(),
1113 source: "${body}".into(),
1114 };
1115
1116 assert_eq!(format!("{:?}", BuilderStep::Stop), "Stop");
1119 assert_eq!(
1120 format!(
1121 "{:?}",
1122 BuilderStep::Processor(OpaqueProcessor(BoxProcessor::new(IdentityProcessor)))
1123 ),
1124 "Processor(BoxProcessor(...))"
1125 );
1126 assert_eq!(
1127 format!("{:?}", BuilderStep::To("mock:out".into())),
1128 "To(\"mock:out\")"
1129 );
1130
1131 assert_eq!(
1134 format!(
1135 "{:?}",
1136 BuilderStep::Log {
1137 level: camel_processor::LogLevel::Info,
1138 message: "hello".into(),
1139 }
1140 ),
1141 "Log { level: Info, message: \"hello\" }"
1142 );
1143
1144 assert_eq!(
1145 format!(
1146 "{:?}",
1147 BuilderStep::DeclarativeSetHeader {
1148 key: "k".into(),
1149 value: ValueSourceDef::Literal(Value::String("v".into())),
1150 }
1151 ),
1152 "DeclarativeSetHeader { key: \"k\", value: Literal(String(\"v\")) }"
1153 );
1154
1155 assert_eq!(
1156 format!(
1157 "{:?}",
1158 BuilderStep::DeclarativeSetHeaderIfAbsent {
1159 key: "k".into(),
1160 value: ValueSourceDef::Literal(Value::String("v".into())),
1161 }
1162 ),
1163 "DeclarativeSetHeaderIfAbsent { key: \"k\", value: Literal(String(\"v\")) }"
1164 );
1165
1166 assert_eq!(
1167 format!(
1168 "{:?}",
1169 BuilderStep::DeclarativeSetBody {
1170 value: ValueSourceDef::Literal(Value::String("v".into())),
1171 }
1172 ),
1173 "DeclarativeSetBody { value: Literal(String(\"v\")) }"
1174 );
1175
1176 assert_eq!(
1177 format!(
1178 "{:?}",
1179 BuilderStep::DeclarativeSetProperty {
1180 key: "prop".into(),
1181 value_source: ValueSourceDef::Literal(Value::String("v".into())),
1182 }
1183 ),
1184 "DeclarativeSetProperty { key: \"prop\", value_source: Literal(String(\"v\")) }"
1185 );
1186
1187 assert_eq!(
1188 format!(
1189 "{:?}",
1190 BuilderStep::DeclarativeScript {
1191 expression: expr.clone(),
1192 }
1193 ),
1194 "DeclarativeScript { expression: LanguageExpressionDef { language: \"simple\", source: \"${body}\" } }"
1195 );
1196
1197 let func_def = FunctionDefinition {
1199 id: FunctionId("test-id".into()),
1200 runtime: "my_runtime".into(),
1201 source: "${body}".into(),
1202 timeout_ms: 5000,
1203 route_id: None,
1204 step_index: None,
1205 };
1206 assert_eq!(
1207 format!(
1208 "{:?}",
1209 BuilderStep::DeclarativeFunction {
1210 definition: func_def,
1211 }
1212 ),
1213 "DeclarativeFunction { definition: FunctionDefinition { id: FunctionId(\"test-id\"), runtime: \"my_runtime\", source: \"${body}\", timeout_ms: 5000, route_id: None, step_index: None } }"
1214 );
1215
1216 assert_eq!(
1217 format!(
1218 "{:?}",
1219 BuilderStep::WireTap {
1220 uri: "mock:tap".into(),
1221 }
1222 ),
1223 "WireTap { uri: \"mock:tap\" }"
1224 );
1225
1226 assert_eq!(
1227 format!(
1228 "{:?}",
1229 BuilderStep::DeclarativeLog {
1230 level: camel_processor::LogLevel::Info,
1231 message: ValueSourceDef::Expression(expr.clone()),
1232 }
1233 ),
1234 "DeclarativeLog { level: Info, message: Expression(LanguageExpressionDef { language: \"simple\", source: \"${body}\" }) }"
1235 );
1236
1237 assert_eq!(
1238 format!(
1239 "{:?}",
1240 BuilderStep::Bean {
1241 name: "myBean".into(),
1242 method: "process".into(),
1243 }
1244 ),
1245 "Bean { name: \"myBean\", method: \"process\" }"
1246 );
1247
1248 assert_eq!(
1249 format!(
1250 "{:?}",
1251 BuilderStep::Script {
1252 language: "js".into(),
1253 script: "body".into(),
1254 }
1255 ),
1256 "Script { language: \"js\", script: \"body\" }"
1257 );
1258
1259 assert_eq!(
1260 format!(
1261 "{:?}",
1262 BuilderStep::Aggregate {
1263 config: camel_api::AggregatorConfig::correlate_by("id")
1264 .complete_when_size(1)
1265 .build()
1266 .unwrap(),
1267 }
1268 ),
1269 "Aggregate { config: AggregatorConfig { header_name: \"id\", completion: Single(Size(1)), correlation: HeaderName(\"id\"), strategy: CollectAll, max_buckets: Some(10000), max_bucket_size: Some(10000), bucket_ttl: Some(300s), force_completion_on_stop: false, discard_on_timeout: false, max_timeout_tasks: 1024 } }"
1270 );
1271
1272 assert_eq!(
1273 format!(
1274 "{:?}",
1275 BuilderStep::DynamicRouter {
1276 config: DynamicRouterConfig::new(Arc::new(|_: &Exchange| Some(
1277 "mock:dr".into()
1278 ))),
1279 }
1280 ),
1281 "DynamicRouter { config: DynamicRouterConfig { uri_delimiter: \",\", cache_size: 1000, ignore_invalid_endpoints: false, max_iterations: 1000, timeout: Some(60s) } }"
1282 );
1283
1284 assert_eq!(
1285 format!(
1286 "{:?}",
1287 BuilderStep::RoutingSlip {
1288 config: RoutingSlipConfig::new(Arc::new(|_: &Exchange| Some("mock:rs".into()))),
1289 }
1290 ),
1291 "RoutingSlip { config: RoutingSlipConfig { uri_delimiter: \",\", cache_size: 1000, ignore_invalid_endpoints: false } }"
1292 );
1293
1294 assert_eq!(
1295 format!(
1296 "{:?}",
1297 BuilderStep::RecipientList {
1298 config: RecipientListConfig::new(Arc::new(|_: &Exchange| String::new())),
1299 }
1300 ),
1301 "RecipientList { config: RecipientListConfig { delimiter: \",\", parallel: false, parallel_limit: None, stop_on_exception: false, max_recipients: 1000 } }"
1302 );
1303
1304 assert_eq!(
1305 format!(
1306 "{:?}",
1307 BuilderStep::Enrich {
1308 uri: "mock:enrich".into(),
1309 strategy: Some("agg".into()),
1310 timeout_ms: Some(1000),
1311 }
1312 ),
1313 "Enrich { uri: \"mock:enrich\", strategy: Some(\"agg\"), timeout_ms: Some(1000) }"
1314 );
1315
1316 assert_eq!(
1317 format!(
1318 "{:?}",
1319 BuilderStep::PollEnrich {
1320 uri: "mock:poll".into(),
1321 strategy: None,
1322 timeout_ms: None,
1323 }
1324 ),
1325 "PollEnrich { uri: \"mock:poll\", strategy: None, timeout_ms: None }"
1326 );
1327
1328 assert_eq!(
1329 format!(
1330 "{:?}",
1331 BuilderStep::Validate {
1332 predicate: expr.clone(),
1333 }
1334 ),
1335 "Validate { predicate: LanguageExpressionDef { language: \"simple\", source: \"${body}\" } }"
1336 );
1337
1338 assert_eq!(
1339 format!("{:?}", BuilderStep::Sampling { period: 100 }),
1340 "Sampling { period: 100 }"
1341 );
1342
1343 assert_eq!(
1344 format!(
1345 "{:?}",
1346 BuilderStep::Resequence {
1347 policy_config: Default::default(),
1348 }
1349 ),
1350 "Resequence { policy_config: ResequencePolicyConfig { mode: Batch { correlation: \"header.id\", sort: \"header.id\", completion: SizeOrTimeout(100, 30000) } } }"
1351 );
1352
1353 assert_eq!(
1355 format!(
1356 "{:?}",
1357 BuilderStep::DeclarativeFilter {
1358 predicate: expr.clone(),
1359 steps: vec![BuilderStep::Stop],
1360 }
1361 ),
1362 "DeclarativeFilter { predicate: LanguageExpressionDef { language: \"simple\", source: \"${body}\" }, steps: [Stop] }"
1363 );
1364
1365 assert_eq!(
1366 format!(
1367 "{:?}",
1368 BuilderStep::DeclarativeSplit {
1369 expression: expr.clone(),
1370 aggregation: AggregationStrategy::Original,
1371 parallel: false,
1372 parallel_limit: Some(2),
1373 trace_item_threshold: None,
1374 stop_on_exception: true,
1375 steps: vec![BuilderStep::Stop],
1376 }
1377 ),
1378 "DeclarativeSplit { expression: LanguageExpressionDef { language: \"simple\", source: \"${body}\" }, aggregation: Original, parallel: false, parallel_limit: Some(2), trace_item_threshold: None, stop_on_exception: true, steps: [Stop] }"
1379 );
1380
1381 assert_eq!(
1382 format!(
1383 "{:?}",
1384 BuilderStep::Split {
1385 config: camel_api::splitter::SplitterConfig::new(
1386 camel_api::splitter::split_body_lines()
1387 ),
1388 steps: vec![BuilderStep::Stop],
1389 }
1390 ),
1391 "Split { config: SplitterConfig { expression: \"<split-expression>\", aggregation: LastWins, parallel: false, parallel_limit: None, stop_on_exception: true, max_fragments: 100000, trace_item_threshold: 100 }, steps: [Stop] }"
1392 );
1393
1394 assert_eq!(
1395 format!(
1396 "{:?}",
1397 BuilderStep::Filter {
1398 predicate: FilterPredicate::new(|_: &Exchange| true),
1399 steps: vec![BuilderStep::Stop],
1400 }
1401 ),
1402 "Filter { predicate: FilterPredicate(..), steps: [Stop] }"
1403 );
1404
1405 assert_eq!(
1406 format!(
1407 "{:?}",
1408 BuilderStep::Throttle {
1409 config: camel_api::ThrottlerConfig::new(
1410 10,
1411 std::time::Duration::from_millis(10)
1412 ),
1413 steps: vec![BuilderStep::Stop],
1414 }
1415 ),
1416 "Throttle { config: ThrottlerConfig { max_requests: 10, period: 10ms, strategy: Delay }, steps: [Stop] }"
1417 );
1418
1419 assert_eq!(
1420 format!(
1421 "{:?}",
1422 BuilderStep::LoadBalance {
1423 config: camel_api::LoadBalancerConfig::round_robin(),
1424 steps: vec![BuilderStep::To("mock:l1".into())],
1425 }
1426 ),
1427 "LoadBalance { config: LoadBalancerConfig { strategy: RoundRobin }, steps: [To(\"mock:l1\")] }"
1428 );
1429
1430 assert_eq!(
1431 format!(
1432 "{:?}",
1433 BuilderStep::Delay {
1434 config: camel_api::DelayConfig::new(500),
1435 }
1436 ),
1437 "Delay { config: DelayConfig { delay_ms: 500, dynamic_header: None, max_delay_ms: 3600000 } }"
1438 );
1439
1440 assert_eq!(
1444 format!(
1445 "{:?}",
1446 BuilderStep::Choice {
1447 whens: vec![WhenStep {
1448 predicate: FilterPredicate::new(|_: &Exchange| true),
1449 steps: vec![BuilderStep::To("mock:a".into())],
1450 }],
1451 otherwise: None,
1452 }
1453 ),
1454 "Choice { whens: [WhenStep { predicate: FilterPredicate(..), steps: [To(\"mock:a\")] }], otherwise: None }"
1455 );
1456
1457 assert_eq!(
1458 format!(
1459 "{:?}",
1460 BuilderStep::DeclarativeChoice {
1461 whens: vec![DeclarativeWhenStep {
1462 predicate: expr.clone(),
1463 steps: vec![BuilderStep::Stop],
1464 }],
1465 otherwise: Some(vec![BuilderStep::Stop]),
1466 }
1467 ),
1468 "DeclarativeChoice { whens: [DeclarativeWhenStep { predicate: LanguageExpressionDef { language: \"simple\", source: \"${body}\" }, steps: [Stop] }], otherwise: Some([Stop]) }"
1469 );
1470
1471 assert_eq!(
1473 format!(
1474 "{:?}",
1475 BuilderStep::Multicast {
1476 steps: vec![BuilderStep::To("direct:a".into())],
1477 config: MulticastConfig::new(),
1478 }
1479 ),
1480 "Multicast { steps: [To(\"direct:a\")], config: MulticastConfig { parallel: false, parallel_limit: None, stop_on_exception: false, timeout: None, aggregation: LastWins } }"
1481 );
1482
1483 assert_eq!(
1485 format!(
1486 "{:?}",
1487 BuilderStep::DeclarativeDynamicRouter {
1488 expression: expr.clone(),
1489 uri_delimiter: ",".into(),
1490 cache_size: 1000,
1491 ignore_invalid_endpoints: false,
1492 max_iterations: 1000,
1493 }
1494 ),
1495 "DeclarativeDynamicRouter { expression: LanguageExpressionDef { language: \"simple\", source: \"${body}\" }, uri_delimiter: \",\", cache_size: 1000, ignore_invalid_endpoints: false, max_iterations: 1000 }"
1496 );
1497
1498 assert_eq!(
1499 format!(
1500 "{:?}",
1501 BuilderStep::DeclarativeRoutingSlip {
1502 expression: expr.clone(),
1503 uri_delimiter: ",".into(),
1504 cache_size: 1000,
1505 ignore_invalid_endpoints: false,
1506 }
1507 ),
1508 "DeclarativeRoutingSlip { expression: LanguageExpressionDef { language: \"simple\", source: \"${body}\" }, uri_delimiter: \",\", cache_size: 1000, ignore_invalid_endpoints: false }"
1509 );
1510
1511 assert_eq!(
1513 format!(
1514 "{:?}",
1515 BuilderStep::DeclarativeRecipientList {
1516 expression: expr.clone(),
1517 delimiter: ",".into(),
1518 parallel: false,
1519 parallel_limit: None,
1520 stop_on_exception: false,
1521 aggregation: "original".into(),
1522 }
1523 ),
1524 "DeclarativeRecipientList { expression: LanguageExpressionDef { language: \"simple\", source: \"${body}\" }, delimiter: \",\", parallel: false, parallel_limit: None, stop_on_exception: false, aggregation: \"original\" }"
1525 );
1526
1527 assert_eq!(
1529 format!(
1530 "{:?}",
1531 BuilderStep::Loop {
1532 config: camel_api::loop_eip::LoopConfig::new(LoopMode::Count(3)),
1533 steps: vec![],
1534 }
1535 ),
1536 "Loop { config: LoopConfig { mode: Count(3), max_iterations: 10000 }, steps: [] }"
1537 );
1538
1539 assert_eq!(
1540 format!(
1541 "{:?}",
1542 BuilderStep::DeclarativeLoop {
1543 count: Some(5),
1544 while_predicate: None,
1545 steps: vec![],
1546 max_iterations: Some(100),
1547 }
1548 ),
1549 "DeclarativeLoop { count: Some(5), while_predicate: None, steps: [], max_iterations: Some(100) }"
1550 );
1551
1552 assert_eq!(
1554 format!(
1555 "{:?}",
1556 BuilderStep::ClaimCheck {
1557 repository: "myRepo".into(),
1558 operation: "checkout".into(),
1559 key: expr.clone(),
1560 filter: None,
1561 }
1562 ),
1563 "ClaimCheck { repository: \"myRepo\", operation: \"checkout\", key: LanguageExpressionDef { language: \"simple\", source: \"${body}\" }, filter: None }"
1564 );
1565
1566 assert_eq!(
1568 format!(
1569 "{:?}",
1570 BuilderStep::Sort {
1571 expression: expr.clone(),
1572 reverse: false,
1573 }
1574 ),
1575 "Sort { expression: LanguageExpressionDef { language: \"simple\", source: \"${body}\" }, reverse: false }"
1576 );
1577
1578 assert_eq!(
1580 format!(
1581 "{:?}",
1582 BuilderStep::IdempotentConsumer {
1583 repository: "myRepo".into(),
1584 expression: expr.clone(),
1585 steps: vec![],
1586 eager: true,
1587 remove_on_failure: false,
1588 }
1589 ),
1590 "IdempotentConsumer { repository: \"myRepo\", expression: LanguageExpressionDef { language: \"simple\", source: \"${body}\" }, steps: [], eager: true, remove_on_failure: false }"
1591 );
1592
1593 assert_eq!(
1595 format!(
1596 "{:?}",
1597 BuilderStep::DeclarativeDoTry {
1598 try_steps: vec![BuilderStep::Stop],
1599 catch: vec![],
1600 finally: None,
1601 }
1602 ),
1603 "DeclarativeDoTry { try_steps: [Stop], catch: [], finally: None }"
1604 );
1605
1606 assert_eq!(
1608 format!(
1609 "{:?}",
1610 BuilderStep::DeclarativeStreamSplit {
1611 stream_config: StreamSplitConfig {
1612 format: StreamSplitFormat::Ndjson,
1613 max_record_bytes: 1024 * 1024,
1614 batch_size: 1,
1615 chunk_size: None,
1616 include_origin: true,
1617 },
1618 aggregation: AggregationStrategy::Original,
1619 stop_on_exception: true,
1620 steps: vec![BuilderStep::Stop],
1621 }
1622 ),
1623 "DeclarativeStreamSplit { stream_config: StreamSplitConfig { format: Ndjson, max_record_bytes: 1048576, batch_size: 1, chunk_size: None, include_origin: true }, aggregation: Original, stop_on_exception: true, steps: [Stop] }"
1624 );
1625 }
1626
1627 #[test]
1628 fn test_builder_step_multicast_variant() {
1629 use camel_api::MulticastConfig;
1630
1631 let step = BuilderStep::Multicast {
1632 steps: vec![BuilderStep::To("direct:a".into())],
1633 config: MulticastConfig::new(),
1634 };
1635
1636 assert!(matches!(step, BuilderStep::Multicast { .. }));
1637 }
1638
1639 #[test]
1640 fn test_route_definition_defaults() {
1641 let def = RouteDefinition::new("direct:test", vec![]).with_route_id("test-route");
1642 assert_eq!(def.route_id(), "test-route");
1643 assert!(def.auto_startup());
1644 assert_eq!(def.startup_order(), 1000);
1645 }
1646
1647 #[test]
1648 fn test_route_definition_builders() {
1649 let def = RouteDefinition::new("direct:test", vec![])
1650 .with_route_id("my-route")
1651 .with_auto_startup(false)
1652 .with_startup_order(50);
1653 assert_eq!(def.route_id(), "my-route");
1654 assert!(!def.auto_startup());
1655 assert_eq!(def.startup_order(), 50);
1656 }
1657
1658 #[test]
1659 fn test_route_definition_accessors_cover_core_fields() {
1660 let def = RouteDefinition::new("direct:in", vec![BuilderStep::To("mock:out".into())])
1661 .with_route_id("accessor-route");
1662
1663 assert_eq!(def.from_uri(), "direct:in");
1664 assert_eq!(def.steps().len(), 1);
1665 assert!(matches!(def.steps()[0], BuilderStep::To(_)));
1666 }
1667
1668 #[test]
1669 fn test_route_definition_error_handler_circuit_breaker_and_concurrency_accessors() {
1670 use camel_api::circuit_breaker::CircuitBreakerConfig;
1671 use camel_api::error_handler::ErrorHandlerConfig;
1672 use camel_component_api::ConcurrencyModel;
1673
1674 let def = RouteDefinition::new("direct:test", vec![])
1675 .with_route_id("eh-route")
1676 .with_error_handler(ErrorHandlerConfig::dead_letter_channel("log:dlc"))
1677 .with_circuit_breaker(CircuitBreakerConfig::new())
1678 .with_concurrency(ConcurrencyModel::Concurrent { max: Some(4) });
1679
1680 let eh = def
1681 .error_handler_config()
1682 .expect("error handler should be set");
1683 assert_eq!(eh.dlc_uri.as_deref(), Some("log:dlc"));
1684 assert!(def.circuit_breaker_config().is_some());
1685 assert!(matches!(
1686 def.concurrency_override(),
1687 Some(ConcurrencyModel::Concurrent { max: Some(4) })
1688 ));
1689 }
1690
1691 #[test]
1692 fn test_builder_step_debug_covers_many_variants() {
1693 use camel_api::splitter::{AggregationStrategy, SplitterConfig, split_body_lines};
1694 use camel_api::{
1695 BoxProcessor, DynamicRouterConfig, Exchange, FilterPredicate, IdentityProcessor,
1696 OpaqueProcessor, RoutingSlipConfig, Value,
1697 };
1698 use std::sync::Arc;
1699
1700 let expr = LanguageExpressionDef {
1701 language: "simple".into(),
1702 source: "${body}".into(),
1703 };
1704
1705 let steps = vec![
1706 BuilderStep::Processor(OpaqueProcessor(BoxProcessor::new(IdentityProcessor))),
1707 BuilderStep::To("mock:out".into()),
1708 BuilderStep::Stop,
1709 BuilderStep::Log {
1710 level: camel_processor::LogLevel::Info,
1711 message: "hello".into(),
1712 },
1713 BuilderStep::DeclarativeSetHeader {
1714 key: "k".into(),
1715 value: ValueSourceDef::Literal(Value::String("v".into())),
1716 },
1717 BuilderStep::DeclarativeSetBody {
1718 value: ValueSourceDef::Expression(expr.clone()),
1719 },
1720 BuilderStep::DeclarativeFilter {
1721 predicate: expr.clone(),
1722 steps: vec![BuilderStep::Stop],
1723 },
1724 BuilderStep::DeclarativeChoice {
1725 whens: vec![DeclarativeWhenStep {
1726 predicate: expr.clone(),
1727 steps: vec![BuilderStep::Stop],
1728 }],
1729 otherwise: Some(vec![BuilderStep::Stop]),
1730 },
1731 BuilderStep::DeclarativeScript {
1732 expression: expr.clone(),
1733 },
1734 BuilderStep::DeclarativeSplit {
1735 expression: expr.clone(),
1736 aggregation: AggregationStrategy::Original,
1737 parallel: false,
1738 parallel_limit: Some(2),
1739 trace_item_threshold: None,
1740 stop_on_exception: true,
1741 steps: vec![BuilderStep::Stop],
1742 },
1743 BuilderStep::Split {
1744 config: SplitterConfig::new(split_body_lines()),
1745 steps: vec![BuilderStep::Stop],
1746 },
1747 BuilderStep::Aggregate {
1748 config: camel_api::AggregatorConfig::correlate_by("id")
1749 .complete_when_size(1)
1750 .build()
1751 .unwrap(),
1752 },
1753 BuilderStep::Filter {
1754 predicate: FilterPredicate::new(|_: &Exchange| true),
1755 steps: vec![BuilderStep::Stop],
1756 },
1757 BuilderStep::WireTap {
1758 uri: "mock:tap".into(),
1759 },
1760 BuilderStep::DeclarativeLog {
1761 level: camel_processor::LogLevel::Info,
1762 message: ValueSourceDef::Expression(expr.clone()),
1763 },
1764 BuilderStep::Bean {
1765 name: "bean".into(),
1766 method: "call".into(),
1767 },
1768 BuilderStep::Script {
1769 language: "rhai".into(),
1770 script: "body".into(),
1771 },
1772 BuilderStep::Throttle {
1773 config: camel_api::ThrottlerConfig::new(10, std::time::Duration::from_millis(10)),
1774 steps: vec![BuilderStep::Stop],
1775 },
1776 BuilderStep::LoadBalance {
1777 config: camel_api::LoadBalancerConfig::round_robin(),
1778 steps: vec![BuilderStep::To("mock:l1".into())],
1779 },
1780 BuilderStep::DynamicRouter {
1781 config: DynamicRouterConfig::new(Arc::new(|_| Some("mock:dr".into()))),
1782 },
1783 BuilderStep::RoutingSlip {
1784 config: RoutingSlipConfig::new(Arc::new(|_| Some("mock:rs".into()))),
1785 },
1786 ];
1787
1788 for step in steps {
1789 let dbg = format!("{step:?}");
1790 assert!(!dbg.is_empty());
1791 }
1792 }
1793
1794 #[test]
1795 fn test_route_definition_to_info_preserves_metadata() {
1796 let info = RouteDefinition::new("direct:test", vec![])
1797 .with_route_id("meta-route")
1798 .with_auto_startup(false)
1799 .with_startup_order(7)
1800 .to_info();
1801
1802 assert_eq!(info.route_id(), "meta-route");
1803 assert!(!info.auto_startup());
1804 assert_eq!(info.startup_order(), 7);
1805 }
1806
1807 #[test]
1808 fn test_choice_builder_step_debug() {
1809 use camel_api::FilterPredicate;
1810
1811 fn always_true(_: &camel_api::Exchange) -> bool {
1812 true
1813 }
1814
1815 let step = BuilderStep::Choice {
1816 whens: vec![WhenStep {
1817 predicate: FilterPredicate::new(always_true),
1818 steps: vec![BuilderStep::To("mock:a".into())],
1819 }],
1820 otherwise: None,
1821 };
1822 let debug = format!("{step:?}");
1823 assert!(debug.contains("Choice"));
1824 }
1825
1826 #[test]
1827 fn test_route_definition_unit_of_work() {
1828 use camel_api::UnitOfWorkConfig;
1829 let config = UnitOfWorkConfig {
1830 on_complete: Some("log:complete".into()),
1831 on_failure: Some("log:failed".into()),
1832 };
1833 let def = RouteDefinition::new("direct:test", vec![])
1834 .with_route_id("uow-test")
1835 .with_unit_of_work(config.clone());
1836 assert_eq!(
1837 def.unit_of_work_config().unwrap().on_complete.as_deref(),
1838 Some("log:complete")
1839 );
1840 assert_eq!(
1841 def.unit_of_work_config().unwrap().on_failure.as_deref(),
1842 Some("log:failed")
1843 );
1844
1845 let def_no_uow = RouteDefinition::new("direct:test", vec![]).with_route_id("no-uow");
1846 assert!(def_no_uow.unit_of_work_config().is_none());
1847 }
1848
1849 #[test]
1850 fn test_route_definition_security_policy_accessor() {
1851 use async_trait::async_trait;
1852 use camel_api::CamelError;
1853 use camel_api::Exchange;
1854 use camel_api::security_policy::{
1855 AuthContext, AuthorizationDecision, Principal, SecurityPolicy, SecurityPolicyConfig,
1856 };
1857
1858 struct StubPolicy;
1859 #[async_trait]
1860 impl SecurityPolicy for StubPolicy {
1861 async fn evaluate(
1862 &self,
1863 _exchange: &mut Exchange,
1864 _auth: &AuthContext<'_>,
1865 ) -> Result<AuthorizationDecision, CamelError> {
1866 Ok(AuthorizationDecision::Granted {
1867 principal: Principal {
1868 subject: "test".into(),
1869 issuer: "test".into(),
1870 audience: vec![],
1871 scopes: vec![],
1872 roles: vec![],
1873 claims: serde_json::Value::Null,
1874 },
1875 })
1876 }
1877 }
1878
1879 let def_no_sp = RouteDefinition::new("direct:test", vec![]).with_route_id("no-sp");
1880 assert!(def_no_sp.security_policy_config().is_none());
1881
1882 let def = RouteDefinition::new("direct:test", vec![])
1883 .with_route_id("sp-test")
1884 .with_security_policy(SecurityPolicyConfig::new(StubPolicy));
1885 assert!(def.security_policy_config().is_some());
1886 }
1887
1888 #[test]
1889 fn test_route_definition_security_authenticator_accessor() {
1890 use camel_api::security_policy::Principal;
1891
1892 struct TestAuth;
1893 #[async_trait::async_trait]
1894 impl TokenAuthenticator for TestAuth {
1895 async fn authenticate_bearer(
1896 &self,
1897 _token: &str,
1898 ) -> Result<Principal, camel_api::CamelError> {
1899 Ok(Principal {
1900 subject: "test".into(),
1901 issuer: "test".into(),
1902 audience: vec![],
1903 scopes: vec![],
1904 roles: vec![],
1905 claims: serde_json::Value::Null,
1906 })
1907 }
1908 }
1909
1910 let def_no_auth = RouteDefinition::new("direct:test".to_string(), vec![]);
1911 assert!(def_no_auth.security_authenticator().is_none());
1912
1913 let auth = Arc::new(TestAuth);
1914 let def = RouteDefinition::new("direct:test".to_string(), vec![])
1915 .with_security_authenticator(auth);
1916 assert!(def.security_authenticator().is_some());
1917 }
1918
1919 #[test]
1920 fn test_map_steps_swaps_steps_and_preserves_other_fields() {
1921 let original = RouteDefinition::new(
1922 "direct:test".to_string(),
1923 vec![BuilderStep::To("mock:a".into()), BuilderStep::Stop],
1924 )
1925 .with_route_id("my-route");
1926
1927 let mapped = original.map_steps(|steps| {
1928 let mut out = Vec::with_capacity(steps.len() + 1);
1929 out.push(BuilderStep::To("mock:prefix".into()));
1930 out.extend(steps);
1931 out
1932 });
1933
1934 assert_eq!(mapped.steps().len(), 3);
1936 assert!(matches!(mapped.steps()[0], BuilderStep::To(ref s) if s == "mock:prefix"));
1937 assert!(matches!(mapped.steps()[1], BuilderStep::To(ref s) if s == "mock:a"));
1938 assert_eq!(mapped.route_id(), "my-route");
1940 }
1941
1942 #[test]
1943 fn circuit_breaker_fallback_accessor_returns_steps() {
1944 let def = RouteDefinition::new("direct:start", vec![])
1945 .with_circuit_breaker_fallback(vec![BuilderStep::To("mock:out".into())]);
1946 assert_eq!(def.circuit_breaker_fallback().len(), 1);
1947 }
1948}