1use std::collections::{HashMap, HashSet};
19
20use bytes::Bytes;
21use serde::{Deserialize, Serialize};
22
23use crate::Result;
24use crate::compression::CompressionCodec;
25use crate::error::invalid_data;
26use crate::io::FileRead;
27
28pub const CREATED_BY_PROPERTY: &str = "created-by";
31
32#[derive(Debug, PartialEq, Eq, Serialize, Deserialize, Clone)]
35#[serde(rename_all = "kebab-case")]
36pub struct BlobMetadata {
37 pub(crate) r#type: String,
38 pub(crate) fields: Vec<i32>,
39 pub(crate) snapshot_id: i64,
40 pub(crate) sequence_number: i64,
41 pub(crate) offset: u64,
42 pub(crate) length: u64,
43 #[serde(skip_serializing_if = "CompressionCodec::is_none")]
44 #[serde(default)]
45 pub(crate) compression_codec: CompressionCodec,
46 #[serde(skip_serializing_if = "HashMap::is_empty")]
47 #[serde(default)]
48 pub(crate) properties: HashMap<String, String>,
49}
50
51impl BlobMetadata {
52 #[inline]
53 pub fn blob_type(&self) -> &str {
55 &self.r#type
56 }
57
58 #[inline]
59 pub fn fields(&self) -> &[i32] {
61 &self.fields
62 }
63
64 #[inline]
65 pub fn snapshot_id(&self) -> i64 {
67 self.snapshot_id
68 }
69
70 #[inline]
71 pub fn sequence_number(&self) -> i64 {
73 self.sequence_number
74 }
75
76 #[inline]
77 pub fn offset(&self) -> u64 {
79 self.offset
80 }
81
82 #[inline]
83 pub fn length(&self) -> u64 {
85 self.length
86 }
87
88 #[inline]
89 pub fn compression_codec(&self) -> CompressionCodec {
91 self.compression_codec
92 }
93
94 #[inline]
95 pub fn properties(&self) -> &HashMap<String, String> {
97 &self.properties
98 }
99}
100
101#[derive(Clone, Copy, PartialEq, Eq, Hash, Debug)]
102pub(crate) enum Flag {
103 FooterPayloadCompressed = 0,
104}
105
106impl Flag {
107 pub(crate) fn byte_idx(self) -> u8 {
108 (self as u8) / 8
109 }
110
111 pub(crate) fn bit_idx(self) -> u8 {
112 (self as u8) % 8
113 }
114
115 fn matches(self, byte_idx: u8, bit_idx: u8) -> bool {
116 self.byte_idx() == byte_idx && self.bit_idx() == bit_idx
117 }
118
119 fn from(byte_idx: u8, bit_idx: u8) -> Result<Flag> {
120 if Flag::FooterPayloadCompressed.matches(byte_idx, bit_idx) {
121 Ok(Flag::FooterPayloadCompressed)
122 } else {
123 Err(invalid_data!(
124 "Unknown flag byte {byte_idx} and bit {bit_idx} combination"
125 ))
126 }
127 }
128}
129
130#[derive(Debug, PartialEq, Eq, Serialize, Deserialize, Clone)]
134pub struct FileMetadata {
135 pub(crate) blobs: Vec<BlobMetadata>,
136 #[serde(skip_serializing_if = "HashMap::is_empty")]
137 #[serde(default)]
138 pub(crate) properties: HashMap<String, String>,
139}
140
141impl FileMetadata {
142 pub(crate) const MAGIC_LENGTH: u8 = 4;
143 pub(crate) const MAGIC: [u8; FileMetadata::MAGIC_LENGTH as usize] = [0x50, 0x46, 0x41, 0x31];
144
145 const FOOTER_STRUCT_PAYLOAD_LENGTH_OFFSET: u8 = 0;
168 const FOOTER_STRUCT_PAYLOAD_LENGTH_LENGTH: u8 = 4;
169 const FOOTER_STRUCT_FLAGS_OFFSET: u8 = FileMetadata::FOOTER_STRUCT_PAYLOAD_LENGTH_OFFSET
170 + FileMetadata::FOOTER_STRUCT_PAYLOAD_LENGTH_LENGTH;
171 pub(crate) const FOOTER_STRUCT_FLAGS_LENGTH: u8 = 4;
172 const FOOTER_STRUCT_MAGIC_OFFSET: u8 =
173 FileMetadata::FOOTER_STRUCT_FLAGS_OFFSET + FileMetadata::FOOTER_STRUCT_FLAGS_LENGTH;
174 pub(crate) const FOOTER_STRUCT_LENGTH: u8 =
175 FileMetadata::FOOTER_STRUCT_MAGIC_OFFSET + FileMetadata::MAGIC_LENGTH;
176
177 const MIN_FILE_LENGTH: u64 =
180 (FileMetadata::MAGIC_LENGTH as u64) * 2 + (FileMetadata::FOOTER_STRUCT_LENGTH as u64);
181
182 pub fn new(blobs: Vec<BlobMetadata>, properties: HashMap<String, String>) -> Self {
184 Self { blobs, properties }
185 }
186
187 fn check_magic(bytes: &[u8]) -> Result<()> {
188 if bytes == FileMetadata::MAGIC {
189 Ok(())
190 } else {
191 Err(invalid_data!(
192 "Bad magic value: {:?} should be {:?}",
193 bytes,
194 FileMetadata::MAGIC
195 ))
196 }
197 }
198
199 async fn read_footer_payload_length(
200 file_read: &dyn FileRead,
201 input_file_length: u64,
202 ) -> Result<u32> {
203 let start = input_file_length - FileMetadata::FOOTER_STRUCT_LENGTH as u64;
204 let end = start + FileMetadata::FOOTER_STRUCT_PAYLOAD_LENGTH_LENGTH as u64;
205 let footer_payload_length_bytes = file_read.read(start..end).await?;
206 let mut buf = [0; 4];
207 buf.copy_from_slice(&footer_payload_length_bytes);
208 let footer_payload_length = u32::from_le_bytes(buf);
209 Ok(footer_payload_length)
210 }
211
212 async fn read_footer_bytes(
213 file_read: &dyn FileRead,
214 input_file_length: u64,
215 footer_payload_length: u32,
216 ) -> Result<Bytes> {
217 let footer_length = footer_payload_length as u64
218 + FileMetadata::FOOTER_STRUCT_LENGTH as u64
219 + FileMetadata::MAGIC_LENGTH as u64;
220 let start = input_file_length
221 .checked_sub(footer_length)
222 .ok_or_else(|| {
223 invalid_data!(
224 "Footer length {footer_length} exceeds file length {input_file_length}"
225 )
226 })?;
227 let end = input_file_length;
228 file_read.read(start..end).await
229 }
230
231 fn decode_flags(footer_bytes: &[u8]) -> Result<HashSet<Flag>> {
232 let mut flags = HashSet::new();
233
234 for byte_idx in 0..FileMetadata::FOOTER_STRUCT_FLAGS_LENGTH {
235 let byte_offset = footer_bytes.len()
236 - FileMetadata::MAGIC_LENGTH as usize
237 - FileMetadata::FOOTER_STRUCT_FLAGS_LENGTH as usize
238 + byte_idx as usize;
239
240 let flag_byte = *footer_bytes
241 .get(byte_offset)
242 .ok_or_else(|| invalid_data!("Index range is out of bounds."))?;
243
244 for bit_idx in 0..8 {
245 if ((flag_byte >> bit_idx) & 1) != 0 {
246 let flag = Flag::from(byte_idx, bit_idx)?;
247 flags.insert(flag);
248 }
249 }
250 }
251
252 Ok(flags)
253 }
254
255 fn extract_footer_payload_as_str(
256 footer_bytes: &[u8],
257 footer_payload_length: u32,
258 ) -> Result<String> {
259 let flags = FileMetadata::decode_flags(footer_bytes)?;
260 let footer_compression_codec = if flags.contains(&Flag::FooterPayloadCompressed) {
261 CompressionCodec::Lz4
262 } else {
263 CompressionCodec::None
264 };
265
266 let start_offset = FileMetadata::MAGIC_LENGTH as usize;
267 let end_offset =
268 FileMetadata::MAGIC_LENGTH as usize + usize::try_from(footer_payload_length)?;
269 let footer_payload_bytes = footer_bytes
270 .get(start_offset..end_offset)
271 .ok_or_else(|| invalid_data!("Index range is out of bounds."))?;
272 let decompressed_footer_payload_bytes =
273 footer_compression_codec.decompress(footer_payload_bytes.into())?;
274
275 String::from_utf8(decompressed_footer_payload_bytes)
276 .map_err(|src| invalid_data!("Footer is not a valid UTF-8 string").with_source(src))
277 }
278
279 fn from_json_str(string: &str) -> Result<FileMetadata> {
280 serde_json::from_str::<FileMetadata>(string)
281 .map_err(|src| invalid_data!("Given string is not valid JSON").with_source(src))
282 }
283
284 pub(crate) async fn read(file_read: &dyn FileRead, file_length: u64) -> Result<FileMetadata> {
286 if file_length < FileMetadata::MIN_FILE_LENGTH {
287 return Err(invalid_data!(
288 "File length {} is too short to be a Puffin file, expected at least {} bytes",
289 file_length,
290 FileMetadata::MIN_FILE_LENGTH
291 ));
292 }
293
294 let first_four_bytes = file_read.read(0..FileMetadata::MAGIC_LENGTH.into()).await?;
295 FileMetadata::check_magic(&first_four_bytes)?;
296
297 let footer_payload_length =
298 FileMetadata::read_footer_payload_length(file_read, file_length).await?;
299 let footer_bytes =
300 FileMetadata::read_footer_bytes(file_read, file_length, footer_payload_length).await?;
301
302 let magic_length = FileMetadata::MAGIC_LENGTH as usize;
303 FileMetadata::check_magic(&footer_bytes[..magic_length])?;
305 FileMetadata::check_magic(&footer_bytes[footer_bytes.len() - magic_length..])?;
307
308 let footer_payload_str =
309 FileMetadata::extract_footer_payload_as_str(&footer_bytes, footer_payload_length)?;
310
311 FileMetadata::from_json_str(&footer_payload_str)
312 }
313
314 #[allow(dead_code)]
320 pub(crate) async fn read_with_prefetch(
321 file_read: &dyn FileRead,
322 file_length: u64,
323 prefetch_hint: u8,
324 ) -> Result<FileMetadata> {
325 if prefetch_hint > 16 {
326 if prefetch_hint as u64 > file_length {
328 return FileMetadata::read(file_read, file_length).await;
329 }
330
331 let first_four_bytes = file_read.read(0..FileMetadata::MAGIC_LENGTH.into()).await?;
333 FileMetadata::check_magic(&first_four_bytes)?;
334
335 let start = file_length - prefetch_hint as u64;
337 let end = file_length;
338 let footer_bytes = file_read.read(start..end).await?;
339
340 let payload_length_start =
341 footer_bytes.len() - (FileMetadata::FOOTER_STRUCT_LENGTH as usize);
342 let payload_length_end =
343 payload_length_start + (FileMetadata::FOOTER_STRUCT_PAYLOAD_LENGTH_LENGTH as usize);
344 let payload_length_bytes = &footer_bytes[payload_length_start..payload_length_end];
345
346 let mut buf = [0; 4];
347 buf.copy_from_slice(payload_length_bytes);
348 let footer_payload_length = u32::from_le_bytes(buf);
349
350 let footer_length = (footer_payload_length as usize)
354 + FileMetadata::FOOTER_STRUCT_LENGTH as usize
355 + FileMetadata::MAGIC_LENGTH as usize;
356 if footer_length > prefetch_hint as usize {
357 return FileMetadata::read(file_read, file_length).await;
358 }
359
360 let footer_start = footer_bytes.len() - footer_length;
362 let footer_end = footer_bytes.len();
363 let footer_bytes = &footer_bytes[footer_start..footer_end];
364
365 let magic_length = FileMetadata::MAGIC_LENGTH as usize;
366 FileMetadata::check_magic(&footer_bytes[..magic_length])?;
368 FileMetadata::check_magic(&footer_bytes[footer_bytes.len() - magic_length..])?;
370
371 let footer_payload_str =
372 FileMetadata::extract_footer_payload_as_str(footer_bytes, footer_payload_length)?;
373 return FileMetadata::from_json_str(&footer_payload_str);
374 }
375
376 FileMetadata::read(file_read, file_length).await
377 }
378
379 #[inline]
380 pub fn blobs(&self) -> &[BlobMetadata] {
382 &self.blobs
383 }
384
385 #[inline]
386 pub fn properties(&self) -> &HashMap<String, String> {
388 &self.properties
389 }
390}
391
392#[cfg(test)]
393mod tests {
394 use std::collections::HashMap;
395
396 use bytes::Bytes;
397 use tempfile::TempDir;
398
399 use crate::ErrorKind;
400 use crate::io::{FileIO, InputFile};
401 use crate::puffin::metadata::{BlobMetadata, CompressionCodec, FileMetadata};
402 use crate::puffin::test_utils::{
403 empty_footer_payload, empty_footer_payload_bytes, empty_footer_payload_bytes_length_bytes,
404 java_empty_uncompressed_input_file, java_uncompressed_metric_input_file,
405 java_zstd_compressed_metric_input_file, read_file_metadata,
406 read_file_metadata_with_prefetch, uncompressed_metric_file_metadata,
407 zstd_compressed_metric_file_metadata,
408 };
409
410 const INVALID_MAGIC_VALUE: [u8; 4] = [80, 70, 65, 0];
411
412 async fn input_file_with_bytes(temp_dir: &TempDir, slice: &[u8]) -> InputFile {
413 let file_io = FileIO::new_with_fs();
414
415 let path_buf = temp_dir.path().join("abc.puffin");
416 let temp_path = path_buf.to_str().unwrap();
417 let output_file = file_io.new_output(temp_path).unwrap();
418
419 output_file
420 .write(Bytes::copy_from_slice(slice))
421 .await
422 .unwrap();
423
424 output_file.to_input_file()
425 }
426
427 async fn input_file_with_payload(temp_dir: &TempDir, payload_str: &str) -> InputFile {
428 let payload_bytes = payload_str.as_bytes();
429
430 let mut bytes = vec![];
431 bytes.extend(FileMetadata::MAGIC.to_vec());
432 bytes.extend(FileMetadata::MAGIC.to_vec());
433 bytes.extend(payload_bytes);
434 bytes.extend(u32::to_le_bytes(payload_bytes.len() as u32));
435 bytes.extend(vec![0, 0, 0, 0]);
436 bytes.extend(FileMetadata::MAGIC);
437
438 input_file_with_bytes(temp_dir, &bytes).await
439 }
440
441 #[tokio::test]
442 async fn test_file_starting_with_invalid_magic_returns_error() {
443 let temp_dir = TempDir::new().unwrap();
444
445 let mut bytes = vec![];
446 bytes.extend(INVALID_MAGIC_VALUE.to_vec());
447 bytes.extend(FileMetadata::MAGIC.to_vec());
448 bytes.extend(empty_footer_payload_bytes());
449 bytes.extend(empty_footer_payload_bytes_length_bytes());
450 bytes.extend(vec![0, 0, 0, 0]);
451 bytes.extend(FileMetadata::MAGIC);
452
453 let input_file = input_file_with_bytes(&temp_dir, &bytes).await;
454
455 assert_eq!(
456 read_file_metadata(&input_file)
457 .await
458 .unwrap_err()
459 .to_string(),
460 "DataInvalid => Bad magic value: [80, 70, 65, 0] should be [80, 70, 65, 49]",
461 )
462 }
463
464 #[tokio::test]
465 async fn test_file_with_invalid_magic_at_start_of_footer_returns_error() {
466 let temp_dir = TempDir::new().unwrap();
467
468 let mut bytes = vec![];
469 bytes.extend(FileMetadata::MAGIC.to_vec());
470 bytes.extend(INVALID_MAGIC_VALUE.to_vec());
471 bytes.extend(empty_footer_payload_bytes());
472 bytes.extend(empty_footer_payload_bytes_length_bytes());
473 bytes.extend(vec![0, 0, 0, 0]);
474 bytes.extend(FileMetadata::MAGIC);
475
476 let input_file = input_file_with_bytes(&temp_dir, &bytes).await;
477
478 assert_eq!(
479 read_file_metadata(&input_file)
480 .await
481 .unwrap_err()
482 .to_string(),
483 "DataInvalid => Bad magic value: [80, 70, 65, 0] should be [80, 70, 65, 49]",
484 )
485 }
486
487 #[tokio::test]
488 async fn test_file_ending_with_invalid_magic_returns_error() {
489 let temp_dir = TempDir::new().unwrap();
490
491 let mut bytes = vec![];
492 bytes.extend(FileMetadata::MAGIC.to_vec());
493 bytes.extend(FileMetadata::MAGIC.to_vec());
494 bytes.extend(empty_footer_payload_bytes());
495 bytes.extend(empty_footer_payload_bytes_length_bytes());
496 bytes.extend(vec![0, 0, 0, 0]);
497 bytes.extend(INVALID_MAGIC_VALUE);
498
499 let input_file = input_file_with_bytes(&temp_dir, &bytes).await;
500
501 assert_eq!(
502 read_file_metadata(&input_file)
503 .await
504 .unwrap_err()
505 .to_string(),
506 "DataInvalid => Bad magic value: [80, 70, 65, 0] should be [80, 70, 65, 49]",
507 )
508 }
509
510 #[tokio::test]
511 async fn test_encoded_payload_length_larger_than_actual_payload_length_returns_error() {
512 let temp_dir = TempDir::new().unwrap();
513
514 let mut bytes = vec![];
515 bytes.extend(FileMetadata::MAGIC.to_vec());
516 bytes.extend(FileMetadata::MAGIC.to_vec());
517 bytes.extend(empty_footer_payload_bytes());
518 bytes.extend(u32::to_le_bytes(
519 empty_footer_payload_bytes().len() as u32 + 1,
520 ));
521 bytes.extend(vec![0, 0, 0, 0]);
522 bytes.extend(FileMetadata::MAGIC.to_vec());
523
524 let input_file = input_file_with_bytes(&temp_dir, &bytes).await;
525
526 assert_eq!(
527 read_file_metadata(&input_file)
528 .await
529 .unwrap_err()
530 .to_string(),
531 "DataInvalid => Bad magic value: [49, 80, 70, 65] should be [80, 70, 65, 49]",
532 )
533 }
534
535 #[tokio::test]
536 async fn test_encoded_payload_length_smaller_than_actual_payload_length_returns_error() {
537 let temp_dir = TempDir::new().unwrap();
538
539 let mut bytes = vec![];
540 bytes.extend(FileMetadata::MAGIC.to_vec());
541 bytes.extend(FileMetadata::MAGIC.to_vec());
542 bytes.extend(empty_footer_payload_bytes());
543 bytes.extend(u32::to_le_bytes(
544 empty_footer_payload_bytes().len() as u32 - 1,
545 ));
546 bytes.extend(vec![0, 0, 0, 0]);
547 bytes.extend(FileMetadata::MAGIC.to_vec());
548
549 let input_file = input_file_with_bytes(&temp_dir, &bytes).await;
550
551 assert_eq!(
552 read_file_metadata(&input_file)
553 .await
554 .unwrap_err()
555 .to_string(),
556 "DataInvalid => Bad magic value: [70, 65, 49, 123] should be [80, 70, 65, 49]",
557 )
558 }
559
560 #[tokio::test]
561 async fn test_lz4_compressed_footer_returns_error() {
562 let temp_dir = TempDir::new().unwrap();
563
564 let mut bytes = vec![];
565 bytes.extend(FileMetadata::MAGIC.to_vec());
566 bytes.extend(FileMetadata::MAGIC.to_vec());
567 bytes.extend(empty_footer_payload_bytes());
568 bytes.extend(empty_footer_payload_bytes_length_bytes());
569 bytes.extend(vec![0b00000001, 0, 0, 0]);
570 bytes.extend(FileMetadata::MAGIC.to_vec());
571
572 let input_file = input_file_with_bytes(&temp_dir, &bytes).await;
573
574 assert_eq!(
575 read_file_metadata(&input_file)
576 .await
577 .unwrap_err()
578 .to_string(),
579 "FeatureUnsupported => LZ4 decompression is not supported currently",
580 )
581 }
582
583 #[tokio::test]
584 async fn test_unknown_byte_bit_combination_returns_error() {
585 let temp_dir = TempDir::new().unwrap();
586
587 let mut bytes = vec![];
588 bytes.extend(FileMetadata::MAGIC.to_vec());
589 bytes.extend(FileMetadata::MAGIC.to_vec());
590 bytes.extend(empty_footer_payload_bytes());
591 bytes.extend(empty_footer_payload_bytes_length_bytes());
592 bytes.extend(vec![0b00000010, 0, 0, 0]);
593 bytes.extend(FileMetadata::MAGIC.to_vec());
594
595 let input_file = input_file_with_bytes(&temp_dir, &bytes).await;
596
597 assert_eq!(
598 read_file_metadata(&input_file)
599 .await
600 .unwrap_err()
601 .to_string(),
602 "DataInvalid => Unknown flag byte 0 and bit 1 combination",
603 )
604 }
605
606 #[tokio::test]
607 async fn test_non_utf8_string_payload_returns_error() {
608 let temp_dir = TempDir::new().unwrap();
609
610 let payload_bytes: [u8; 4] = [0, 159, 146, 150];
611 let payload_bytes_length_bytes: [u8; 4] = u32::to_le_bytes(payload_bytes.len() as u32);
612
613 let mut bytes = vec![];
614 bytes.extend(FileMetadata::MAGIC.to_vec());
615 bytes.extend(FileMetadata::MAGIC.to_vec());
616 bytes.extend(payload_bytes);
617 bytes.extend(payload_bytes_length_bytes);
618 bytes.extend(vec![0, 0, 0, 0]);
619 bytes.extend(FileMetadata::MAGIC.to_vec());
620
621 let input_file = input_file_with_bytes(&temp_dir, &bytes).await;
622
623 assert_eq!(
624 read_file_metadata(&input_file)
625 .await
626 .unwrap_err()
627 .to_string(),
628 "DataInvalid => Footer is not a valid UTF-8 string, source: invalid utf-8 sequence of 1 bytes from index 1",
629 )
630 }
631
632 #[tokio::test]
633 async fn test_file_shorter_than_minimum_length_returns_error() {
634 let temp_dir = TempDir::new().unwrap();
635
636 let input_file = input_file_with_bytes(&temp_dir, &FileMetadata::MAGIC).await;
638
639 let err = read_file_metadata(&input_file).await.unwrap_err();
640 assert_eq!(err.kind(), ErrorKind::DataInvalid);
641 assert!(
642 err.to_string().contains("too short to be a Puffin file"),
643 "unexpected error: {err}"
644 );
645 }
646
647 #[tokio::test]
648 async fn test_footer_payload_length_larger_than_file_returns_error() {
649 let temp_dir = TempDir::new().unwrap();
650
651 let mut bytes = vec![];
652 bytes.extend(FileMetadata::MAGIC.to_vec());
653 bytes.extend(FileMetadata::MAGIC.to_vec());
654 bytes.extend(empty_footer_payload_bytes());
655 bytes.extend(u32::to_le_bytes(u32::MAX));
657 bytes.extend(vec![0, 0, 0, 0]);
658 bytes.extend(FileMetadata::MAGIC);
659
660 let input_file = input_file_with_bytes(&temp_dir, &bytes).await;
661
662 let err = read_file_metadata(&input_file).await.unwrap_err();
663 assert_eq!(err.kind(), ErrorKind::DataInvalid);
664 assert!(
665 err.to_string().contains("exceeds file length"),
666 "unexpected error: {err}"
667 );
668 }
669
670 #[tokio::test]
671 async fn test_minimal_valid_file_returns_file_metadata() {
672 let temp_dir = TempDir::new().unwrap();
673
674 let mut bytes = vec![];
675 bytes.extend(FileMetadata::MAGIC.to_vec());
676 bytes.extend(FileMetadata::MAGIC.to_vec());
677 bytes.extend(empty_footer_payload_bytes());
678 bytes.extend(empty_footer_payload_bytes_length_bytes());
679 bytes.extend(vec![0, 0, 0, 0]);
680 bytes.extend(FileMetadata::MAGIC);
681
682 let input_file = input_file_with_bytes(&temp_dir, &bytes).await;
683
684 assert_eq!(
685 read_file_metadata(&input_file).await.unwrap(),
686 FileMetadata {
687 blobs: vec![],
688 properties: HashMap::new(),
689 }
690 )
691 }
692
693 #[tokio::test]
694 async fn test_returns_file_metadata_property() {
695 let temp_dir = TempDir::new().unwrap();
696
697 let input_file = input_file_with_payload(
698 &temp_dir,
699 r#"{
700 "blobs" : [ ],
701 "properties" : {
702 "a property" : "a property value"
703 }
704 }"#,
705 )
706 .await;
707
708 assert_eq!(
709 read_file_metadata(&input_file).await.unwrap(),
710 FileMetadata {
711 blobs: vec![],
712 properties: {
713 let mut map = HashMap::new();
714 map.insert("a property".to_string(), "a property value".to_string());
715 map
716 },
717 }
718 )
719 }
720
721 #[tokio::test]
722 async fn test_returns_file_metadata_properties() {
723 let temp_dir = TempDir::new().unwrap();
724
725 let input_file = input_file_with_payload(
726 &temp_dir,
727 r#"{
728 "blobs" : [ ],
729 "properties" : {
730 "a property" : "a property value",
731 "another one": "also with value"
732 }
733 }"#,
734 )
735 .await;
736
737 assert_eq!(
738 read_file_metadata(&input_file).await.unwrap(),
739 FileMetadata {
740 blobs: vec![],
741 properties: {
742 let mut map = HashMap::new();
743 map.insert("a property".to_string(), "a property value".to_string());
744 map.insert("another one".to_string(), "also with value".to_string());
745 map
746 },
747 }
748 )
749 }
750
751 #[tokio::test]
752 async fn test_returns_error_if_blobs_field_is_missing() {
753 let temp_dir = TempDir::new().unwrap();
754
755 let input_file = input_file_with_payload(
756 &temp_dir,
757 r#"{
758 "properties" : {}
759 }"#,
760 )
761 .await;
762
763 assert_eq!(
764 read_file_metadata(&input_file)
765 .await
766 .unwrap_err()
767 .to_string(),
768 format!(
769 "DataInvalid => Given string is not valid JSON, source: missing field `blobs` at line 3 column 13"
770 ),
771 )
772 }
773
774 #[tokio::test]
775 async fn test_returns_error_if_blobs_field_is_bad() {
776 let temp_dir = TempDir::new().unwrap();
777
778 let input_file = input_file_with_payload(
779 &temp_dir,
780 r#"{
781 "blobs" : {}
782 }"#,
783 )
784 .await;
785
786 assert_eq!(
787 read_file_metadata(&input_file)
788 .await
789 .unwrap_err()
790 .to_string(),
791 format!(
792 "DataInvalid => Given string is not valid JSON, source: invalid type: map, expected a sequence at line 2 column 26"
793 ),
794 )
795 }
796
797 #[tokio::test]
798 async fn test_returns_blobs_metadatas() {
799 let temp_dir = TempDir::new().unwrap();
800
801 let input_file = input_file_with_payload(
802 &temp_dir,
803 r#"{
804 "blobs" : [
805 {
806 "type" : "type-a",
807 "fields" : [ 1 ],
808 "snapshot-id" : 14,
809 "sequence-number" : 3,
810 "offset" : 4,
811 "length" : 16
812 },
813 {
814 "type" : "type-bbb",
815 "fields" : [ 2, 3, 4 ],
816 "snapshot-id" : 77,
817 "sequence-number" : 4,
818 "offset" : 21474836470000,
819 "length" : 79834
820 }
821 ]
822 }"#,
823 )
824 .await;
825
826 assert_eq!(
827 read_file_metadata(&input_file).await.unwrap(),
828 FileMetadata {
829 blobs: vec![
830 BlobMetadata {
831 r#type: "type-a".to_string(),
832 fields: vec![1],
833 snapshot_id: 14,
834 sequence_number: 3,
835 offset: 4,
836 length: 16,
837 compression_codec: CompressionCodec::None,
838 properties: HashMap::new(),
839 },
840 BlobMetadata {
841 r#type: "type-bbb".to_string(),
842 fields: vec![2, 3, 4],
843 snapshot_id: 77,
844 sequence_number: 4,
845 offset: 21474836470000,
846 length: 79834,
847 compression_codec: CompressionCodec::None,
848 properties: HashMap::new(),
849 },
850 ],
851 properties: HashMap::new(),
852 }
853 )
854 }
855
856 #[tokio::test]
857 async fn test_returns_properties_in_blob_metadata() {
858 let temp_dir = TempDir::new().unwrap();
859
860 let input_file = input_file_with_payload(
861 &temp_dir,
862 r#"{
863 "blobs" : [
864 {
865 "type" : "type-a",
866 "fields" : [ 1 ],
867 "snapshot-id" : 14,
868 "sequence-number" : 3,
869 "offset" : 4,
870 "length" : 16,
871 "properties" : {
872 "some key" : "some value"
873 }
874 }
875 ]
876 }"#,
877 )
878 .await;
879
880 assert_eq!(
881 read_file_metadata(&input_file).await.unwrap(),
882 FileMetadata {
883 blobs: vec![BlobMetadata {
884 r#type: "type-a".to_string(),
885 fields: vec![1],
886 snapshot_id: 14,
887 sequence_number: 3,
888 offset: 4,
889 length: 16,
890 compression_codec: CompressionCodec::None,
891 properties: {
892 let mut map = HashMap::new();
893 map.insert("some key".to_string(), "some value".to_string());
894 map
895 },
896 }],
897 properties: HashMap::new(),
898 }
899 )
900 }
901
902 #[tokio::test]
903 async fn test_returns_error_if_blobs_fields_value_is_outside_i32_range() {
904 let temp_dir = TempDir::new().unwrap();
905
906 let out_of_i32_range_number: i64 = i32::MAX as i64 + 1;
907
908 let input_file = input_file_with_payload(
909 &temp_dir,
910 &format!(
911 r#"{{
912 "blobs" : [
913 {{
914 "type" : "type-a",
915 "fields" : [ {out_of_i32_range_number} ],
916 "snapshot-id" : 14,
917 "sequence-number" : 3,
918 "offset" : 4,
919 "length" : 16
920 }}
921 ]
922 }}"#
923 ),
924 )
925 .await;
926
927 assert_eq!(
928 read_file_metadata(&input_file)
929 .await
930 .unwrap_err()
931 .to_string(),
932 format!(
933 "DataInvalid => Given string is not valid JSON, source: invalid value: integer `{out_of_i32_range_number}`, expected i32 at line 5 column 51"
934 ),
935 )
936 }
937
938 #[tokio::test]
939 async fn test_returns_errors_if_footer_payload_is_not_encoded_in_json_format() {
940 let temp_dir = TempDir::new().unwrap();
941
942 let input_file = input_file_with_payload(&temp_dir, r#""blobs" = []"#).await;
943
944 assert_eq!(
945 read_file_metadata(&input_file)
946 .await
947 .unwrap_err()
948 .to_string(),
949 "DataInvalid => Given string is not valid JSON, source: invalid type: string \"blobs\", expected struct FileMetadata at line 1 column 7",
950 )
951 }
952
953 #[tokio::test]
954 async fn test_read_file_metadata_of_uncompressed_empty_file() {
955 let input_file = java_empty_uncompressed_input_file();
956
957 let file_metadata = read_file_metadata(&input_file).await.unwrap();
958 assert_eq!(file_metadata, empty_footer_payload())
959 }
960
961 #[tokio::test]
962 async fn test_read_file_metadata_of_uncompressed_metric_data() {
963 let input_file = java_uncompressed_metric_input_file();
964
965 let file_metadata = read_file_metadata(&input_file).await.unwrap();
966 assert_eq!(file_metadata, uncompressed_metric_file_metadata())
967 }
968
969 #[tokio::test]
970 async fn test_read_file_metadata_of_zstd_compressed_metric_data() {
971 let input_file = java_zstd_compressed_metric_input_file();
972
973 let file_metadata = read_file_metadata_with_prefetch(&input_file, 64)
974 .await
975 .unwrap();
976 assert_eq!(file_metadata, zstd_compressed_metric_file_metadata())
977 }
978
979 #[tokio::test]
980 async fn test_read_file_metadata_of_empty_file_with_prefetching() {
981 let input_file = java_empty_uncompressed_input_file();
982 let file_metadata = read_file_metadata_with_prefetch(&input_file, 64)
983 .await
984 .unwrap();
985
986 assert_eq!(file_metadata, empty_footer_payload());
987 }
988
989 #[tokio::test]
990 async fn test_read_file_metadata_of_uncompressed_metric_data_with_prefetching() {
991 let input_file = java_uncompressed_metric_input_file();
992 let file_metadata = read_file_metadata_with_prefetch(&input_file, 64)
993 .await
994 .unwrap();
995
996 assert_eq!(file_metadata, uncompressed_metric_file_metadata());
997 }
998
999 #[tokio::test]
1000 async fn test_read_file_metadata_of_zstd_compressed_metric_data_with_prefetching() {
1001 let input_file = java_zstd_compressed_metric_input_file();
1002 let file_metadata = read_file_metadata_with_prefetch(&input_file, 64)
1003 .await
1004 .unwrap();
1005
1006 assert_eq!(file_metadata, zstd_compressed_metric_file_metadata());
1007 }
1008
1009 #[tokio::test]
1010 async fn test_read_with_incorrect_header_magic() {
1011 let temp_dir = TempDir::new().unwrap();
1012
1013 let prefetch_hint: u8 = 64;
1014 let mut bytes = vec![];
1015 bytes.extend([0x00, 0x00, 0x00, 0x00]);
1017 bytes.extend(vec![0u8; prefetch_hint as usize]);
1019 bytes.extend(FileMetadata::MAGIC);
1021 bytes.extend(empty_footer_payload_bytes());
1022 bytes.extend(empty_footer_payload_bytes_length_bytes());
1023 bytes.extend(vec![0, 0, 0, 0]); bytes.extend(FileMetadata::MAGIC);
1025
1026 let input_file = input_file_with_bytes(&temp_dir, &bytes).await;
1027
1028 assert_eq!(
1029 read_file_metadata(&input_file).await.unwrap_err().kind(),
1030 ErrorKind::DataInvalid,
1031 );
1032 assert_eq!(
1033 read_file_metadata_with_prefetch(&input_file, prefetch_hint)
1034 .await
1035 .unwrap_err()
1036 .kind(),
1037 ErrorKind::DataInvalid,
1038 );
1039 }
1040
1041 #[tokio::test]
1042 async fn test_gzip_compression_allowed_in_metadata() {
1043 let temp_dir = TempDir::new().unwrap();
1044
1045 let payload = r#"{
1048 "blobs": [
1049 {
1050 "type": "test-type",
1051 "fields": [1],
1052 "snapshot-id": 1,
1053 "sequence-number": 1,
1054 "offset": 4,
1055 "length": 10,
1056 "compression-codec": "gzip"
1057 }
1058 ]
1059 }"#;
1060
1061 let input_file = input_file_with_payload(&temp_dir, payload).await;
1062
1063 let result = read_file_metadata(&input_file).await;
1065 assert!(result.is_ok());
1066 let metadata = result.unwrap();
1067 assert_eq!(metadata.blobs.len(), 1);
1068 assert_eq!(
1069 metadata.blobs[0].compression_codec,
1070 CompressionCodec::gzip_default()
1071 );
1072 }
1073}