Skip to main content

iceberg/puffin/
metadata.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18use 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
28/// Human-readable identification of the application writing the file, along with its version.
29/// Example: "Trino version 381"
30pub const CREATED_BY_PROPERTY: &str = "created-by";
31
32/// Metadata about a blob.
33/// For more information, see: <https://iceberg.apache.org/puffin-spec/#blobmetadata>
34#[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    /// See blob types: <https://iceberg.apache.org/puffin-spec/#blob-types>
54    pub fn blob_type(&self) -> &str {
55        &self.r#type
56    }
57
58    #[inline]
59    /// List of field IDs the blob was computed for; the order of items is used to compute sketches stored in the blob.
60    pub fn fields(&self) -> &[i32] {
61        &self.fields
62    }
63
64    #[inline]
65    /// ID of the Iceberg table's snapshot the blob was computed from
66    pub fn snapshot_id(&self) -> i64 {
67        self.snapshot_id
68    }
69
70    #[inline]
71    /// Sequence number of the Iceberg table's snapshot the blob was computed from
72    pub fn sequence_number(&self) -> i64 {
73        self.sequence_number
74    }
75
76    #[inline]
77    /// The offset in the file where the blob contents start
78    pub fn offset(&self) -> u64 {
79        self.offset
80    }
81
82    #[inline]
83    /// The length of the blob stored in the file (after compression, if compressed)
84    pub fn length(&self) -> u64 {
85        self.length
86    }
87
88    #[inline]
89    /// The compression codec used to compress the data
90    pub fn compression_codec(&self) -> CompressionCodec {
91        self.compression_codec
92    }
93
94    #[inline]
95    /// Arbitrary meta-information about the blob
96    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/// Metadata about a puffin file.
131///
132/// For more information, see: <https://iceberg.apache.org/puffin-spec/#filemetadata>
133#[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    /// We use the term FOOTER_STRUCT to refer to the fixed-length portion of the Footer.
146    /// The structure of the Footer specification is illustrated below:
147    ///
148    /// ```text                                             
149    ///        Footer
150    ///        ┌────────────────────┐                 
151    ///        │  Magic (4 bytes)   │                 
152    ///        │                    │                 
153    ///        ├────────────────────┤                 
154    ///        │   FooterPayload    │                 
155    ///        │  (PAYLOAD_LENGTH)  │                 
156    ///        ├────────────────────┤ ◀─┐             
157    ///        │ FooterPayloadSize  │   │             
158    ///        │     (4 bytes)      │   │             
159    ///        ├────────────────────┤                 
160    ///        │  Flags (4 bytes)   │  FOOTER_STRUCT  
161    ///        │                    │                 
162    ///        ├────────────────────┤   │             
163    ///        │  Magic (4 bytes)   │   │             
164    ///        │                    │   │             
165    ///        └────────────────────┘ ◀─┘  
166    /// ```                      
167    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    /// Smallest possible Puffin file: the file header magic, followed by a footer
178    /// holding an empty payload (footer magic + FOOTER_STRUCT).
179    const MIN_FILE_LENGTH: u64 =
180        (FileMetadata::MAGIC_LENGTH as u64) * 2 + (FileMetadata::FOOTER_STRUCT_LENGTH as u64);
181
182    /// Constructs new puffin `FileMetadata`
183    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    /// Returns the file metadata about a Puffin file
285    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        // check first four bytes of footer
304        FileMetadata::check_magic(&footer_bytes[..magic_length])?;
305        // check last four bytes of footer
306        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    /// Reads file_metadata in puffin file with a prefetch hint
315    ///
316    /// `prefetch_hint` is used to try to fetch the entire footer in one read. If
317    /// the entire footer isn't fetched in one read the function will call the regular
318    /// read option.
319    #[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            // Hint cannot be larger than input file
327            if prefetch_hint as u64 > file_length {
328                return FileMetadata::read(file_read, file_length).await;
329            }
330
331            // Validate file header magic
332            let first_four_bytes = file_read.read(0..FileMetadata::MAGIC_LENGTH.into()).await?;
333            FileMetadata::check_magic(&first_four_bytes)?;
334
335            // Read footer based on prefetch hint
336            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            // If the (footer payload length + FOOTER_STRUCT_LENGTH + MAGIC_LENGTH) is greater
351            // than the fetched footer then you can have it read regularly from a read with no
352            // prefetch while passing in the footer_payload_length.
353            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            // Read footer bytes
361            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            // check first four bytes of footer
367            FileMetadata::check_magic(&footer_bytes[..magic_length])?;
368            // check last four bytes of footer
369            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    /// Metadata about blobs in file
381    pub fn blobs(&self) -> &[BlobMetadata] {
382        &self.blobs
383    }
384
385    #[inline]
386    /// Arbitrary meta-information, like writer identification/version.
387    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        // Only the file header magic, nothing else.
637        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        // Declared footer payload length is far larger than the file itself.
656        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        // Invalid header magic
1016        bytes.extend([0x00, 0x00, 0x00, 0x00]);
1017        // Intentionally keep file size larger than prefetch_hint.
1018        bytes.extend(vec![0u8; prefetch_hint as usize]);
1019        // Valid footer: magic + payload + footer struct
1020        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]); // flags
1024        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        // Create a JSON payload with Gzip compression codec
1046        // Metadata should be readable, but accessing the blob will fail
1047        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        // Reading metadata should succeed (lazy validation)
1064        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}