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                "" => Err(CamelError::Config(
53                    "stream split format=Auto but stream has no content_type".into(),
54                )),
55                other => Err(CamelError::Config(format!(
56                    "stream split format=Auto: unknown content type '{}'",
57                    other
58                ))),
59            }
60        }
61        other => Ok(other.clone()),
62    }
63}
64
65#[derive(Debug)]
66pub enum ArchiveSplitKind {
67    Zip,
68}
69
70pub enum ResolvedStreamSplit {
71    Incremental(Box<dyn StreamSplitCodec>),
72    MaterializedArchive(ArchiveSplitKind),
73}
74
75pub fn resolve_incremental_codec(
76    format: &StreamSplitFormat,
77) -> Result<Box<dyn StreamSplitCodec>, CamelError> {
78    match format {
79        StreamSplitFormat::Ndjson => Ok(Box::new(ndjson::NdjsonCodec)),
80        StreamSplitFormat::Lines => Ok(Box::new(lines::LinesCodec)),
81        StreamSplitFormat::Chunks => Ok(Box::new(chunks::ChunksCodec)),
82        StreamSplitFormat::Zip => Err(CamelError::Config(
83            "Zip is a materialized archive format, not an incremental codec".into(),
84        )),
85        StreamSplitFormat::Auto => Err(CamelError::Config(
86            "resolve_incremental_codec requires a resolved format, not Auto".into(),
87        )),
88        _ => Err(CamelError::Config("unsupported stream split format".into())),
89    }
90}
91
92pub fn resolve_split(
93    format: &StreamSplitFormat,
94    metadata: &StreamMetadata,
95) -> Result<ResolvedStreamSplit, CamelError> {
96    let resolved = resolve_format(format, metadata)?;
97    match resolved {
98        StreamSplitFormat::Zip => Ok(ResolvedStreamSplit::MaterializedArchive(
99            ArchiveSplitKind::Zip,
100        )),
101        _ => Ok(ResolvedStreamSplit::Incremental(resolve_incremental_codec(
102            &resolved,
103        )?)),
104    }
105}
106
107pub fn fragment_stream_exchange(parent: &Exchange, body: Body) -> Exchange {
108    let mut ex = fragment_exchange(parent, body);
109    ex.input.headers.remove("Content-Length");
110    ex.input.headers.remove("Content-Type");
111    ex
112}
113
114#[cfg(test)]
115mod tests {
116    use super::*;
117    use camel_api::StreamMetadata;
118
119    fn default_metadata() -> StreamMetadata {
120        StreamMetadata::default()
121    }
122
123    #[test]
124    fn test_resolve_split_ndjson_is_incremental() {
125        let result = resolve_split(&StreamSplitFormat::Ndjson, &default_metadata()).unwrap();
126        assert!(matches!(result, ResolvedStreamSplit::Incremental(_)));
127    }
128
129    #[test]
130    fn test_resolve_split_zip_is_materialized() {
131        let result = resolve_split(&StreamSplitFormat::Zip, &default_metadata()).unwrap();
132        assert!(matches!(
133            result,
134            ResolvedStreamSplit::MaterializedArchive(ArchiveSplitKind::Zip)
135        ));
136    }
137
138    #[test]
139    fn test_resolve_incremental_codec_returns_codec_for_ndjson() {
140        let result = resolve_incremental_codec(&StreamSplitFormat::Ndjson);
141        assert!(result.is_ok());
142    }
143
144    #[test]
145    fn test_resolve_incremental_codec_rejects_zip() {
146        let result = resolve_incremental_codec(&StreamSplitFormat::Zip);
147        assert!(result.is_err());
148    }
149
150    #[test]
151    fn test_auto_application_zip_suggests_explicit_format() {
152        let meta = StreamMetadata {
153            content_type: Some("application/zip".to_string()),
154            ..Default::default()
155        };
156        let result = resolve_split(&StreamSplitFormat::Auto, &meta);
157        let msg = match result {
158            Err(e) => e.to_string(),
159            Ok(_) => panic!("expected Err for Auto + application/zip"),
160        };
161        assert!(
162            msg.contains("format: zip"),
163            "expected suggestion in error, got: {msg}"
164        );
165    }
166}