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