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
12pub type SplitExpression =
18 Arc<dyn Fn(&Exchange) -> Result<Vec<Exchange>, CamelError> + Send + Sync>;
19
20#[derive(Clone)]
28#[non_exhaustive]
29pub enum SplitSource {
30 Sync(SplitExpression),
32 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 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
52fn 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
87pub type StreamingSplitExpression = Arc<
94 dyn Fn(Exchange) -> Pin<Box<dyn Stream<Item = Result<Exchange, CamelError>> + Send>>
95 + Send
96 + Sync,
97>;
98
99pub 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#[derive(Clone, Default)]
111#[non_exhaustive]
112pub enum AggregationStrategy {
113 #[default]
115 LastWins,
116 CollectAll,
118 Original,
120 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#[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 #[default]
153 Auto,
154 Ndjson,
156 Lines,
158 Chunks,
160 Zip,
162 Tar,
164 #[serde(rename = "tar.gz")]
166 #[ts(rename = "tar.gz")]
167 TarGz,
168}
169
170#[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 pub format: StreamSplitFormat,
189 pub max_record_bytes: usize,
191 pub batch_size: usize,
193 pub chunk_size: Option<usize>,
195 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 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 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
278pub const DEFAULT_TRACE_ITEM_THRESHOLD: usize = 100;
280
281#[derive(Clone)]
283pub struct SplitterConfig {
284 pub expression: SplitSource,
286 pub aggregation: AggregationStrategy,
288 pub parallel: bool,
290 pub parallel_limit: Option<usize>,
292 pub stop_on_exception: bool,
298 pub max_fragments: usize,
304 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 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 pub fn aggregation(mut self, strategy: AggregationStrategy) -> Self {
338 self.aggregation = strategy;
339 self
340 }
341
342 pub fn parallel(mut self, parallel: bool) -> Self {
344 self.parallel = parallel;
345 self
346 }
347
348 pub fn parallel_limit(mut self, limit: usize) -> Self {
350 self.parallel_limit = Some(limit);
351 self
352 }
353
354 pub fn stop_on_exception(mut self, stop: bool) -> Self {
359 self.stop_on_exception = stop;
360 self
361 }
362
363 pub fn max_fragments(mut self, max: usize) -> Self {
365 self.max_fragments = max;
366 self
367 }
368
369 pub fn trace_item_threshold(mut self, threshold: usize) -> Self {
371 self.trace_item_threshold = threshold;
372 self
373 }
374
375 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
394pub 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 ex.otel_context = parent.otel_context.clone();
426 ex
427}
428
429pub 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
453pub 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
491pub 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 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()); 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 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 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()); 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()); 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 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 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 let fragments = split_body_lines()(&parent).unwrap();
785 assert!(!fragments.is_empty(), "Should have at least one fragment");
786
787 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 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 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 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}