Skip to main content

camel_builder/
lib.rs

1//! Fluent builder API for constructing Camel routes programmatically with EIP patterns.
2//!
3//! Main types: `RouteBuilder`, `StepAccumulator`, `SplitBuilder`, `ChoiceBuilder`, `MulticastBuilder`,
4//! `ThrottleBuilder`, `LoopBuilder`, `LoadBalancerBuilder`, `OnExceptionBuilder`.
5
6use 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
40// ── Module declarations ─────────────────────────────────────────────────────
41pub mod do_try;
42pub use do_try::{DoCatchBuilder, DoFinallyBuilder, DoTryBuilder};
43
44/// Shared step-accumulation methods for all builder types.
45///
46/// Implementors provide `steps_mut()` and get step-adding methods for free.
47/// `filter()` and other branching methods are NOT included — they return
48/// different types per builder and stay as per-builder methods.
49pub 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    /// Apache Camel-compatible alias for [`set_body`](Self::set_body).
111    ///
112    /// Transforms the message body using the given value. Semantically identical
113    /// to `set_body` — provided for familiarity with Apache Camel route DSLs.
114    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    /// Stop processing this exchange immediately. No further steps in the
151    /// current pipeline will run.
152    ///
153    /// Can be used at any point in the route: directly on RouteBuilder,
154    /// inside `.filter()`, inside `.split()`, etc.
155    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    /// Log a message at the specified level.
179    ///
180    /// The message will be logged when an exchange passes through this step.
181    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    /// Convert the message body to the target type.
190    ///
191    /// Supported: Text ↔ Json ↔ Bytes. `Body::Stream` always fails.
192    /// Returns `TypeConversionFailed` if conversion is not possible.
193    ///
194    /// # Example
195    /// ```ignore
196    /// route.set_body(Value::String(r#"{"x":1}"#.into()))
197    ///      .convert_body_to(BodyType::Json)
198    ///      .to("direct:next")
199    /// ```
200    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    /// Materialize `Body::Stream` into `Body::Bytes` using the default threshold (128 KB).
220    ///
221    /// Equivalent to `.stream_cache(camel_api::stream_cache::DEFAULT_STREAM_CACHE_THRESHOLD)`.
222    fn stream_cache_default(self) -> Self {
223        self.stream_cache(camel_api::stream_cache::DEFAULT_STREAM_CACHE_THRESHOLD)
224    }
225
226    /// Marshal the message body using the specified data format.
227    ///
228    /// Supported formats: `"json"`, `"xml"`, `"csv"`, `"zip"`, `"tar"`, `"gzip"`, `"tar.gz"`.
229    /// Returns `Err(CamelError::Config)` if the format name is unknown.
230    /// Converts a structured body (e.g., `Body::Json`) to a wire-format body (e.g., `Body::Text`).
231    ///
232    /// # Example
233    /// ```ignore
234    /// route.marshal("json")?.to("direct:next")
235    /// ```
236    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    /// Unmarshal the message body using the specified data format.
249    ///
250    /// Supported formats: `"json"`, `"xml"`, `"csv"`, `"zip"`, `"tar"`, `"gzip"`, `"tar.gz"`.
251    /// Returns `Err(CamelError::Config)` if the format name is unknown.
252    /// Converts a wire-format body (e.g., `Body::Text`) to a structured body (e.g., `Body::Json`).
253    ///
254    /// # Example
255    /// ```ignore
256    /// route.unmarshal("json")?.to("direct:next")
257    /// ```
258    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    /// Validate the exchange using a predicate expression.
271    ///
272    /// If the expression evaluates to `true`, the exchange continues.
273    /// If `false`, a `CamelError::ValidationError` is returned into the route error handler.
274    ///
275    /// # Example
276    /// ```ignore
277    /// route.validate("${body.size()} > 0").to("direct:out")
278    /// ```
279    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    /// Execute a script that can modify the exchange (headers, properties, body).
292    ///
293    /// The script has access to `headers`, `properties`, and `body` variables
294    /// and can modify them with assignment syntax: `headers["k"] = v`.
295    ///
296    /// # Example
297    /// ```ignore
298    /// // ignore: requires full CamelContext setup with registered language
299    /// route.script("rhai", r#"headers["tenant"] = "acme"; body = body + "_processed""#)
300    /// ```
301    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    /// EIP-7 enrich: synchronous content enrichment via a resolved producer.
310    ///
311    /// Calls the given endpoint URI as a producer, then merges the response
312    /// back into the original exchange body using the default `UseEnrichedBody`
313    /// strategy (original headers/properties are preserved).
314    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    /// EIP-7 pollEnrich: blocking poll of a PollingConsumer with timeout.
324    ///
325    /// Reads from a polling endpoint (e.g., file) and merges the result
326    /// into the exchange body using the default `UseEnrichedBody` strategy.
327    /// `timeout_ms` controls how long to wait for data.
328    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/// Identifies the endpoint slot a `.parameters()` call attaches to.
347#[derive(Debug, Clone, PartialEq, Eq)]
348enum EndpointSlot {
349    /// The route's `from` endpoint.
350    From,
351    /// A step in `RouteBuilder::steps`, by index.
352    Step(usize),
353}
354
355/// A fluent builder for constructing routes.
356///
357/// # Example
358///
359/// ```ignore
360/// let definition = RouteBuilder::from("timer:tick?period=1000")
361///     .set_header("source", Value::String("timer".into()))
362///     .filter(|ex| ex.input.body.as_text().is_some())
363///     .to("log:info?showHeaders=true")
364///     .build()?;
365/// ```
366/// `RouteBuilder` is `Clone`: a partially-built route can be cloned and reused as a
367/// template for multiple routes (mirrors Apache Camel's cloneable `RouteBuilder`).
368/// Clone is a deep copy of the step list; step closures live behind `Arc`/`BoxProcessor`
369/// so cloning shares the closure and duplicates only the light wrapper (rc-8m5o).
370#[derive(Clone)]
371pub struct RouteBuilder {
372    from_uri: String,
373    steps: Vec<BuilderStep>,
374    /// Pending `.parameters()` maps, each attached to a specific endpoint slot.
375    parameter_assignments: Vec<(EndpointSlot, BTreeMap<String, String>)>,
376    /// Misuse recorded at `.parameters()` call time and surfaced at `build()`.
377    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    /// Start building a route from the given source endpoint URI.
411    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    /// Open a filter scope. Only exchanges matching `predicate` will be processed
431    /// by the steps inside the scope. Non-matching exchanges skip the scope entirely
432    /// and continue to steps after `.end_filter()`.
433    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    /// Open a choice scope for content-based routing.
445    ///
446    /// Within the choice, you can define multiple `.when()` clauses and an
447    /// optional `.otherwise()` clause. The first matching `when` predicate
448    /// determines which sub-pipeline executes.
449    pub fn choice(self) -> ChoiceBuilder {
450        ChoiceBuilder {
451            parent: self,
452            whens: vec![],
453            _otherwise: None,
454        }
455    }
456
457    /// Add a WireTap step that sends a clone of the exchange to the given
458    /// endpoint URI (fire-and-forget). The original exchange continues
459    /// downstream unchanged.
460    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    /// Attach a `parameters:` map to the most recent endpoint slot — the `from`
468    /// endpoint when called before any step, otherwise the last added step.
469    ///
470    /// Parameters on different endpoints each persist independently; none
471    /// overwrites another. Misuse (a second `.parameters()` on the same slot,
472    /// or a call with no pending endpoint step) is deferred and surfaced as a
473    /// `CamelError::RouteError` at `build()` / `build_canonical()`, never a panic.
474    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    /// Set a per-route error handler. Overrides the global error handler on `CamelContext`.
498    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    /// Set a dead letter channel URI for shorthand error handler mode.
510    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    /// Add a shorthand exception policy scope. Call `.end_on_exception()` to return to route builder.
527    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    /// Set a circuit breaker for this route.
551    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    /// Named provider registry for security plan compilation (ADR-0061):
573    /// routes declaring security resolve their provider here; staging
574    /// aborts when the registry cannot satisfy the declaration.
575    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    /// Override the consumer's default concurrency model.
584    ///
585    /// When set, the pipeline spawns a task per exchange, processing them
586    /// concurrently. `max` limits the number of simultaneously active
587    /// pipeline executions (0 = unbounded, channel buffer is backpressure).
588    ///
589    /// # Example
590    /// ```ignore
591    /// RouteBuilder::from("http://0.0.0.0:8080/api")
592    ///     .concurrent(16)  // max 16 in-flight pipeline executions
593    ///     .process(handle_request)
594    ///     .build()
595    /// ```
596    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    /// Force sequential processing, overriding a concurrent-capable consumer.
603    ///
604    /// Useful for HTTP routes that mutate shared state and need ordering
605    /// guarantees.
606    pub fn sequential(mut self) -> Self {
607        self.concurrency = Some(ConcurrencyModel::Sequential);
608        self
609    }
610
611    /// Set the route ID for this route.
612    ///
613    /// If not set, the route will be assigned an auto-generated ID.
614    pub fn route_id(mut self, id: impl Into<String>) -> Self {
615        self.route_id = Some(id.into());
616        self
617    }
618
619    /// Set whether this route should automatically start when the context starts.
620    ///
621    /// Default is `true`.
622    pub fn auto_startup(mut self, auto: bool) -> Self {
623        self.auto_startup = Some(auto);
624        self
625    }
626
627    /// Set the startup order for this route.
628    ///
629    /// Routes with lower values start first. Default is 1000.
630    pub fn startup_order(mut self, order: i32) -> Self {
631        self.startup_order = Some(order);
632        self
633    }
634
635    /// Begin a Splitter sub-pipeline. Steps added after this call (until
636    /// `.end_split()`) will be executed per-fragment.
637    ///
638    /// Returns a `SplitBuilder` — you cannot call `.build()` until
639    /// `.end_split()` closes the split scope (enforced by the type system).
640    pub fn split(self, config: SplitterConfig) -> SplitBuilder {
641        SplitBuilder {
642            parent: self,
643            config,
644            steps: Vec::new(),
645        }
646    }
647
648    /// Begin a Multicast sub-pipeline. Steps added after this call (until
649    /// `.end_multicast()`) will each receive a copy of the exchange.
650    ///
651    /// Returns a `MulticastBuilder` — you cannot call `.build()` until
652    /// `.end_multicast()` closes the multicast scope (enforced by the type system).
653    pub fn multicast(self) -> MulticastBuilder {
654        MulticastBuilder {
655            parent: self,
656            steps: Vec::new(),
657            config: MulticastConfig::new(),
658        }
659    }
660
661    /// Begin a Throttle sub-pipeline. Rate limits message processing to at most
662    /// `max_requests` per `period`. Steps inside the throttle scope are only
663    /// executed when the rate limit allows.
664    ///
665    /// Returns a `ThrottleBuilder` — you cannot call `.build()` until
666    /// `.end_throttle()` closes the throttle scope (enforced by the type system).
667    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    /// Begin a Loop sub-pipeline that iterates a fixed number of times.
676    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    /// Begin a Loop sub-pipeline that iterates while a predicate is true.
685    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    /// Begin a LoadBalance sub-pipeline. Distributes exchanges across multiple
697    /// endpoints using a configurable strategy (round-robin, random, weighted, failover).
698    ///
699    /// Returns a `LoadBalancerBuilder` — you cannot call `.build()` until
700    /// `.end_load_balance()` closes the load balance scope (enforced by the type system).
701    pub fn load_balance(self) -> LoadBalancerBuilder {
702        LoadBalancerBuilder {
703            parent: self,
704            config: LoadBalancerConfig::round_robin(),
705            steps: Vec::new(),
706        }
707    }
708
709    /// Add a dynamic router step that routes exchanges dynamically based on
710    /// expression evaluation at runtime.
711    ///
712    /// The expression receives the exchange and returns `Some(uri)` to route to
713    /// the next endpoint, or `None` to stop routing.
714    ///
715    /// # Example
716    /// ```ignore
717    /// RouteBuilder::from("timer:tick")
718    ///     .route_id("test-route")
719    ///     .dynamic_router(|ex| {
720    ///         ex.input.header("dest").and_then(|v| v.as_str().map(|s| s.to_string()))
721    ///     })
722    ///     .build()
723    /// ```
724    pub fn dynamic_router(self, expression: RouterExpression) -> Self {
725        self.dynamic_router_with_config(DynamicRouterConfig::new(expression))
726    }
727
728    /// Add a dynamic router step with full configuration.
729    ///
730    /// Allows customization of URI delimiter, cache size, timeout, and other options.
731    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    /// Consume the builder and produce a [`RouteDefinition`].
755    // Duplicate route IDs are detected at `CamelContext::add_route_definition` time
756    // (RouteController rejects atomically with CamelError::RouteError).
757    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        // Deferred `.parameters()` validation and merge: misuse is a RouteError,
810        // URI merge failures surface as CamelError::EndpointUri via `?`.
811        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    /// Compile this builder route into canonical spec.
864    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        // Deferred `.parameters()` validation and merge (same path as `build()`).
877        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
961/// True when `slot` names an endpoint-bearing position: the `from` endpoint, or
962/// a step of one of the endpoint kinds (`To`, `WireTap`, `Enrich`, `PollEnrich`).
963fn 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
978/// Merge each pending `.parameters()` assignment into its endpoint slot's URI.
979///
980/// A recorded misuse flag aborts first (builder misuse → `RouteError`). Each
981/// endpoint slot is then verified endpoint-bearing and its URI re-rendered via
982/// [`EndpointUri::try_from_uri_and_params`]; URI merge failures propagate as
983/// `CamelError::EndpointUri` (never folded into `RouteError`). Empty parameter
984/// maps are skipped entirely, preserving URI bytes (parity with the DSL
985/// lowering); misuse for such calls is still surfaced via `misuse`.
986fn 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        // Empty maps are no-ops: re-rendering would still route through
998        // `EndpointUri` and could normalize URI bytes, diverging from the DSL
999        // lowering which skips empty maps for byte-identity/passthrough.
1000        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
1031/// Validate that a URI is non-empty and contains a scheme component.
1032fn 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            // The programmatic language-split builder has no knob source;
1101            // the default applies at wrap time (asymmetry deferred).
1102            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        // INVARIANT: every other correlation variant (Fn and any future
1220        // variant) is rejected by the `header` match above, which returns
1221        // early — so this branch is unreachable here.
1222        _ => 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
1341/// Builder for the sub-pipeline within a `.split()` ... `.end_split()` block.
1342///
1343/// Exposes the same step methods as `RouteBuilder` (to, process, filter, etc.)
1344/// but NOT `.build()` and NOT `.split()` (no nested splits).
1345///
1346/// Calling `.end_split()` packages the sub-steps into a `BuilderStep::Split`
1347/// and returns the parent `RouteBuilder`.
1348pub struct SplitBuilder {
1349    parent: RouteBuilder,
1350    config: SplitterConfig,
1351    steps: Vec<BuilderStep>,
1352}
1353
1354impl SplitBuilder {
1355    /// Open a filter scope within the split sub-pipeline.
1356    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    /// Close the split scope. Packages the accumulated sub-steps into a
1368    /// `BuilderStep::Split` and returns the parent `RouteBuilder`.
1369    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
1385/// Builder for the sub-pipeline within a `.filter()` ... `.end_filter()` block.
1386pub struct FilterBuilder {
1387    parent: RouteBuilder,
1388    predicate: FilterPredicate,
1389    steps: Vec<BuilderStep>,
1390}
1391
1392impl FilterBuilder {
1393    /// Close the filter scope. Packages the accumulated sub-steps into a
1394    /// `BuilderStep::Filter` and returns the parent `RouteBuilder`.
1395    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
1411/// Builder for a filter scope nested inside a `.split()` block.
1412pub struct FilterInSplitBuilder {
1413    parent: SplitBuilder,
1414    predicate: FilterPredicate,
1415    steps: Vec<BuilderStep>,
1416}
1417
1418impl FilterInSplitBuilder {
1419    /// Close the filter scope and return the parent `SplitBuilder`.
1420    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
1436// ── Choice/When/Otherwise builders ─────────────────────────────────────────
1437
1438/// Builder for a `.choice()` ... `.end_choice()` block.
1439///
1440/// Accumulates `when` clauses and an optional `otherwise` clause.
1441/// Cannot call `.build()` until `.end_choice()` is called.
1442pub struct ChoiceBuilder {
1443    parent: RouteBuilder,
1444    whens: Vec<WhenStep>,
1445    _otherwise: Option<Vec<BuilderStep>>,
1446}
1447
1448impl ChoiceBuilder {
1449    /// Open a `when` clause. Only exchanges matching `predicate` will be
1450    /// processed by the steps inside the `.when()` ... `.end_when()` scope.
1451    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    /// Open an `otherwise` clause. Executed when no `when` predicate matched.
1463    ///
1464    /// Only one `otherwise` is allowed per `choice`. Call this after all `.when()` clauses.
1465    pub fn otherwise(self) -> OtherwiseBuilder {
1466        OtherwiseBuilder {
1467            parent: self,
1468            steps: vec![],
1469        }
1470    }
1471
1472    /// Close the choice scope. Packages all accumulated `when` clauses and
1473    /// optional `otherwise` into a `BuilderStep::Choice` and returns the
1474    /// parent `RouteBuilder`.
1475    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
1485/// Builder for the sub-pipeline within a `.when()` ... `.end_when()` block.
1486pub struct WhenBuilder {
1487    parent: ChoiceBuilder,
1488    predicate: camel_api::FilterPredicate,
1489    steps: Vec<BuilderStep>,
1490}
1491
1492impl WhenBuilder {
1493    /// Close the when scope. Packages the accumulated sub-steps into a
1494    /// `WhenStep` and returns the parent `ChoiceBuilder`.
1495    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
1510/// Builder for the sub-pipeline within an `.otherwise()` ... `.end_otherwise()` block.
1511pub struct OtherwiseBuilder {
1512    parent: ChoiceBuilder,
1513    steps: Vec<BuilderStep>,
1514}
1515
1516impl OtherwiseBuilder {
1517    /// Close the otherwise scope and return the parent `ChoiceBuilder`.
1518    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
1531/// Builder for the sub-pipeline within a `.multicast()` ... `.end_multicast()` block.
1532///
1533/// Exposes the same step methods as `RouteBuilder` (to, process, filter, etc.)
1534/// but NOT `.build()` and NOT `.multicast()` (no nested multicasts).
1535///
1536/// Calling `.end_multicast()` packages the sub-steps into a `BuilderStep::Multicast`
1537/// and returns the parent `RouteBuilder`.
1538pub 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
1586/// Builder for the sub-pipeline within a `.throttle()` ... `.end_throttle()` block.
1587///
1588/// Exposes the same step methods as `RouteBuilder` (to, process, filter, etc.)
1589/// but NOT `.build()` and NOT `.throttle()` (no nested throttles).
1590///
1591/// Calling `.end_throttle()` packages the sub-steps into a `BuilderStep::Throttle`
1592/// and returns the parent `RouteBuilder`.
1593pub struct ThrottleBuilder {
1594    parent: RouteBuilder,
1595    config: ThrottlerConfig,
1596    steps: Vec<BuilderStep>,
1597}
1598
1599impl ThrottleBuilder {
1600    /// Set the throttle strategy. Default is `Delay`.
1601    ///
1602    /// - `Delay`: Queue messages until capacity available
1603    /// - `Reject`: Return error immediately when throttled
1604    /// - `Drop`: Silently discard excess messages
1605    pub fn strategy(mut self, strategy: ThrottleStrategy) -> Self {
1606        self.config = self.config.strategy(strategy);
1607        self
1608    }
1609
1610    /// Close the throttle scope. Packages the accumulated sub-steps into a
1611    /// `BuilderStep::Throttle` and returns the parent `RouteBuilder`.
1612    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
1628/// Builder for the sub-pipeline within a `.loop_count()` / `.loop_while()` ... `.end_loop()` block.
1629pub 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
1694/// Builder for the sub-pipeline within a `.load_balance()` ... `.end_load_balance()` block.
1695///
1696/// Exposes the same step methods as `RouteBuilder` (to, process, filter, etc.)
1697/// but NOT `.build()` and NOT `.load_balance()` (no nested load balancers).
1698///
1699/// Calling `.end_load_balance()` packages the sub-steps into a `BuilderStep::LoadBalance`
1700/// and returns the parent `RouteBuilder`.
1701pub struct LoadBalancerBuilder {
1702    parent: RouteBuilder,
1703    config: LoadBalancerConfig,
1704    steps: Vec<BuilderStep>,
1705}
1706
1707impl LoadBalancerBuilder {
1708    /// Set the load balance strategy to round-robin (default).
1709    pub fn round_robin(mut self) -> Self {
1710        self.config = LoadBalancerConfig::round_robin();
1711        self
1712    }
1713
1714    /// Set the load balance strategy to random selection.
1715    pub fn random(mut self) -> Self {
1716        self.config = LoadBalancerConfig::random();
1717        self
1718    }
1719
1720    /// Set the load balance strategy to weighted selection.
1721    ///
1722    /// Each endpoint is assigned a weight that determines its probability
1723    /// of being selected.
1724    pub fn weighted(mut self, weights: Vec<(String, u32)>) -> Self {
1725        self.config = LoadBalancerConfig::weighted(weights);
1726        self
1727    }
1728
1729    /// Set the load balance strategy to failover.
1730    ///
1731    /// Exchanges are sent to the first endpoint; on failure, the next endpoint
1732    /// is tried.
1733    pub fn failover(mut self) -> Self {
1734        self.config = LoadBalancerConfig::failover();
1735        self
1736    }
1737
1738    /// Close the load balance scope. Packages the accumulated sub-steps into a
1739    /// `BuilderStep::LoadBalance` and returns the parent `RouteBuilder`.
1740    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// ---------------------------------------------------------------------------
1757// Tests
1758// ---------------------------------------------------------------------------
1759
1760#[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        // We can verify steps were added by checking the structure
1847        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); // set_header + Filter + To("mock:result")
1909        assert!(matches!(&definition.steps()[0], BuilderStep::Processor(_))); // set_header
1910        assert!(matches!(&definition.steps()[1], BuilderStep::Filter { .. })); // filter
1911        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    // -----------------------------------------------------------------------
1996    // Processor behavior tests — exercise the real Tower services directly
1997    // -----------------------------------------------------------------------
1998
1999    #[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    // -----------------------------------------------------------------------
2064    // Sequential pipeline test
2065    // -----------------------------------------------------------------------
2066
2067    #[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    // -----------------------------------------------------------------------
2122    // Circuit breaker builder tests
2123    // -----------------------------------------------------------------------
2124
2125    #[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        // Route definition was built successfully with both configs.
2163    }
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    // --- Splitter builder tests ---
2312
2313    #[test]
2314    fn test_split_builder_typestate() {
2315        use camel_api::splitter::{SplitterConfig, split_body_lines};
2316
2317        // .split() returns SplitBuilder, .end_split() returns RouteBuilder
2318        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        // Should have 2 top-level steps: Split + To("mock:final")
2328        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        // Should have 1 top-level step: Split (containing 2 sub-steps)
2345        assert_eq!(definition.steps().len(), 1);
2346        match &definition.steps()[0] {
2347            BuilderStep::Split { steps, .. } => {
2348                assert_eq!(steps.len(), 2); // SetHeader + To
2349            }
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    // ── set_body / set_body_fn / set_header_fn builder tests ────────────────────
2435
2436    #[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    // ── FilterBuilder typestate tests ─────────────────────────────────────
2597
2598    #[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    // ── MulticastBuilder typestate tests ─────────────────────────────────────
2641
2642    #[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); // Multicast + To("mock:result")
2655    }
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    // ── Concurrency builder tests ─────────────────────────────────────
2677
2678    #[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    // ── Route lifecycle builder tests ─────────────────────────────────────
2741
2742    #[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    // ── Choice typestate tests ──────────────────────────────────────────────────
2816
2817    #[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        // Steps after end_choice() are added to the outer pipeline, not inside choice.
2878        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") // must be step[1], not inside choice
2886            .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    // ── Throttle typestate tests ──────────────────────────────────────────────────
2893
2894    #[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); // SetHeader + To
2943            }
2944            other => panic!("Expected Throttle, got {:?}", other),
2945        }
2946    }
2947
2948    #[test]
2949    fn test_throttle_step_after_throttle() {
2950        // Steps after end_throttle() are added to the outer pipeline, not inside throttle.
2951        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    // ── LoadBalance typestate tests ──────────────────────────────────────────────────
2965
2966    #[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); // SetHeader + To
3017            }
3018            other => panic!("Expected LoadBalance, got {:?}", other),
3019        }
3020    }
3021
3022    #[test]
3023    fn test_load_balance_step_after_load_balance() {
3024        // Steps after end_load_balance() are added to the outer pipeline, not inside load_balance.
3025        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    // ── DynamicRouter typestate tests ──────────────────────────────────────────────────
3039
3040    #[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        // Steps after dynamic_router() are added to the outer pipeline.
3079        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    // ── LoadBalance strategy-specific tests ─────────────────────────────────────
3387
3388    #[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    // ── FilterInSplitBuilder tests ──────────────────────────────────────────────
3429
3430    #[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            // To("before-filter") + Filter{...} + To("after-filter") = 3
3471            assert_eq!(steps.len(), 3);
3472        } else {
3473            panic!("Expected Split step");
3474        }
3475    }
3476
3477    // ── build_canonical tests ───────────────────────────────────────────────────
3478
3479    #[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        // Split with closure-based expression is rejected in canonical v2.
3509        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        // Builder-based split uses closure expressions, which are not serializable.
3650        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    // ── OnExceptionBuilder full chain tests ─────────────────────────────────────
3664
3665    #[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    // ── StepAccumulator: process_fn, convert_body_to, bean ──────────────────────
3712
3713    #[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    // ── Throttle strategy-specific tests ────────────────────────────────────────
3752
3753    #[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    // ── LoopInLoopBuilder with loop_while ───────────────────────────────────────
3790
3791    #[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    // ── Choice with multiple whens + otherwise ──────────────────────────────────
3820
3821    #[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    // ── Multicast individual config tests ───────────────────────────────────────
3852
3853    #[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    // ── extract_completion_fields: Any mode with multiple conditions ────────────
3909
3910    #[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    // ── canonicalize_aggregate: discard_on_timeout and force_completion_on_stop ─
3953
3954    #[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    // ── build_canonical: max_buckets and bucket_ttl ─────────────────────────────
4001
4002    #[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    // ── canonicalize_aggregate: correlation strategy → correlation_key mapping ──
4028    //
4029    // These tests pin the canonicalize half of the expression round-trip
4030    // scenario. The recompile half is
4031    // `canonical_recompile_of_canonicalized_expression_config_yields_expression`
4032    // in camel-dsl compile.rs (task 2.4) — camel-builder has no camel-dsl
4033    // dependency, so the two halves live in their own crates.
4034
4035    #[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    // ── SplitBuilder with filter inside ─────────────────────────────────────────
4067
4068    #[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    // ── WireTap additional tests ────────────────────────────────────────────────
4091
4092    #[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    // ── Error handler: explicit config after shorthand → Mixed mode ─────────────
4112
4113    #[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    // ── build_canonical: empty from_uri error ───────────────────────────────────
4127
4128    #[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    // ── SplitBuilder: aggregate inside split ────────────────────────────────────
4143
4144    #[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    // ── Throttle: steps collected inside throttle scope ─────────────────────────
4171
4172    #[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    // ── LoadBalance: steps collected inside scope ───────────────────────────────
4191
4192    #[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    // ── Multicast: steps collected inside scope ─────────────────────────────────
4212
4213    #[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    // ── LoopBuilder: steps collected inside loop scope ──────────────────────────
4232
4233    #[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    // ── canonical_step_name coverage for remaining variants ─────────────────────
4252
4253    #[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    // ── extract_completion_fields: Any mode with predicate → error ──────────────
4392
4393    #[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    // ── BUILDER-004: Validation errors for missing required fields ────────────
4423
4424    #[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") // no scheme
4443            .build();
4444        // The builder itself accepts any URI string; validation happens at
4445        // resolution time. Verify the build succeeds (step URI is deferred).
4446        assert!(
4447            result.is_ok(),
4448            "builder should accept opaque step URIs; resolution happens later"
4449        );
4450    }
4451
4452    // ── rc-p9vq: Duplicate route IDs ──────────────────────────────────────
4453
4454    #[test]
4455    fn test_builder_duplicate_route_ids_produce_identical_definitions() {
4456        // The builder itself doesn't check for duplicates (that's context-level).
4457        // Verify both builds succeed with the same ID — detection is TODO(rc-p9vq).
4458        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        // rc-8m5o: a partially-built RouteBuilder can be cloned and reused as a
4475        // template, then each clone varied independently before build().
4476        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        // Shared template steps (set_header + log) are present in both, plus the
4495        // per-clone `to` step: the clone is a deep copy, not an alias.
4496        assert_eq!(route_a.steps().len(), route_b.steps().len());
4497        assert_eq!(route_a.from_uri(), route_b.from_uri());
4498    }
4499}