Skip to main content

camel_api/
splitter.rs

1use std::pin::Pin;
2use std::sync::Arc;
3
4use futures::Stream;
5
6use crate::body::{Body, body_type_name};
7use crate::error::CamelError;
8use crate::exchange::Exchange;
9use crate::message::Message;
10use crate::value::Value;
11
12/// A function that splits a single exchange into multiple fragment exchanges.
13///
14/// Built-in expressions return [`CamelError::TypeConversionFailed`] when the
15/// body type does not match their input contract; empty-content bodies yield
16/// `Ok(Vec::new())` (pass-through).
17pub type SplitExpression =
18    Arc<dyn Fn(&Exchange) -> Result<Vec<Exchange>, CamelError> + Send + Sync>;
19
20/// Fallible split source: a synchronous [`SplitExpression`] or a
21/// language-backed asynchronous expression.
22///
23/// The async arm evaluates to a [`Value`]; fragments are derived with the
24/// canonical split rules (string → non-empty line fragments, array → one
25/// fragment per element). Evaluation errors propagate before any fragment
26/// is materialized.
27#[derive(Clone)]
28#[non_exhaustive]
29pub enum SplitSource {
30    /// Programmatic synchronous splitter.
31    Sync(SplitExpression),
32    /// Language-backed asynchronous expression.
33    Async(Arc<dyn Fn(&Exchange) -> crate::filter::BoxValueFuture + Send + Sync>),
34}
35
36impl From<SplitExpression> for SplitSource {
37    fn from(f: SplitExpression) -> Self {
38        Self::Sync(f)
39    }
40}
41
42impl SplitSource {
43    /// Split `exchange` into fragments, propagating evaluation failures.
44    pub async fn split(&self, exchange: &Exchange) -> Result<Vec<Exchange>, CamelError> {
45        match self {
46            Self::Sync(f) => f(exchange),
47            Self::Async(f) => derive_fragments(exchange, f(exchange).await?),
48        }
49    }
50}
51
52/// Derive fragment exchanges from an evaluated language value.
53///
54/// Mirrors the rules previously inlined in camel-core `splitting.rs`
55/// (DeclarativeSplit): a string yields one fragment per non-empty line; an
56/// array yields one fragment per element (string elements become text
57/// bodies, everything else JSON); anything else is a type error.
58fn derive_fragments(exchange: &Exchange, value: Value) -> Result<Vec<Exchange>, CamelError> {
59    match value {
60        Value::String(s) => Ok(s
61            .lines()
62            .filter(|line| !line.is_empty())
63            .map(|line| {
64                let mut fragment = exchange.clone();
65                fragment.input.body = Body::from(line.to_string());
66                fragment
67            })
68            .collect()),
69        Value::Array(arr) => Ok(arr
70            .into_iter()
71            .map(|v| {
72                let mut fragment = exchange.clone();
73                fragment.input.body = match v {
74                    Value::String(s) => Body::Text(s),
75                    other => Body::Json(other),
76                };
77                fragment
78            })
79            .collect()),
80        other => Err(CamelError::TypeConversionFailed(format!(
81            "declarative split requires a text or array value, got {received}; add an unmarshal step before split",
82            received = crate::value_type_name(&other)
83        ))),
84    }
85}
86
87/// A function that lazily produces a stream of exchange fragments.
88///
89/// Used by `StreamingSplitterService` (camel-processor) for v1 sequential streaming split
90/// (e.g., ZIP entry extraction, CSV/JSON streaming in future work).
91///
92/// Each call returns a `Stream` that yields fragments one at a time.
93pub type StreamingSplitExpression = Arc<
94    dyn Fn(Exchange) -> Pin<Box<dyn Stream<Item = Result<Exchange, CamelError>> + Send>>
95        + Send
96        + Sync,
97>;
98
99/// Typed error returned when a streaming split receives a body that is not
100/// `Body::Stream`. Shared by the compiled production expression, the test
101/// mirrors, and examples so the message cannot drift between copies.
102pub fn streaming_split_type_error(body: &Body) -> CamelError {
103    CamelError::TypeConversionFailed(format!(
104        "streaming split requires body type stream, got {}; add an unmarshal step before split",
105        body_type_name(body)
106    ))
107}
108
109/// Strategy for aggregating fragment results back into a single exchange.
110#[derive(Clone, Default)]
111#[non_exhaustive]
112pub enum AggregationStrategy {
113    /// Result is the last fragment's exchange (default).
114    #[default]
115    LastWins,
116    /// Collects all fragment bodies into a JSON array.
117    CollectAll,
118    /// Returns the original exchange unchanged.
119    Original,
120    /// Custom aggregation function: `(accumulated, next) -> merged`.
121    Custom(Arc<dyn Fn(Exchange, Exchange) -> Exchange + Send + Sync>),
122}
123
124impl std::fmt::Debug for AggregationStrategy {
125    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
126        match self {
127            AggregationStrategy::LastWins => f.write_str("LastWins"),
128            AggregationStrategy::CollectAll => f.write_str("CollectAll"),
129            AggregationStrategy::Original => f.write_str("Original"),
130            AggregationStrategy::Custom(_) => f.write_str("Custom(..)"),
131        }
132    }
133}
134
135/// The streaming format to use when splitting a stream body.
136#[derive(
137    Clone,
138    Debug,
139    Default,
140    PartialEq,
141    Eq,
142    serde::Serialize,
143    serde::Deserialize,
144    schemars::JsonSchema,
145    ts_rs::TS,
146)]
147#[serde(rename_all = "snake_case")]
148#[ts(rename_all = "snake_case")]
149#[non_exhaustive]
150pub enum StreamSplitFormat {
151    /// Auto-detect the format from the body content.
152    #[default]
153    Auto,
154    /// Newline-delimited JSON — each line is a complete JSON value.
155    Ndjson,
156    /// Split by newlines, each line becomes a text fragment.
157    Lines,
158    /// Split into fixed-size byte chunks.
159    Chunks,
160    /// ZIP archive — materialized format, each entry becomes a fragment exchange.
161    Zip,
162    /// TAR archive — materialized format, each regular-file entry becomes a fragment exchange.
163    Tar,
164    /// GZIP-compressed TAR archive — materialized format, each regular-file entry becomes a fragment exchange.
165    #[serde(rename = "tar.gz")]
166    #[ts(rename = "tar.gz")]
167    TarGz,
168}
169
170/// Configuration for splitting a streaming body into fragments.
171///
172/// Controls how the stream splitter processes the body, including format
173/// detection, sizing limits, and metadata propagation.
174#[derive(
175    Clone,
176    Debug,
177    PartialEq,
178    Eq,
179    serde::Serialize,
180    serde::Deserialize,
181    schemars::JsonSchema,
182    ts_rs::TS,
183)]
184#[serde(rename_all = "snake_case")]
185#[ts(rename_all = "snake_case")]
186pub struct StreamSplitConfig {
187    /// The streaming format to use.
188    pub format: StreamSplitFormat,
189    /// Maximum size (in bytes) of a single record or chunk.
190    pub max_record_bytes: usize,
191    /// Number of records/chunks to collect into a single exchange batch.
192    pub batch_size: usize,
193    /// Explicit chunk size in bytes (required when format is [`Chunks`](StreamSplitFormat::Chunks)).
194    pub chunk_size: Option<usize>,
195    /// Whether to include origin metadata in each fragment.
196    pub include_origin: bool,
197}
198
199impl Default for StreamSplitConfig {
200    fn default() -> Self {
201        Self {
202            format: StreamSplitFormat::Auto,
203            max_record_bytes: 1024 * 1024,
204            batch_size: 1,
205            chunk_size: None,
206            include_origin: true,
207        }
208    }
209}
210
211impl StreamSplitConfig {
212    /// Validates the configuration.
213    ///
214    /// # Errors
215    ///
216    /// Returns [`CamelError::Config`] if:
217    /// - `batch_size` is `0`
218    /// - `max_record_bytes` is `0`
219    /// - `format` is [`Chunks`](StreamSplitFormat::Chunks) but `chunk_size` is `None`
220    /// - `format` is a materialized archive format ([`Zip`](StreamSplitFormat::Zip),
221    ///   [`Tar`](StreamSplitFormat::Tar), [`TarGz`](StreamSplitFormat::TarGz)) but
222    ///   `chunk_size` is `Some(...)`
223    /// - `chunk_size` is `Some(0)`
224    pub fn validate(&self) -> Result<(), CamelError> {
225        if self.batch_size == 0 {
226            return Err(CamelError::Config(
227                "stream split batch_size must be > 0".into(),
228            ));
229        }
230        if self.max_record_bytes == 0 {
231            return Err(CamelError::Config(
232                "stream split max_record_bytes must be > 0".into(),
233            ));
234        }
235        if self.format == StreamSplitFormat::Chunks && self.chunk_size.is_none() {
236            return Err(CamelError::Config(
237                "stream split format=Chunks requires chunk_size".into(),
238            ));
239        }
240        // Materialized archive formats split whole entries, not byte chunks,
241        // so chunk_size is meaningless. The Zip+chunk_size check must come
242        // before the generic chunk_size zero/exceeds checks so that
243        // `Zip + Some(0)` yields the more specific error.
244        if self.format == StreamSplitFormat::Zip && self.chunk_size.is_some() {
245            return Err(CamelError::Config(
246                "stream split format=Zip does not support chunk_size".into(),
247            ));
248        }
249        if matches!(
250            self.format,
251            StreamSplitFormat::Tar | StreamSplitFormat::TarGz
252        ) && self.chunk_size.is_some()
253        {
254            return Err(CamelError::Config(format!(
255                "stream split format={:?} is a materialized archive format and does not support chunk_size",
256                self.format
257            )));
258        }
259        if let Some(cs) = self.chunk_size
260            && cs == 0
261        {
262            return Err(CamelError::Config(
263                "stream split chunk_size must be > 0".into(),
264            ));
265        }
266        if self.format == StreamSplitFormat::Chunks
267            && let Some(cs) = self.chunk_size
268            && cs > self.max_record_bytes
269        {
270            return Err(CamelError::Config(
271                "stream split chunk_size must be <= max_record_bytes".into(),
272            ));
273        }
274        Ok(())
275    }
276}
277
278/// Default threshold above which split fragments start new traces (0 = off).
279pub const DEFAULT_TRACE_ITEM_THRESHOLD: usize = 100;
280
281/// Configuration for the Splitter EIP.
282#[derive(Clone)]
283pub struct SplitterConfig {
284    /// Expression that splits an exchange into fragments.
285    pub expression: SplitSource,
286    /// How to aggregate fragment results.
287    pub aggregation: AggregationStrategy,
288    /// Whether to process fragments in parallel.
289    pub parallel: bool,
290    /// Maximum number of parallel fragments (None = unlimited).
291    pub parallel_limit: Option<usize>,
292    /// Whether to stop processing on the first exception.
293    ///
294    /// In parallel mode this only affects aggregation (the first error is
295    /// propagated), **not** in-flight futures — `join_all` cannot cancel
296    /// already-spawned work.
297    pub stop_on_exception: bool,
298    /// Maximum number of fragments materialized by `expression` (DoS cap, R3-M4).
299    ///
300    /// The eager splitter materializes the whole `Vec<Exchange>` before
301    /// processing; this cap rejects a split that would explode memory.
302    /// Default 100_000. For unbounded/lazy input use `StreamingSplitter`.
303    pub max_fragments: usize,
304    /// Threshold above which split fragments start new traces (0 = off).
305    pub trace_item_threshold: usize,
306}
307
308impl std::fmt::Debug for SplitterConfig {
309    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
310        f.debug_struct("SplitterConfig")
311            .field("expression", &"<split-expression>")
312            .field("aggregation", &self.aggregation)
313            .field("parallel", &self.parallel)
314            .field("parallel_limit", &self.parallel_limit)
315            .field("stop_on_exception", &self.stop_on_exception)
316            .field("max_fragments", &self.max_fragments)
317            .field("trace_item_threshold", &self.trace_item_threshold)
318            .finish()
319    }
320}
321
322impl SplitterConfig {
323    /// Create a new splitter config with the given split expression.
324    pub fn new(expression: impl Into<SplitSource>) -> Self {
325        Self {
326            expression: expression.into(),
327            aggregation: AggregationStrategy::default(),
328            parallel: false,
329            parallel_limit: None,
330            stop_on_exception: true,
331            max_fragments: 100_000,
332            trace_item_threshold: DEFAULT_TRACE_ITEM_THRESHOLD,
333        }
334    }
335
336    /// Set the aggregation strategy for combining fragment results.
337    pub fn aggregation(mut self, strategy: AggregationStrategy) -> Self {
338        self.aggregation = strategy;
339        self
340    }
341
342    /// Enable or disable parallel fragment processing.
343    pub fn parallel(mut self, parallel: bool) -> Self {
344        self.parallel = parallel;
345        self
346    }
347
348    /// Set the maximum number of concurrent fragments in parallel mode.
349    pub fn parallel_limit(mut self, limit: usize) -> Self {
350        self.parallel_limit = Some(limit);
351        self
352    }
353
354    /// Control whether processing stops on the first fragment error.
355    ///
356    /// In parallel mode this only affects aggregation — see the field-level
357    /// doc comment for details.
358    pub fn stop_on_exception(mut self, stop: bool) -> Self {
359        self.stop_on_exception = stop;
360        self
361    }
362
363    /// Set the maximum number of fragments the eager splitter will materialize.
364    pub fn max_fragments(mut self, max: usize) -> Self {
365        self.max_fragments = max;
366        self
367    }
368
369    /// Set the threshold above which split fragments start new traces (0 = off).
370    pub fn trace_item_threshold(mut self, threshold: usize) -> Self {
371        self.trace_item_threshold = threshold;
372        self
373    }
374
375    /// Validates the configuration.
376    ///
377    /// Returns `Err(CamelError::Config)` if `parallel_limit` is set to 0,
378    /// which would cause a `Semaphore::new(0)` panic at runtime.
379    pub fn validate(&self) -> Result<(), CamelError> {
380        if self.parallel && self.parallel_limit == Some(0) {
381            return Err(CamelError::Config(
382                "splitter parallel_limit must be > 0".to_string(),
383            ));
384        }
385        if self.max_fragments == 0 {
386            return Err(CamelError::Config(
387                "splitter max_fragments must be > 0".to_string(),
388            ));
389        }
390        Ok(())
391    }
392}
393
394// ---------------------------------------------------------------------------
395// Helpers
396// ---------------------------------------------------------------------------
397
398/// Create a fragment exchange that inherits headers, properties, and OTel context
399/// from the parent, but with a new body.
400///
401/// # OpenTelemetry Trace Propagation
402///
403/// Each fragment inherits the live segment step span context: the splitter step's
404/// span, not the route root and not a previous fragment's span. A fragment-driven
405/// sub-route root therefore opens as a child of that segment span in the same trace,
406/// creating a natural fan-out relationship in the distributed trace:
407///
408/// ```text
409/// Segment step (span A)
410///   ├─ Fragment 1 sub-route root (span B, child of A)
411///   ├─ Fragment 2 sub-route root (span C, child of A)
412///   └─ Fragment N sub-route root (span N, child of A)
413/// ```
414///
415/// Restoring the entry context after the segment completes is the segment wrapper's
416/// job (`route_compiler::TracedSegmentStep` in `camel-core`), not
417/// `fragment_exchange`'s job.
418pub fn fragment_exchange(parent: &Exchange, body: Body) -> Exchange {
419    let mut msg = Message::new(body);
420    msg.headers = parent.input.headers.clone();
421    let mut ex = Exchange::new(msg);
422    ex.properties = parent.properties.clone();
423    ex.pattern = parent.pattern;
424    // Inherit OTel context so fragment spans are children of the parent span
425    ex.otel_context = parent.otel_context.clone();
426    ex
427}
428
429/// Split the exchange body by newlines. Returns one fragment per line.
430///
431/// Empty bodies pass through with zero fragments. Wrong-type bodies return a
432/// [`CamelError::TypeConversionFailed`] naming the received body type and the
433/// expected `text` input.
434pub fn split_body_lines() -> SplitExpression {
435    Arc::new(|exchange: &Exchange| {
436        let text = match &exchange.input.body {
437            Body::Text(s) => s.as_str(),
438            Body::Empty => return Ok(Vec::new()),
439            _ => {
440                return Err(CamelError::TypeConversionFailed(format!(
441                    "split expression 'body_lines' requires body type text, got {received}; add an unmarshal step before split",
442                    received = body_type_name(&exchange.input.body)
443                )));
444            }
445        };
446        Ok(text
447            .lines()
448            .map(|line| fragment_exchange(exchange, Body::Text(line.to_string())))
449            .collect())
450    })
451}
452
453/// Split a JSON array body into one fragment per element.
454///
455/// Fragment typing is element-driven: string elements produce [`Body::Text`]
456/// fragments carrying the raw string (so `${body}` renders it unquoted);
457/// number, boolean, object, nested-array, and null elements produce
458/// [`Body::Json`] fragments.
459///
460/// Empty bodies and empty arrays pass through with zero fragments. Non-array
461/// JSON and wrong-type bodies return a [`CamelError::TypeConversionFailed`]
462/// naming the received body type and the expected `json (array)` input.
463pub fn split_body_json_array() -> SplitExpression {
464    Arc::new(|exchange: &Exchange| {
465        let arr = match &exchange.input.body {
466            Body::Json(serde_json::Value::Array(arr)) => arr,
467            Body::Empty => return Ok(Vec::new()),
468            Body::Json(_) => {
469                return Err(CamelError::TypeConversionFailed(
470                    "split expression 'body_json_array' requires body type json (array), got json (non-array); add an unmarshal step before split"
471                        .to_string(),
472                ))
473            }
474            _ => {
475                return Err(CamelError::TypeConversionFailed(format!(
476                    "split expression 'body_json_array' requires body type json (array), got {received}; add an unmarshal step before split",
477                    received = body_type_name(&exchange.input.body)
478                )))
479            }
480        };
481        Ok(arr
482            .iter()
483            .map(|val| match val {
484                serde_json::Value::String(s) => fragment_exchange(exchange, Body::Text(s.clone())),
485                other => fragment_exchange(exchange, Body::Json(other.clone())),
486            })
487            .collect())
488    })
489}
490
491/// Split the exchange body using a custom function that operates on the body.
492///
493/// Custom closures stay infallible: they own their body-type policy and an
494/// empty `Vec` keeps pass-through semantics.
495pub fn split_body<F>(f: F) -> SplitExpression
496where
497    F: Fn(&Body) -> Vec<Body> + Send + Sync + 'static,
498{
499    Arc::new(move |exchange: &Exchange| {
500        Ok(f(&exchange.input.body)
501            .into_iter()
502            .map(|body| fragment_exchange(exchange, body))
503            .collect())
504    })
505}
506
507#[cfg(test)]
508mod tests {
509    use super::*;
510    use crate::value::Value;
511
512    #[test]
513    fn test_split_body_lines() {
514        let mut ex = Exchange::new(Message::new("a\nb\nc"));
515        ex.input.set_header("source", Value::String("test".into()));
516        ex.set_property("trace", Value::Bool(true));
517
518        let fragments = split_body_lines()(&ex).unwrap();
519        assert_eq!(fragments.len(), 3);
520        assert_eq!(fragments[0].input.body.as_text(), Some("a"));
521        assert_eq!(fragments[1].input.body.as_text(), Some("b"));
522        assert_eq!(fragments[2].input.body.as_text(), Some("c"));
523
524        // Verify headers and properties inherited
525        for frag in &fragments {
526            assert_eq!(
527                frag.input.header("source"),
528                Some(&Value::String("test".into()))
529            );
530            assert_eq!(frag.property("trace"), Some(&Value::Bool(true)));
531        }
532    }
533
534    #[test]
535    fn test_split_body_lines_empty() {
536        let ex = Exchange::new(Message::default()); // Body::Empty
537        let fragments = split_body_lines()(&ex).unwrap();
538        assert!(fragments.is_empty());
539    }
540
541    #[test]
542    fn test_split_body_json_array() {
543        let arr = serde_json::json!([1, 2, 3]);
544        let ex = Exchange::new(Message::new(arr));
545
546        let fragments = split_body_json_array()(&ex).unwrap();
547        assert_eq!(fragments.len(), 3);
548        assert!(matches!(&fragments[0].input.body, Body::Json(v) if *v == serde_json::json!(1)));
549        assert!(matches!(&fragments[1].input.body, Body::Json(v) if *v == serde_json::json!(2)));
550        assert!(matches!(&fragments[2].input.body, Body::Json(v) if *v == serde_json::json!(3)));
551    }
552
553    #[test]
554    fn split_body_json_array_string_elements_become_text() {
555        // The leading empty-string element pins the accepted delta from
556        // mission 251 (bd rc-etf0q): an empty element yields an empty
557        // `Body::Text`, not a quoted `""` JSON string.
558        let ex = Exchange::new(Message::new(serde_json::json!(["", "a", "b"])));
559
560        let fragments = split_body_json_array()(&ex).unwrap();
561        assert_eq!(fragments.len(), 3);
562        assert!(
563            matches!(&fragments[0].input.body, Body::Text(s) if s.is_empty()),
564            "fragment 0 must be Body::Text(\"\"), got {:?}",
565            fragments[0].input.body
566        );
567        assert!(matches!(&fragments[1].input.body, Body::Text(s) if s == "a"));
568        assert!(matches!(&fragments[2].input.body, Body::Text(s) if s == "b"));
569        // No fragment body carries a literal quote character: string elements
570        // are raw text, so `${body}` renders them unquoted.
571        for frag in &fragments {
572            let text = match &frag.input.body {
573                Body::Text(s) => s.as_str(),
574                other => panic!("expected Body::Text fragment, got {other:?}"),
575            };
576            assert!(
577                !text.contains('"'),
578                "fragment body must not carry a quote character, got {text:?}"
579            );
580        }
581    }
582
583    #[test]
584    fn test_split_body_json_array_non_string_elements_stay_json() {
585        let ex = Exchange::new(Message::new(serde_json::json!([1, {"k": "v"}, null])));
586
587        let fragments = split_body_json_array()(&ex).unwrap();
588        assert_eq!(fragments.len(), 3);
589        assert!(matches!(&fragments[0].input.body, Body::Json(v) if *v == serde_json::json!(1)));
590        assert!(matches!(&fragments[1].input.body, Body::Json(v)
591                if *v == serde_json::json!({"k": "v"})));
592        assert!(matches!(&fragments[2].input.body, Body::Json(v) if v.is_null()));
593    }
594
595    #[test]
596    fn test_split_body_json_array_not_array() {
597        let obj = serde_json::json!({"not": "array"});
598        let ex = Exchange::new(Message::new(obj));
599
600        let err = split_body_json_array()(&ex).unwrap_err();
601        assert!(matches!(err, CamelError::TypeConversionFailed(_)));
602        assert!(err.to_string().contains("json (non-array)"));
603    }
604
605    #[test]
606    fn test_split_body_lines_wrong_type_json_errors() {
607        let ex = Exchange::new(Message::new(serde_json::json!({"a": 1})));
608
609        let err = split_body_lines()(&ex).unwrap_err();
610        let msg = err.to_string();
611        assert!(matches!(err, CamelError::TypeConversionFailed(_)));
612        for needle in [
613            "body_lines",
614            "json",
615            "text",
616            "add an unmarshal step before split",
617        ] {
618            assert!(msg.contains(needle), "message '{msg}' missing '{needle}'");
619        }
620    }
621
622    #[test]
623    fn test_split_body_json_array_wrong_type_text_errors() {
624        let ex = Exchange::new(Message::new("x"));
625
626        let err = split_body_json_array()(&ex).unwrap_err();
627        let msg = err.to_string();
628        assert!(matches!(err, CamelError::TypeConversionFailed(_)));
629        for needle in [
630            "body_json_array",
631            "text",
632            "json (array)",
633            "add an unmarshal step before split",
634        ] {
635            assert!(msg.contains(needle), "message '{msg}' missing '{needle}'");
636        }
637    }
638
639    #[test]
640    fn test_split_body_json_array_non_array_json_errors() {
641        let ex = Exchange::new(Message::new(serde_json::json!({"o": 1})));
642
643        let err = split_body_json_array()(&ex).unwrap_err();
644        let msg = err.to_string();
645        assert!(matches!(err, CamelError::TypeConversionFailed(_)));
646        assert!(msg.contains("json (non-array)"));
647    }
648
649    #[test]
650    fn test_split_body_lines_empty_body_ok() {
651        let ex = Exchange::new(Message::default()); // Body::Empty
652        let fragments = split_body_lines()(&ex).unwrap();
653        assert!(fragments.is_empty());
654    }
655
656    #[test]
657    fn test_split_body_json_array_empty_body_ok() {
658        let ex = Exchange::new(Message::default()); // Body::Empty
659        let fragments = split_body_json_array()(&ex).unwrap();
660        assert!(fragments.is_empty());
661    }
662
663    #[test]
664    fn test_split_body_json_array_empty_array_ok() {
665        let ex = Exchange::new(Message::new(serde_json::json!([])));
666        let fragments = split_body_json_array()(&ex).unwrap();
667        assert!(fragments.is_empty());
668    }
669
670    #[test]
671    fn test_split_body_lines_empty_text_ok() {
672        let ex = Exchange::new(Message::new(""));
673        let fragments = split_body_lines()(&ex).unwrap();
674        assert!(fragments.is_empty());
675    }
676
677    #[test]
678    fn test_split_error_omits_payload() {
679        let ex = Exchange::new(Message::new(serde_json::json!({
680            "secret": "SECRET-8f31a"
681        })));
682
683        let err = split_body_lines()(&ex).unwrap_err();
684        let msg = err.to_string();
685        assert!(matches!(err, CamelError::TypeConversionFailed(_)));
686        for needle in [
687            "body_lines",
688            "json",
689            "text",
690            "add an unmarshal step before split",
691        ] {
692            assert!(msg.contains(needle), "message '{msg}' missing '{needle}'");
693        }
694        assert!(
695            !msg.contains("SECRET-8f31a"),
696            "message '{msg}' leaks payload"
697        );
698    }
699
700    #[test]
701    fn test_split_body_custom() {
702        let splitter = split_body(|body: &Body| match body {
703            Body::Text(s) => s
704                .split(',')
705                .map(|part| Body::Text(part.trim().to_string()))
706                .collect(),
707            _ => Vec::new(),
708        });
709
710        let mut ex = Exchange::new(Message::new("x, y, z"));
711        ex.set_property("id", Value::from(42));
712
713        let fragments = splitter(&ex).unwrap();
714        assert_eq!(fragments.len(), 3);
715        assert_eq!(fragments[0].input.body.as_text(), Some("x"));
716        assert_eq!(fragments[1].input.body.as_text(), Some("y"));
717        assert_eq!(fragments[2].input.body.as_text(), Some("z"));
718
719        // Properties inherited
720        for frag in &fragments {
721            assert_eq!(frag.property("id"), Some(&Value::from(42)));
722        }
723    }
724
725    #[test]
726    fn test_splitter_config_defaults() {
727        let config = SplitterConfig::new(split_body_lines());
728        assert!(matches!(config.aggregation, AggregationStrategy::LastWins));
729        assert!(!config.parallel);
730        assert!(config.parallel_limit.is_none());
731        assert!(config.stop_on_exception);
732    }
733
734    #[test]
735    fn test_splitter_config_builder() {
736        let config = SplitterConfig::new(split_body_lines())
737            .aggregation(AggregationStrategy::CollectAll)
738            .parallel(true)
739            .parallel_limit(4)
740            .stop_on_exception(false);
741
742        assert!(matches!(
743            config.aggregation,
744            AggregationStrategy::CollectAll
745        ));
746        assert!(config.parallel);
747        assert_eq!(config.parallel_limit, Some(4));
748        assert!(!config.stop_on_exception);
749    }
750
751    #[test]
752    fn test_splitter_config_default_max_fragments() {
753        let cfg = SplitterConfig::new(Arc::new(|_: &Exchange| Ok(Vec::new())) as SplitExpression);
754        assert_eq!(cfg.max_fragments, 100_000);
755    }
756
757    #[test]
758    fn test_splitter_config_rejects_zero_max_fragments() {
759        let cfg = SplitterConfig::new(Arc::new(|_: &Exchange| Ok(Vec::new())) as SplitExpression)
760            .max_fragments(0);
761        assert!(cfg.validate().is_err());
762    }
763
764    #[test]
765    fn test_fragment_exchange_inherits_otel_context() {
766        use opentelemetry::Context;
767        use opentelemetry::trace::{SpanContext, SpanId, TraceContextExt, TraceFlags, TraceId};
768
769        // Create parent exchange with a valid span context
770        let mut parent = Exchange::new(Message::new("test"));
771        let trace_id = TraceId::from_bytes([0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 123]);
772        let span_id = SpanId::from_bytes([0, 0, 0, 0, 0, 0, 1, 200]);
773        let span_context = SpanContext::new(
774            trace_id,
775            span_id,
776            TraceFlags::SAMPLED,
777            true,
778            Default::default(),
779        );
780        let expected_trace_id = span_context.trace_id();
781        parent.otel_context = Context::current().with_remote_span_context(span_context);
782
783        // Create fragment via split_body_lines
784        let fragments = split_body_lines()(&parent).unwrap();
785        assert!(!fragments.is_empty(), "Should have at least one fragment");
786
787        // Verify each fragment has the same span context as parent
788        for fragment in &fragments {
789            let span = fragment.otel_context.span();
790            let frag_span_ctx = span.span_context();
791            assert!(
792                frag_span_ctx.is_valid(),
793                "Fragment should have valid span context"
794            );
795            assert_eq!(
796                frag_span_ctx.trace_id(),
797                expected_trace_id,
798                "Fragment should have same trace ID as parent"
799            );
800        }
801    }
802
803    #[test]
804    fn test_stream_split_config_defaults_valid() {
805        let config = StreamSplitConfig::default();
806        assert!(config.validate().is_ok());
807    }
808
809    #[test]
810    fn test_stream_split_config_batch_size_zero_rejected() {
811        let config = StreamSplitConfig {
812            batch_size: 0,
813            ..Default::default()
814        };
815        let err = config.validate().unwrap_err();
816        assert!(err.to_string().contains("batch_size"));
817    }
818
819    #[test]
820    fn test_stream_split_config_max_record_bytes_zero_rejected() {
821        let config = StreamSplitConfig {
822            max_record_bytes: 0,
823            ..Default::default()
824        };
825        let err = config.validate().unwrap_err();
826        assert!(err.to_string().contains("max_record_bytes"));
827    }
828
829    #[test]
830    fn test_stream_split_config_chunks_requires_chunk_size() {
831        let config = StreamSplitConfig {
832            format: StreamSplitFormat::Chunks,
833            chunk_size: None,
834            ..Default::default()
835        };
836        let err = config.validate().unwrap_err();
837        assert!(err.to_string().contains("Chunks requires chunk_size"));
838    }
839
840    #[test]
841    fn test_stream_split_config_chunk_size_zero_rejected() {
842        let config = StreamSplitConfig {
843            format: StreamSplitFormat::Chunks,
844            chunk_size: Some(0),
845            ..Default::default()
846        };
847        let err = config.validate().unwrap_err();
848        assert!(err.to_string().contains("chunk_size must be > 0"));
849    }
850
851    #[test]
852    fn test_stream_split_config_chunk_size_exceeds_max_record_bytes() {
853        let config = StreamSplitConfig {
854            format: StreamSplitFormat::Chunks,
855            chunk_size: Some(2000),
856            max_record_bytes: 1000,
857            ..Default::default()
858        };
859        let err = config.validate().unwrap_err();
860        assert!(
861            err.to_string()
862                .contains("chunk_size must be <= max_record_bytes")
863        );
864    }
865
866    #[test]
867    fn test_stream_split_config_zip_rejects_chunk_size() {
868        let config = StreamSplitConfig {
869            format: StreamSplitFormat::Zip,
870            chunk_size: Some(1024),
871            ..Default::default()
872        };
873        let err = config.validate().unwrap_err();
874        assert!(err.to_string().contains("Zip does not support chunk_size"));
875    }
876
877    #[test]
878    fn test_all_fragments_share_same_trace_context() {
879        use opentelemetry::Context;
880        use opentelemetry::trace::{SpanContext, SpanId, TraceContextExt, TraceFlags, TraceId};
881
882        // Create parent with a specific trace ID
883        let mut parent = Exchange::new(Message::new("line1\nline2\nline3"));
884        let trace_id =
885            TraceId::from_bytes([0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0x3B, 0x9A, 0xCA, 0x09]);
886        let span_id = SpanId::from_bytes([0, 0, 0, 0, 0, 0, 0, 111]);
887        let span_context = SpanContext::new(
888            trace_id,
889            span_id,
890            TraceFlags::SAMPLED,
891            true,
892            Default::default(),
893        );
894        parent.otel_context = Context::current().with_remote_span_context(span_context);
895
896        let fragments = split_body_lines()(&parent).unwrap();
897        assert_eq!(fragments.len(), 3);
898
899        // All fragments should share the same trace ID
900        let trace_ids: Vec<_> = fragments
901            .iter()
902            .map(|f| {
903                let span = f.otel_context.span();
904                span.span_context().trace_id()
905            })
906            .collect();
907
908        assert!(
909            trace_ids.iter().all(|&id| id == trace_id),
910            "All fragments should have the same trace ID"
911        );
912    }
913
914    #[test]
915    fn trace_item_threshold_defaults_to_100() {
916        let config = SplitterConfig::new(split_body_lines());
917        assert_eq!(config.trace_item_threshold, 100);
918    }
919
920    #[test]
921    fn trace_item_threshold_builder_sets_value() {
922        let config = SplitterConfig::new(split_body_lines()).trace_item_threshold(0);
923        assert_eq!(config.trace_item_threshold, 0);
924
925        let config = SplitterConfig::new(split_body_lines()).trace_item_threshold(7);
926        assert_eq!(config.trace_item_threshold, 7);
927    }
928
929    #[test]
930    fn canonical_split_spec_carries_trace_item_threshold() {
931        use crate::runtime::{CanonicalSplitAggregationSpec, CanonicalSplitExpressionSpec};
932
933        let with_threshold = crate::runtime::CanonicalStepSpec::Split {
934            expression: CanonicalSplitExpressionSpec::BodyLines,
935            aggregation: CanonicalSplitAggregationSpec::CollectAll,
936            parallel: false,
937            parallel_limit: None,
938            stop_on_exception: true,
939            trace_item_threshold: Some(5),
940            steps: Vec::new(),
941        };
942        let without_threshold = crate::runtime::CanonicalStepSpec::Split {
943            expression: CanonicalSplitExpressionSpec::BodyLines,
944            aggregation: CanonicalSplitAggregationSpec::CollectAll,
945            parallel: false,
946            parallel_limit: None,
947            stop_on_exception: true,
948            trace_item_threshold: None,
949            steps: Vec::new(),
950        };
951
952        // Mirror the derive(PartialEq)/destructure path `parallel_limit` uses.
953        let crate::runtime::CanonicalStepSpec::Split {
954            trace_item_threshold,
955            ..
956        } = &with_threshold
957        else {
958            unreachable!()
959        };
960        assert_eq!(trace_item_threshold, &Some(5));
961
962        let crate::runtime::CanonicalStepSpec::Split {
963            trace_item_threshold,
964            ..
965        } = &without_threshold
966        else {
967            unreachable!()
968        };
969        assert_eq!(trace_item_threshold, &None);
970
971        assert_ne!(with_threshold, without_threshold);
972    }
973}