Skip to main content

camel_processor/
stream_codec.rs

1use bytes::Bytes;
2use camel_api::{
3    Body, CamelError, Exchange, StreamMetadata, StreamSplitConfig, StreamSplitFormat,
4    fragment_exchange,
5};
6use futures::Stream;
7use std::pin::Pin;
8
9pub mod chunks;
10pub mod lines;
11pub mod ndjson;
12
13pub const CAMEL_STREAM_ORIGIN: &str = "CamelStreamOrigin";
14pub const CAMEL_STREAM_SOURCE_CONTENT_TYPE: &str = "CamelStreamSourceContentType";
15pub const CAMEL_STREAM_OFFSET: &str = "CamelStreamOffset";
16pub const CAMEL_STREAM_BATCH_SIZE: &str = "CamelStreamBatchSize";
17
18pub struct StreamSplitInput {
19    pub parent: Exchange,
20    pub stream: Pin<Box<dyn Stream<Item = Result<Bytes, CamelError>> + Send>>,
21    pub metadata: StreamMetadata,
22}
23
24pub trait StreamSplitCodec: Send + Sync {
25    fn split(
26        &self,
27        input: StreamSplitInput,
28        config: StreamSplitConfig,
29    ) -> Pin<Box<dyn Stream<Item = Result<Exchange, CamelError>> + Send>>;
30}
31
32pub fn resolve_format(
33    format: &StreamSplitFormat,
34    metadata: &StreamMetadata,
35) -> Result<StreamSplitFormat, CamelError> {
36    match format {
37        StreamSplitFormat::Auto => {
38            let ct = metadata
39                .content_type
40                .as_deref()
41                .unwrap_or("")
42                .to_lowercase();
43            let ct = ct.split(';').next().unwrap_or("").trim();
44            match ct {
45                "application/x-ndjson" => Ok(StreamSplitFormat::Ndjson),
46                "text/plain" => Ok(StreamSplitFormat::Lines),
47                "application/octet-stream" => Ok(StreamSplitFormat::Chunks),
48                "application/zip" | "application/x-zip-compressed" => Err(CamelError::Config(
49                    "stream split format=Auto: ZIP archives require explicit stream.format: zip"
50                        .into(),
51                )),
52                "application/x-tar" => Err(CamelError::Config(
53                    "stream split format=Auto: TAR archives require explicit stream.format: tar"
54                        .into(),
55                )),
56                "application/gzip" | "application/x-gzip" => Err(CamelError::Config(
57                    "stream split format=Auto: GZIP archives require explicit stream.format: tar.gz"
58                        .into(),
59                )),
60                "" => Err(CamelError::Config(
61                    "stream split format=Auto but stream has no content_type".into(),
62                )),
63                other => Err(CamelError::Config(format!(
64                    "stream split format=Auto: unknown content type '{}'",
65                    other
66                ))),
67            }
68        }
69        other => Ok(other.clone()),
70    }
71}
72
73#[derive(Debug)]
74pub enum ArchiveSplitKind {
75    Zip,
76    Tar,
77    TarGz,
78}
79
80pub enum ResolvedStreamSplit {
81    Incremental(Box<dyn StreamSplitCodec>),
82    MaterializedArchive(ArchiveSplitKind),
83}
84
85pub fn resolve_incremental_codec(
86    format: &StreamSplitFormat,
87) -> Result<Box<dyn StreamSplitCodec>, CamelError> {
88    match format {
89        StreamSplitFormat::Ndjson => Ok(Box::new(ndjson::NdjsonCodec)),
90        StreamSplitFormat::Lines => Ok(Box::new(lines::LinesCodec)),
91        StreamSplitFormat::Chunks => Ok(Box::new(chunks::ChunksCodec)),
92        StreamSplitFormat::Zip => Err(CamelError::Config(
93            "Zip is a materialized archive format, not an incremental codec".into(),
94        )),
95        StreamSplitFormat::Tar | StreamSplitFormat::TarGz => Err(CamelError::Config(
96            "Tar and TarGz are materialized archive formats, not incremental codecs".into(),
97        )),
98        StreamSplitFormat::Auto => Err(CamelError::Config(
99            "resolve_incremental_codec requires a resolved format, not Auto".into(),
100        )),
101        _ => Err(CamelError::Config("unsupported stream split format".into())),
102    }
103}
104
105pub fn resolve_split(
106    format: &StreamSplitFormat,
107    metadata: &StreamMetadata,
108) -> Result<ResolvedStreamSplit, CamelError> {
109    let resolved = resolve_format(format, metadata)?;
110    match resolved {
111        StreamSplitFormat::Zip => Ok(ResolvedStreamSplit::MaterializedArchive(
112            ArchiveSplitKind::Zip,
113        )),
114        StreamSplitFormat::Tar => Ok(ResolvedStreamSplit::MaterializedArchive(
115            ArchiveSplitKind::Tar,
116        )),
117        StreamSplitFormat::TarGz => Ok(ResolvedStreamSplit::MaterializedArchive(
118            ArchiveSplitKind::TarGz,
119        )),
120        _ => Ok(ResolvedStreamSplit::Incremental(resolve_incremental_codec(
121            &resolved,
122        )?)),
123    }
124}
125
126pub fn fragment_stream_exchange(parent: &Exchange, body: Body) -> Exchange {
127    let mut ex = fragment_exchange(parent, body);
128    ex.input.headers.remove("Content-Length");
129    ex.input.headers.remove("Content-Type");
130    ex
131}
132
133#[cfg(test)]
134mod tests {
135    use super::*;
136    use camel_api::StreamMetadata;
137
138    fn default_metadata() -> StreamMetadata {
139        StreamMetadata::default()
140    }
141
142    #[test]
143    fn test_resolve_split_ndjson_is_incremental() {
144        let result = resolve_split(&StreamSplitFormat::Ndjson, &default_metadata()).unwrap();
145        assert!(matches!(result, ResolvedStreamSplit::Incremental(_)));
146    }
147
148    #[test]
149    fn test_resolve_split_zip_is_materialized() {
150        let result = resolve_split(&StreamSplitFormat::Zip, &default_metadata()).unwrap();
151        assert!(matches!(
152            result,
153            ResolvedStreamSplit::MaterializedArchive(ArchiveSplitKind::Zip)
154        ));
155    }
156
157    #[test]
158    fn test_resolve_split_tar_is_materialized() {
159        let result = resolve_split(&StreamSplitFormat::Tar, &default_metadata()).unwrap();
160        assert!(matches!(
161            result,
162            ResolvedStreamSplit::MaterializedArchive(ArchiveSplitKind::Tar)
163        ));
164    }
165
166    #[test]
167    fn test_resolve_split_tar_gz_is_materialized() {
168        let result = resolve_split(&StreamSplitFormat::TarGz, &default_metadata()).unwrap();
169        assert!(matches!(
170            result,
171            ResolvedStreamSplit::MaterializedArchive(ArchiveSplitKind::TarGz)
172        ));
173    }
174
175    #[test]
176    fn test_resolve_incremental_codec_rejects_tar_formats() {
177        for format in [StreamSplitFormat::Tar, StreamSplitFormat::TarGz] {
178            let result = resolve_incremental_codec(&format);
179            let msg = match result {
180                Err(e) => e.to_string(),
181                Ok(_) => panic!("expected Err for {format:?}"),
182            };
183            assert!(
184                msg.contains("materialized archive format"),
185                "expected materialized-archive rejection, got: {msg}"
186            );
187        }
188    }
189
190    #[test]
191    fn test_auto_tar_content_types_suggest_explicit_format() {
192        for (ct, expected) in [
193            ("application/x-tar", "format: tar"),
194            ("application/gzip", "format: tar.gz"),
195        ] {
196            let meta = StreamMetadata {
197                content_type: Some(ct.to_string()),
198                ..Default::default()
199            };
200            let result = resolve_split(&StreamSplitFormat::Auto, &meta);
201            let msg = match result {
202                Err(e) => e.to_string(),
203                Ok(_) => panic!("expected Err for Auto + {ct}"),
204            };
205            assert!(
206                msg.contains(expected),
207                "expected suggestion '{expected}' in error, got: {msg}"
208            );
209        }
210    }
211
212    #[test]
213    fn test_resolve_incremental_codec_returns_codec_for_ndjson() {
214        let result = resolve_incremental_codec(&StreamSplitFormat::Ndjson);
215        assert!(result.is_ok());
216    }
217
218    #[test]
219    fn test_resolve_incremental_codec_rejects_zip() {
220        let result = resolve_incremental_codec(&StreamSplitFormat::Zip);
221        assert!(result.is_err());
222    }
223
224    #[test]
225    fn test_auto_application_zip_suggests_explicit_format() {
226        let meta = StreamMetadata {
227            content_type: Some("application/zip".to_string()),
228            ..Default::default()
229        };
230        let result = resolve_split(&StreamSplitFormat::Auto, &meta);
231        let msg = match result {
232            Err(e) => e.to_string(),
233            Ok(_) => panic!("expected Err for Auto + application/zip"),
234        };
235        assert!(
236            msg.contains("format: zip"),
237            "expected suggestion in error, got: {msg}"
238        );
239    }
240}