camel_processor/data_format/
zip.rs1use 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;
11const DEFAULT_MAX_INPUT_SIZE: u64 = 64 * 1024 * 1024; const 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 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}