Skip to main content

camel_processor/data_format/
zip.rs

1use bytes::Bytes;
2use camel_api::body::Body;
3use camel_api::data_format::DataFormat;
4use camel_api::error::CamelError;
5use serde::Deserialize;
6use std::io::Read;
7use std::io::Write;
8use zip::ZipArchive;
9
10const DEFAULT_MAX_DECOMPRESSED_SIZE: u64 = 1_073_741_824;
11/// Default cap on the materialized input size of `marshal` (R3-L1). The eager
12/// marshal collects the whole body into a `Vec<u8>` before compression; this
13/// bounds that allocation.
14const DEFAULT_MAX_INPUT_SIZE: u64 = 64 * 1024 * 1024; // 64 MiB
15const ENTRY_NAME: &str = "payload";
16
17#[derive(Debug, Clone, Deserialize)]
18#[serde(default, deny_unknown_fields)]
19pub struct ZipConfig {
20    pub max_decompressed_size: u64,
21    /// Maximum materialized input size accepted by `marshal` (DoS cap, R3-L1).
22    pub max_input_size: u64,
23    pub compression_level: Option<i32>,
24    pub allow_multi_entry: bool,
25}
26
27impl Default for ZipConfig {
28    fn default() -> Self {
29        Self {
30            max_decompressed_size: DEFAULT_MAX_DECOMPRESSED_SIZE,
31            max_input_size: DEFAULT_MAX_INPUT_SIZE,
32            compression_level: None,
33            allow_multi_entry: false,
34        }
35    }
36}
37
38#[derive(Debug, Clone, Default)]
39pub struct ZipDataFormat {
40    config: ZipConfig,
41}
42
43impl ZipDataFormat {
44    pub fn new(config: ZipConfig) -> Self {
45        Self { config }
46    }
47}
48
49impl DataFormat for ZipDataFormat {
50    fn name(&self) -> &str {
51        "zip"
52    }
53
54    fn marshal(&self, body: Body) -> Result<Body, CamelError> {
55        let content: Vec<u8> = match &body {
56            Body::Text(s) => s.as_bytes().to_vec(),
57            Body::Json(v) => serde_json::to_vec(v).map_err(|e| {
58                CamelError::TypeConversionFailed(format!(
59                    "ZipDataFormat::marshal cannot serialize JSON: {e}"
60                ))
61            })?,
62            Body::Bytes(b) => b.to_vec(),
63            Body::Xml(s) => s.as_bytes().to_vec(),
64            Body::Empty => {
65                return Err(CamelError::TypeConversionFailed(
66                    "ZipDataFormat::marshal requires non-empty body".to_string(),
67                ));
68            }
69            Body::Stream(_) => {
70                return Err(CamelError::TypeConversionFailed(
71                    "cannot marshal Body::Stream — add 'stream_cache' or 'convert_body_to' before this step".to_string(),
72                ));
73            }
74            _ => {
75                return Err(CamelError::TypeConversionFailed(
76                    "ZipDataFormat::marshal does not support this body type".to_string(),
77                ));
78            }
79        };
80
81        if content.len() as u64 > self.config.max_input_size {
82            return Err(CamelError::TypeConversionFailed(format!(
83                "ZipDataFormat::marshal input {} bytes exceeds max_input_size {}",
84                content.len(),
85                self.config.max_input_size
86            )));
87        }
88
89        let mut buf = Vec::new();
90        {
91            let mut writer = zip::ZipWriter::new(std::io::Cursor::new(&mut buf));
92            let mut options = zip::write::SimpleFileOptions::default()
93                .compression_method(zip::CompressionMethod::Deflated);
94            if let Some(level) = self.config.compression_level {
95                if !(0..=9).contains(&level) {
96                    return Err(CamelError::TypeConversionFailed(format!(
97                        "ZipDataFormat::marshal compression_level must be 0-9, got {level}"
98                    )));
99                }
100                options = options.compression_level(Some(level as i64));
101            }
102            writer.start_file(ENTRY_NAME, options).map_err(|e| {
103                CamelError::TypeConversionFailed(format!(
104                    "ZipDataFormat::marshal failed to start entry: {e}"
105                ))
106            })?;
107            writer.write_all(&content).map_err(|e| {
108                CamelError::TypeConversionFailed(format!(
109                    "ZipDataFormat::marshal failed to write entry: {e}"
110                ))
111            })?;
112            writer.finish().map_err(|e| {
113                CamelError::TypeConversionFailed(format!(
114                    "ZipDataFormat::marshal failed to finalize archive: {e}"
115                ))
116            })?;
117        }
118
119        Ok(Body::Bytes(Bytes::from(buf)))
120    }
121
122    fn unmarshal(&self, body: Body) -> Result<Body, CamelError> {
123        let raw: Vec<u8> = match &body {
124            Body::Bytes(b) => b.to_vec(),
125            Body::Text(s) => s.as_bytes().to_vec(),
126            Body::Empty => {
127                return Err(CamelError::TypeConversionFailed(
128                    "ZipDataFormat::unmarshal requires non-empty body".to_string(),
129                ));
130            }
131            Body::Stream(_) => {
132                return Err(CamelError::TypeConversionFailed(
133                    "cannot unmarshal Body::Stream — use UnmarshalService which auto-materializes"
134                        .to_string(),
135                ));
136            }
137            Body::Json(_) | Body::Xml(_) => {
138                return Err(CamelError::TypeConversionFailed(
139                    "ZipDataFormat::unmarshal only supports Body::Bytes and Body::Text (ZIP data)"
140                        .to_string(),
141                ));
142            }
143            _ => {
144                return Err(CamelError::TypeConversionFailed(
145                    "ZipDataFormat::unmarshal does not support this body type".to_string(),
146                ));
147            }
148        };
149
150        let reader = std::io::Cursor::new(&raw);
151        let mut archive = ZipArchive::new(reader).map_err(|e| {
152            CamelError::TypeConversionFailed(format!("ZipDataFormat::unmarshal invalid ZIP: {e}"))
153        })?;
154
155        if archive.is_empty() {
156            return Err(CamelError::TypeConversionFailed(
157                "ZipDataFormat::unmarshal ZIP archive has no entries".to_string(),
158            ));
159        }
160
161        if archive.len() > 1 && !self.config.allow_multi_entry {
162            return Err(CamelError::TypeConversionFailed(format!(
163                "ZipDataFormat::unmarshal ZIP has {} entries but allow_multi_entry is false",
164                archive.len()
165            )));
166        }
167
168        if archive.len() > 1 {
169            tracing::warn!(
170                entries = archive.len(),
171                "ZIP archive has multiple entries, extracting first only"
172            );
173        }
174
175        let mut entry = archive.by_index(0).map_err(|e| {
176            CamelError::TypeConversionFailed(format!(
177                "ZipDataFormat::unmarshal failed to read entry: {e}"
178            ))
179        })?;
180
181        let mut decompressed = Vec::new();
182        let limit = self.config.max_decompressed_size.saturating_add(1);
183        let mut limited = std::io::Read::take(&mut entry, limit);
184        limited.read_to_end(&mut decompressed).map_err(|e| {
185            CamelError::TypeConversionFailed(format!(
186                "ZipDataFormat::unmarshal failed to decompress: {e}"
187            ))
188        })?;
189
190        if decompressed.len() as u64 > self.config.max_decompressed_size {
191            return Err(CamelError::TypeConversionFailed(format!(
192                "ZipDataFormat::unmarshal decompressed size exceeds max {}",
193                self.config.max_decompressed_size
194            )));
195        }
196
197        Ok(Body::Bytes(Bytes::from(decompressed)))
198    }
199}
200
201#[cfg(test)]
202mod tests {
203    use super::*;
204    use bytes::Bytes;
205    use serde_json::json;
206    use std::io::Cursor;
207    use std::io::Read;
208    use zip::ZipArchive;
209
210    fn extract_single_entry(zip_bytes: &[u8]) -> Vec<u8> {
211        let reader = Cursor::new(zip_bytes);
212        let mut archive = ZipArchive::new(reader).unwrap();
213        let mut entry = archive.by_index(0).unwrap();
214        let name = entry.name().to_string();
215        assert_eq!(name, "payload");
216        let mut buf = Vec::new();
217        entry.read_to_end(&mut buf).unwrap();
218        buf
219    }
220
221    #[test]
222    fn test_name() {
223        let df = ZipDataFormat::default();
224        assert_eq!(df.name(), "zip");
225    }
226
227    #[test]
228    fn test_marshal_text_to_zip() {
229        let df = ZipDataFormat::default();
230        let body = Body::Text("hello world".to_string());
231        let result = df.marshal(body).unwrap();
232        match result {
233            Body::Bytes(b) => {
234                let decompressed = extract_single_entry(&b);
235                assert_eq!(decompressed, b"hello world");
236            }
237            _ => panic!("expected Body::Bytes"),
238        }
239    }
240
241    #[test]
242    fn test_marshal_json_to_zip() {
243        let df = ZipDataFormat::default();
244        let body = Body::Json(json!({"key": "value"}));
245        let result = df.marshal(body).unwrap();
246        match result {
247            Body::Bytes(b) => {
248                let decompressed = extract_single_entry(&b);
249                let original = serde_json::to_vec(&json!({"key": "value"})).unwrap();
250                assert_eq!(decompressed, original);
251            }
252            _ => panic!("expected Body::Bytes"),
253        }
254    }
255
256    #[test]
257    fn test_marshal_bytes_to_zip() {
258        let df = ZipDataFormat::default();
259        let body = Body::Bytes(Bytes::from_static(b"raw bytes"));
260        let result = df.marshal(body).unwrap();
261        match result {
262            Body::Bytes(b) => {
263                let decompressed = extract_single_entry(&b);
264                assert_eq!(decompressed, b"raw bytes");
265            }
266            _ => panic!("expected Body::Bytes"),
267        }
268    }
269
270    #[test]
271    fn test_marshal_xml_to_zip() {
272        let df = ZipDataFormat::default();
273        let body = Body::Xml("<root><item>1</item></root>".to_string());
274        let result = df.marshal(body).unwrap();
275        match result {
276            Body::Bytes(b) => {
277                let decompressed = extract_single_entry(&b);
278                assert_eq!(decompressed, b"<root><item>1</item></root>");
279            }
280            _ => panic!("expected Body::Bytes"),
281        }
282    }
283
284    #[test]
285    fn test_marshal_empty_error() {
286        let df = ZipDataFormat::default();
287        let result = df.marshal(Body::Empty);
288        assert!(result.is_err());
289    }
290
291    #[test]
292    fn test_marshal_stream_error() {
293        use camel_api::body::{StreamBody, StreamMetadata};
294        use futures::stream;
295        use std::sync::Arc;
296        use tokio::sync::Mutex;
297
298        let stream = stream::iter(vec![Ok(Bytes::from_static(b"data"))]);
299        let body = Body::Stream(StreamBody {
300            stream: Arc::new(Mutex::new(Some(Box::pin(stream)))),
301            metadata: StreamMetadata::default(),
302        });
303        let df = ZipDataFormat::default();
304        let result = df.marshal(body);
305        assert!(result.is_err());
306    }
307
308    fn make_zip(content: &[u8]) -> Vec<u8> {
309        let mut buf = Vec::new();
310        {
311            let mut writer = zip::ZipWriter::new(std::io::Cursor::new(&mut buf));
312            let options = zip::write::SimpleFileOptions::default()
313                .compression_method(zip::CompressionMethod::Deflated);
314            writer.start_file("payload", options).unwrap();
315            writer.write_all(content).unwrap();
316            writer.finish().unwrap();
317        }
318        buf
319    }
320
321    #[test]
322    fn test_unmarshal_zip_bytes() {
323        let df = ZipDataFormat::default();
324        let zip_data = make_zip(b"decompressed content");
325        let body = Body::Bytes(Bytes::from(zip_data));
326        let result = df.unmarshal(body).unwrap();
327        match result {
328            Body::Bytes(b) => assert_eq!(b.as_ref(), b"decompressed content"),
329            _ => panic!("expected Body::Bytes"),
330        }
331    }
332
333    #[test]
334    fn test_unmarshal_zip_text() {
335        let df = ZipDataFormat::default();
336        let content = b"text from text body";
337        let zip_data = make_zip(content);
338        let body = Body::Bytes(Bytes::from(zip_data));
339        let result = df.unmarshal(body).unwrap();
340        match result {
341            Body::Bytes(b) => assert_eq!(b.as_ref(), content),
342            _ => panic!("expected Body::Bytes"),
343        }
344    }
345
346    #[test]
347    fn test_unmarshal_invalid_zip_error() {
348        let df = ZipDataFormat::default();
349        let body = Body::Bytes(Bytes::from_static(b"not a zip file"));
350        let result = df.unmarshal(body);
351        assert!(result.is_err());
352    }
353
354    #[test]
355    fn test_unmarshal_empty_zip_error() {
356        let mut buf = Vec::new();
357        {
358            let writer = zip::ZipWriter::new(std::io::Cursor::new(&mut buf));
359            writer.finish().unwrap();
360        }
361        let df = ZipDataFormat::default();
362        let body = Body::Bytes(Bytes::from(buf));
363        let result = df.unmarshal(body);
364        assert!(result.is_err());
365    }
366
367    #[test]
368    fn test_unmarshal_json_error() {
369        let df = ZipDataFormat::default();
370        let body = Body::Json(json!({"not": "zip"}));
371        let result = df.unmarshal(body);
372        assert!(result.is_err());
373    }
374
375    #[test]
376    fn test_unmarshal_xml_error() {
377        let df = ZipDataFormat::default();
378        let body = Body::Xml("<root/>".to_string());
379        let result = df.unmarshal(body);
380        assert!(result.is_err());
381    }
382
383    #[test]
384    fn test_unmarshal_multi_entry_error() {
385        let mut buf = Vec::new();
386        {
387            let mut writer = zip::ZipWriter::new(std::io::Cursor::new(&mut buf));
388            let options = zip::write::SimpleFileOptions::default();
389            writer.start_file("file1.txt", options).unwrap();
390            writer.write_all(b"one").unwrap();
391            writer.start_file("file2.txt", options).unwrap();
392            writer.write_all(b"two").unwrap();
393            writer.finish().unwrap();
394        }
395        let df = ZipDataFormat::default();
396        let body = Body::Bytes(Bytes::from(buf));
397        let result = df.unmarshal(body);
398        assert!(result.is_err());
399    }
400
401    #[test]
402    fn test_unmarshal_multi_entry_allowed() {
403        let mut buf = Vec::new();
404        {
405            let mut writer = zip::ZipWriter::new(std::io::Cursor::new(&mut buf));
406            let options = zip::write::SimpleFileOptions::default();
407            writer.start_file("file1.txt", options).unwrap();
408            writer.write_all(b"first").unwrap();
409            writer.start_file("file2.txt", options).unwrap();
410            writer.write_all(b"second").unwrap();
411            writer.finish().unwrap();
412        }
413        let config = ZipConfig {
414            allow_multi_entry: true,
415            ..Default::default()
416        };
417        let df = ZipDataFormat::new(config);
418        let body = Body::Bytes(Bytes::from(buf));
419        let result = df.unmarshal(body).unwrap();
420        match result {
421            Body::Bytes(b) => assert_eq!(b.as_ref(), b"first"),
422            _ => panic!("expected Body::Bytes"),
423        }
424    }
425
426    #[test]
427    fn test_roundtrip_text() {
428        let df = ZipDataFormat::default();
429        let original = Body::Text("roundtrip text content".to_string());
430        let compressed = df.marshal(original).unwrap();
431        let decompressed = df.unmarshal(compressed).unwrap();
432        match decompressed {
433            Body::Bytes(b) => assert_eq!(b.as_ref(), b"roundtrip text content"),
434            _ => panic!("expected Body::Bytes"),
435        }
436    }
437
438    #[test]
439    fn test_roundtrip_json() {
440        let df = ZipDataFormat::default();
441        let original = Body::Json(json!({"round": "trip"}));
442        let compressed = df.marshal(original).unwrap();
443        let decompressed = df.unmarshal(compressed).unwrap();
444        match decompressed {
445            Body::Bytes(b) => {
446                let v: serde_json::Value = serde_json::from_slice(&b).unwrap();
447                assert_eq!(v, json!({"round": "trip"}));
448            }
449            _ => panic!("expected Body::Bytes"),
450        }
451    }
452
453    #[test]
454    fn test_roundtrip_bytes() {
455        let df = ZipDataFormat::default();
456        let original = Body::Bytes(Bytes::from_static(b"\x00\x01\x02\xff"));
457        let compressed = df.marshal(original).unwrap();
458        let decompressed = df.unmarshal(compressed).unwrap();
459        match decompressed {
460            Body::Bytes(b) => assert_eq!(b.as_ref(), b"\x00\x01\x02\xff"),
461            _ => panic!("expected Body::Bytes"),
462        }
463    }
464
465    #[test]
466    fn test_max_decompressed_size_exceeded() {
467        let config = ZipConfig {
468            max_decompressed_size: 10,
469            ..Default::default()
470        };
471        let df = ZipDataFormat::new(config);
472        let zip_data = make_zip(b"this content is way longer than 10 bytes");
473        let body = Body::Bytes(Bytes::from(zip_data));
474        let result = df.unmarshal(body);
475        assert!(result.is_err());
476    }
477
478    #[test]
479    fn test_unmarshal_empty_error() {
480        let df = ZipDataFormat::default();
481        let result = df.unmarshal(Body::Empty);
482        assert!(result.is_err());
483    }
484
485    #[test]
486    fn test_unmarshal_stream_error() {
487        use camel_api::body::{StreamBody, StreamMetadata};
488        use futures::stream;
489        use std::sync::Arc;
490        use tokio::sync::Mutex;
491
492        let stream = stream::iter(vec![Ok(Bytes::from_static(b"data"))]);
493        let body = Body::Stream(StreamBody {
494            stream: Arc::new(Mutex::new(Some(Box::pin(stream)))),
495            metadata: StreamMetadata::default(),
496        });
497        let df = ZipDataFormat::default();
498        let result = df.unmarshal(body);
499        assert!(result.is_err());
500    }
501
502    #[test]
503    fn test_marshal_invalid_compression_level() {
504        let config = ZipConfig {
505            compression_level: Some(42),
506            ..Default::default()
507        };
508        let df = ZipDataFormat::new(config);
509        let result = df.marshal(Body::Text("test".to_string()));
510        assert!(result.is_err());
511    }
512
513    #[test]
514    fn test_marshal_input_size_cap() {
515        let config = ZipConfig {
516            max_input_size: 16,
517            ..Default::default()
518        };
519        let df = ZipDataFormat::new(config);
520        let body = Body::Text("x".repeat(64));
521        let result = df.marshal(body);
522        assert!(result.is_err());
523        let msg = format!("{}", result.unwrap_err());
524        assert!(
525            msg.contains("max_input_size"),
526            "error should mention max_input_size: {msg}"
527        );
528    }
529
530    #[test]
531    fn test_marshal_default_cap_accepts_small_input() {
532        let df = ZipDataFormat::default();
533        let body = Body::Text("hello".to_string());
534        let result = df.marshal(body).unwrap();
535        assert!(matches!(result, Body::Bytes(_)));
536    }
537
538    #[test]
539    fn test_builtin_zip_registered() {
540        let df = super::super::builtin_data_format("zip");
541        assert!(df.is_some());
542        assert_eq!(df.unwrap().name(), "zip");
543    }
544
545    #[test]
546    fn test_zip_config_deserialize_from_json() {
547        let json = serde_json::json!({
548            "max_decompressed_size": 2147483648u64,
549            "max_input_size": 134217728u64,
550            "compression_level": 6,
551            "allow_multi_entry": true
552        });
553        let cfg: ZipConfig = serde_json::from_value(json).unwrap();
554        assert_eq!(cfg.max_decompressed_size, 2147483648);
555        assert_eq!(cfg.compression_level, Some(6));
556        assert!(cfg.allow_multi_entry);
557    }
558
559    #[test]
560    fn test_zip_config_deny_unknown_fields() {
561        let json = serde_json::json!({"unknown_key": 42});
562        let result: Result<ZipConfig, _> = serde_json::from_value(json);
563        assert!(result.is_err());
564    }
565}