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
11pub type SplitExpression =
17 Arc<dyn Fn(&Exchange) -> Result<Vec<Exchange>, CamelError> + Send + Sync>;
18
19pub type StreamingSplitExpression = Arc<
26 dyn Fn(Exchange) -> Pin<Box<dyn Stream<Item = Result<Exchange, CamelError>> + Send>>
27 + Send
28 + Sync,
29>;
30
31pub 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#[derive(Clone, Default)]
43#[non_exhaustive]
44pub enum AggregationStrategy {
45 #[default]
47 LastWins,
48 CollectAll,
50 Original,
52 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#[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 #[default]
85 Auto,
86 Ndjson,
88 Lines,
90 Chunks,
92 Zip,
94 Tar,
96 #[serde(rename = "tar.gz")]
98 #[ts(rename = "tar.gz")]
99 TarGz,
100}
101
102#[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 pub format: StreamSplitFormat,
121 pub max_record_bytes: usize,
123 pub batch_size: usize,
125 pub chunk_size: Option<usize>,
127 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 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 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
210pub const DEFAULT_TRACE_ITEM_THRESHOLD: usize = 100;
212
213#[derive(Clone)]
215pub struct SplitterConfig {
216 pub expression: SplitExpression,
218 pub aggregation: AggregationStrategy,
220 pub parallel: bool,
222 pub parallel_limit: Option<usize>,
224 pub stop_on_exception: bool,
230 pub max_fragments: usize,
236 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 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 pub fn aggregation(mut self, strategy: AggregationStrategy) -> Self {
270 self.aggregation = strategy;
271 self
272 }
273
274 pub fn parallel(mut self, parallel: bool) -> Self {
276 self.parallel = parallel;
277 self
278 }
279
280 pub fn parallel_limit(mut self, limit: usize) -> Self {
282 self.parallel_limit = Some(limit);
283 self
284 }
285
286 pub fn stop_on_exception(mut self, stop: bool) -> Self {
291 self.stop_on_exception = stop;
292 self
293 }
294
295 pub fn max_fragments(mut self, max: usize) -> Self {
297 self.max_fragments = max;
298 self
299 }
300
301 pub fn trace_item_threshold(mut self, threshold: usize) -> Self {
303 self.trace_item_threshold = threshold;
304 self
305 }
306
307 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
326pub 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 ex.otel_context = parent.otel_context.clone();
358 ex
359}
360
361pub 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
385pub 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
423pub 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 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()); 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 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 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()); 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()); 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 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 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 let fragments = split_body_lines()(&parent).unwrap();
717 assert!(!fragments.is_empty(), "Should have at least one fragment");
718
719 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 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 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 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}