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}