Skip to main content

grib/
reader.rs

1use std::io::{self, Read, Seek, SeekFrom};
2
3use crate::{SectionBody, SectionInfo, datatypes::*, error::*, helpers::read_as};
4
5const SECT0_IS_MAGIC: &[u8] = b"GRIB";
6const SECT0_IS_MAGIC_SIZE: usize = SECT0_IS_MAGIC.len();
7const SECT0_IS_SIZE: usize = 16;
8const SECT_HEADER_SIZE: usize = 5;
9const SECT8_ES_MAGIC: &[u8] = b"7777";
10pub(crate) const SECT8_ES_SIZE: usize = SECT8_ES_MAGIC.len();
11
12/// # Example
13/// ```
14/// use grib::{Grib2SectionStream, Indicator, SectionBody, SectionInfo, SeekableGrib2Reader};
15///
16/// fn main() -> Result<(), Box<dyn std::error::Error>> {
17///     let f = std::fs::File::open(
18///         "testdata/icon_global_icosahedral_single-level_2021112018_000_TOT_PREC.grib2",
19///     )?;
20///     let f = std::io::BufReader::new(f);
21///     let grib2_reader = SeekableGrib2Reader::new(f);
22///
23///     let mut sect_stream = Grib2SectionStream::new(grib2_reader);
24///     assert_eq!(
25///         sect_stream.next(),
26///         Some(Ok(SectionInfo {
27///             num: 0,
28///             offset: 0,
29///             size: 16,
30///             body: Some(SectionBody::Section0(Indicator {
31///                 discipline: 0,
32///                 total_length: 193,
33///             })),
34///         }))
35///     );
36///     Ok(())
37/// }
38/// ```
39pub struct Grib2SectionStream<R> {
40    reader: R,
41    whole_size: usize,
42    rest_size: usize,
43}
44
45impl<R> Grib2SectionStream<R> {
46    /// # Example
47    /// ```
48    /// use grib::{Grib2SectionStream, SeekableGrib2Reader};
49    ///
50    /// fn main() -> Result<(), Box<dyn std::error::Error>> {
51    ///     let f = std::fs::File::open(
52    ///         "testdata/icon_global_icosahedral_single-level_2021112018_000_TOT_PREC.grib2",
53    ///     )?;
54    ///     let mut f = std::io::BufReader::new(f);
55    ///     let grib2_reader = SeekableGrib2Reader::new(f);
56    ///     let _sect_stream = Grib2SectionStream::new(grib2_reader);
57    ///     Ok(())
58    /// }
59    /// ```
60    pub fn new(reader: R) -> Self {
61        Self {
62            reader,
63            whole_size: 0,
64            rest_size: 0,
65        }
66    }
67
68    pub fn into_reader(self) -> R {
69        self.reader
70    }
71}
72
73impl<R> Grib2SectionStream<R>
74where
75    R: Grib2Read,
76{
77    #[inline]
78    fn next_sect0(&mut self) -> Option<Result<SectionInfo, ParseError>> {
79        if self.whole_size == 0 {
80            // if the offset value is left at the initial value, reset it to the current
81            // position
82            let result = self.reset_pos();
83            if let Err(e) = result {
84                return Some(Err(ParseError::ReadError(format!(
85                    "resetting the initial position failed: {e}"
86                ))));
87            }
88        }
89        let result = self
90            .reader
91            .read_sect0()
92            .transpose()?
93            .map(|(offset, indicator)| {
94                self.whole_size += offset;
95                let offset = self.whole_size;
96                let message_size = indicator.total_length as usize;
97                self.whole_size += message_size;
98                let sect_info = SectionInfo {
99                    num: 0,
100                    offset,
101                    size: SECT0_IS_SIZE,
102                    body: Some(SectionBody::Section0(indicator)),
103                };
104                self.rest_size = message_size - SECT0_IS_SIZE;
105                sect_info
106            });
107        Some(result)
108    }
109
110    fn reset_pos(&mut self) -> Result<(), io::Error> {
111        let pos = self.reader.stream_position()?;
112        self.whole_size = pos as usize;
113        Ok(())
114    }
115
116    #[inline]
117    fn next_sect8(&mut self) -> Option<Result<SectionInfo, ParseError>> {
118        let result = self.reader.read_sect8().transpose()?.map(|_| {
119            let sect_info = SectionInfo {
120                num: 8,
121                offset: self.whole_size - self.rest_size,
122                size: SECT8_ES_SIZE,
123                body: None,
124            };
125            self.rest_size -= SECT8_ES_SIZE;
126            sect_info
127        });
128        Some(result)
129    }
130
131    #[inline]
132    fn next_sect(&mut self) -> Option<Result<SectionInfo, ParseError>> {
133        let result = self.reader.read_sect_header().transpose()?;
134        match result {
135            Ok(header) => {
136                let offset = self.whole_size - self.rest_size;
137                match self.reader.read_sect_payload(&header) {
138                    Ok(body) => {
139                        let body = Some(body);
140                        let (size, num) = header;
141                        self.rest_size -= size;
142                        Some(Ok(SectionInfo {
143                            num,
144                            offset,
145                            size,
146                            body,
147                        }))
148                    }
149                    Err(e) => Some(Err(e)),
150                }
151            }
152            Err(e) => Some(Err(e)),
153        }
154    }
155}
156
157impl<R> Iterator for Grib2SectionStream<R>
158where
159    R: Grib2Read,
160{
161    type Item = Result<SectionInfo, ParseError>;
162
163    fn next(&mut self) -> Option<Self::Item> {
164        match self.rest_size {
165            0 => self.next_sect0(),
166            SECT8_ES_SIZE => self.next_sect8(),
167            _ => self.next_sect(),
168        }
169    }
170}
171
172pub trait Grib2Read: Read + Seek {
173    /// Reads Section 0.
174    fn read_sect0(&mut self) -> Result<Option<(usize, Indicator)>, ParseError>;
175
176    /// Reads Section 8.
177    fn read_sect8(&mut self) -> Result<Option<()>, ParseError>;
178
179    /// Reads a common header for Sections 1-7 and returns the section
180    /// size and number.
181    fn read_sect_as_slice(&mut self, sect: &SectionInfo) -> Result<Vec<u8>, ParseError>;
182    fn read_sect_header(&mut self) -> Result<Option<SectHeader>, ParseError>;
183    fn read_sect_payload(&mut self, header: &SectHeader) -> Result<SectionBody, ParseError>;
184    fn read_sect6_payload(&mut self, size: usize) -> Result<SectionBody, ParseError>;
185    fn skip_sect7_payload(&mut self, size: usize) -> Result<SectionBody, ParseError>;
186    fn read_slice_without_offset_check(&mut self, size: usize) -> Result<Box<[u8]>, ParseError>;
187}
188
189pub struct SeekableGrib2Reader<R> {
190    reader: R,
191}
192
193impl<R> SeekableGrib2Reader<R> {
194    pub fn new(r: R) -> Self {
195        Self { reader: r }
196    }
197}
198
199impl<R: Read> Read for SeekableGrib2Reader<R> {
200    fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
201        self.reader.read(buf)
202    }
203
204    fn read_exact(&mut self, buf: &mut [u8]) -> io::Result<()> {
205        self.reader.read_exact(buf)
206    }
207}
208
209impl<S: Seek> Seek for SeekableGrib2Reader<S> {
210    fn seek(&mut self, pos: SeekFrom) -> io::Result<u64> {
211        self.reader.seek(pos)
212    }
213}
214
215macro_rules! check_size {
216    ($size:expr, $expected_size:expr) => {{
217        if $size == 0 {
218            return Ok(None);
219        }
220        if $size != $expected_size {
221            return Err(io::Error::new(
222                io::ErrorKind::UnexpectedEof,
223                "failed to fill whole buffer",
224            )
225            .into());
226        }
227    }};
228}
229
230impl<R: Read + Seek> Grib2Read for SeekableGrib2Reader<R> {
231    fn read_sect0(&mut self) -> Result<Option<(usize, Indicator)>, ParseError> {
232        let mut buf = [0; 4096];
233        let mut offset = 0;
234
235        loop {
236            let size = self.read(&mut buf[..])?;
237            if size < SECT0_IS_SIZE {
238                return Ok(None);
239            }
240            let next_offset = size - SECT0_IS_SIZE + 1;
241            for pos in 0..next_offset {
242                if &buf[pos..pos + SECT0_IS_MAGIC_SIZE] == SECT0_IS_MAGIC {
243                    offset += pos;
244                    self.seek(SeekFrom::Current(
245                        (pos + SECT0_IS_SIZE) as i64 - size as i64,
246                    ))?;
247
248                    let indicator = Indicator::from_slice(&buf[pos..pos + SECT0_IS_SIZE])?;
249                    return Ok(Some((offset, indicator)));
250                }
251            }
252            self.seek(SeekFrom::Current(next_offset as i64 - size as i64))?;
253            offset += next_offset;
254        }
255    }
256
257    fn read_sect8(&mut self) -> Result<Option<()>, ParseError> {
258        let mut buf = [0; SECT8_ES_SIZE];
259        let size = self.read(&mut buf[..])?;
260        check_size!(size, buf.len());
261
262        if buf[..] != SECT8_ES_MAGIC[..] {
263            return Err(ParseError::EndSectionMismatch);
264        }
265
266        Ok(Some(()))
267    }
268
269    fn read_sect_as_slice(&mut self, sect: &SectionInfo) -> Result<Vec<u8>, ParseError> {
270        self.seek(SeekFrom::Start(sect.offset as u64))?;
271
272        let mut buf = vec![0; sect.size];
273        self.read_exact(buf.as_mut_slice())?;
274
275        Ok(buf)
276    }
277
278    fn read_sect_header(&mut self) -> Result<Option<SectHeader>, ParseError> {
279        let mut buf = [0; SECT_HEADER_SIZE];
280        let size = self.read(&mut buf[..])?;
281        check_size!(size, buf.len());
282
283        let sect_size = read_as!(u32, buf, 0) as usize;
284        let sect_num = buf[4];
285
286        Ok(Some((sect_size, sect_num)))
287    }
288
289    fn read_sect_payload(&mut self, header: &SectHeader) -> Result<SectionBody, ParseError> {
290        let (size, num) = header;
291        let body_size = size - SECT_HEADER_SIZE;
292        let body = match num {
293            1 => SectionBody::Section1(Identification::from_payload(
294                self.read_slice_without_offset_check(body_size)?,
295            )?),
296            2 => SectionBody::Section2(LocalUse::from_payload(
297                self.read_slice_without_offset_check(body_size)?,
298            )),
299            3 => SectionBody::Section3(GridDefinition::from_payload(
300                self.read_slice_without_offset_check(body_size)?,
301            )?),
302            4 => SectionBody::Section4(ProdDefinition::from_payload(
303                self.read_slice_without_offset_check(body_size)?,
304            )?),
305            5 => SectionBody::Section5(ReprDefinition::from_payload(
306                self.read_slice_without_offset_check(body_size)?,
307            )?),
308            6 => self.read_sect6_payload(body_size)?,
309            7 => self.skip_sect7_payload(body_size)?,
310            _ => return Err(ParseError::UnknownSectionNumber(*num)),
311        };
312
313        Ok(body)
314    }
315
316    fn read_sect6_payload(&mut self, body_size: usize) -> Result<SectionBody, ParseError> {
317        let mut buf = [0; 1]; // octet 6
318        self.read_exact(&mut buf[..])?;
319
320        let len_extra = body_size - buf.len();
321        if len_extra > 0 {
322            let mut buf = vec![0; len_extra];
323            self.read_exact(&mut buf[..])?;
324        }
325
326        Ok(SectionBody::Section6(BitMap {
327            bitmap_indicator: buf[0],
328        }))
329    }
330
331    fn skip_sect7_payload(&mut self, body_size: usize) -> Result<SectionBody, ParseError> {
332        self.seek(SeekFrom::Current(body_size as i64))?;
333
334        Ok(SectionBody::Section7)
335    }
336
337    fn read_slice_without_offset_check(&mut self, size: usize) -> Result<Box<[u8]>, ParseError> {
338        let mut buf = vec![0; size];
339        self.read_exact(&mut buf[..])?;
340        Ok(buf.into_boxed_slice())
341    }
342}
343
344type SectHeader = (usize, u8);
345
346#[cfg(test)]
347mod tests {
348    use std::io::{Cursor, Write};
349
350    use super::*;
351    use crate::test_utils::{data::grib2::DWD_ICON, decompress_to_vec};
352
353    #[test]
354    fn read_one_grib2_message() -> Result<(), Box<dyn std::error::Error>> {
355        let f = std::fs::File::open(DWD_ICON)?;
356        let f = std::io::BufReader::new(f);
357
358        let grib2_reader = SeekableGrib2Reader::new(f);
359        let sect_stream = Grib2SectionStream::new(grib2_reader);
360        assert_eq!(
361            sect_stream
362                .take(10)
363                .map(|result| result.map(|sect| (sect.num, sect.offset, sect.size)))
364                .collect::<Vec<_>>(),
365            vec![
366                Ok((0, 0, 16)),
367                Ok((1, 16, 21)),
368                Ok((2, 37, 27)),
369                Ok((3, 64, 35)),
370                Ok((4, 99, 58)),
371                Ok((5, 157, 21)),
372                Ok((6, 178, 6)),
373                Ok((7, 184, 5)),
374                Ok((8, 189, 4)),
375            ]
376        );
377
378        Ok(())
379    }
380
381    #[test]
382    fn read_multiple_grib2_messages() -> Result<(), Box<dyn std::error::Error>> {
383        let buf = decompress_to_vec(DWD_ICON)?;
384        let repeated_message = buf.repeat(2);
385        let f = Cursor::new(repeated_message);
386
387        let grib2_reader = SeekableGrib2Reader::new(f);
388        let sect_stream = Grib2SectionStream::new(grib2_reader);
389        assert_eq!(
390            sect_stream
391                .take(19)
392                .map(|result| result.map(|sect| (sect.num, sect.offset, sect.size)))
393                .collect::<Vec<_>>(),
394            vec![
395                Ok((0, 0, 16)),
396                Ok((1, 16, 21)),
397                Ok((2, 37, 27)),
398                Ok((3, 64, 35)),
399                Ok((4, 99, 58)),
400                Ok((5, 157, 21)),
401                Ok((6, 178, 6)),
402                Ok((7, 184, 5)),
403                Ok((8, 189, 4)),
404                Ok((0, 193, 16)),
405                Ok((1, 209, 21)),
406                Ok((2, 230, 27)),
407                Ok((3, 257, 35)),
408                Ok((4, 292, 58)),
409                Ok((5, 350, 21)),
410                Ok((6, 371, 6)),
411                Ok((7, 377, 5)),
412                Ok((8, 382, 4)),
413            ]
414        );
415
416        Ok(())
417    }
418
419    #[test]
420    fn read_grib2_message_with_incomplete_section_0() -> Result<(), Box<dyn std::error::Error>> {
421        let mut buf = decompress_to_vec(DWD_ICON)?;
422        let mut extra_bytes = "extra".as_bytes().to_vec();
423        buf.append(&mut extra_bytes);
424        let f = Cursor::new(buf);
425
426        let grib2_reader = SeekableGrib2Reader::new(f);
427        let sect_stream = Grib2SectionStream::new(grib2_reader);
428        assert_eq!(
429            sect_stream
430                .take(10)
431                .map(|result| result.map(|sect| (sect.num, sect.offset, sect.size)))
432                .collect::<Vec<_>>(),
433            vec![
434                Ok((0, 0, 16)),
435                Ok((1, 16, 21)),
436                Ok((2, 37, 27)),
437                Ok((3, 64, 35)),
438                Ok((4, 99, 58)),
439                Ok((5, 157, 21)),
440                Ok((6, 178, 6)),
441                Ok((7, 184, 5)),
442                Ok((8, 189, 4)),
443            ]
444        );
445
446        Ok(())
447    }
448
449    #[test]
450    fn read_grib2_message_with_incomplete_section_1() -> Result<(), Box<dyn std::error::Error>> {
451        let mut buf = decompress_to_vec(DWD_ICON)?;
452        let mut message_2_bytes = buf[..(SECT0_IS_SIZE + 1)].to_vec();
453        buf.append(&mut message_2_bytes);
454        let f = Cursor::new(buf);
455
456        let grib2_reader = SeekableGrib2Reader::new(f);
457        let sect_stream = Grib2SectionStream::new(grib2_reader);
458        assert_eq!(
459            sect_stream
460                .take(19)
461                .map(|result| result.map(|sect| (sect.num, sect.offset, sect.size)))
462                .collect::<Vec<_>>(),
463            vec![
464                Ok((0, 0, 16)),
465                Ok((1, 16, 21)),
466                Ok((2, 37, 27)),
467                Ok((3, 64, 35)),
468                Ok((4, 99, 58)),
469                Ok((5, 157, 21)),
470                Ok((6, 178, 6)),
471                Ok((7, 184, 5)),
472                Ok((8, 189, 4)),
473                Ok((0, 193, 16)),
474                Err(ParseError::ReadError(
475                    "failed to fill whole buffer".to_owned()
476                ))
477            ]
478        );
479
480        Ok(())
481    }
482
483    #[test]
484    fn read_grib2_message_with_incomplete_section_8() -> Result<(), Box<dyn std::error::Error>> {
485        let buf = decompress_to_vec(DWD_ICON)?;
486        let mut repeated_message = buf.repeat(2);
487        repeated_message.pop();
488        let f = Cursor::new(repeated_message);
489
490        let grib2_reader = SeekableGrib2Reader::new(f);
491        let sect_stream = Grib2SectionStream::new(grib2_reader);
492        assert_eq!(
493            sect_stream
494                .take(19)
495                .map(|result| result.map(|sect| (sect.num, sect.offset, sect.size)))
496                .collect::<Vec<_>>(),
497            vec![
498                Ok((0, 0, 16)),
499                Ok((1, 16, 21)),
500                Ok((2, 37, 27)),
501                Ok((3, 64, 35)),
502                Ok((4, 99, 58)),
503                Ok((5, 157, 21)),
504                Ok((6, 178, 6)),
505                Ok((7, 184, 5)),
506                Ok((8, 189, 4)),
507                Ok((0, 193, 16)),
508                Ok((1, 209, 21)),
509                Ok((2, 230, 27)),
510                Ok((3, 257, 35)),
511                Ok((4, 292, 58)),
512                Ok((5, 350, 21)),
513                Ok((6, 371, 6)),
514                Ok((7, 377, 5)),
515                Err(ParseError::ReadError(
516                    "failed to fill whole buffer".to_owned()
517                ))
518            ]
519        );
520
521        Ok(())
522    }
523
524    fn create_grib2_message_starting_from_non_zero_position(
525        header: &[u8],
526    ) -> Result<Vec<u8>, Box<dyn std::error::Error>> {
527        let mut buf = Vec::new();
528        buf.write_all(header)?;
529
530        let f = std::fs::File::open(DWD_ICON)?;
531        let mut f = std::io::BufReader::new(f);
532        f.read_to_end(&mut buf)?;
533
534        Ok(buf)
535    }
536
537    #[test]
538    fn read_grib2_message_starting_from_non_zero_position_after_seeking()
539    -> Result<(), Box<dyn std::error::Error>> {
540        let header_bytes_skipped = b"HEADER TO BE SKIPPED\n";
541        let buf = create_grib2_message_starting_from_non_zero_position(header_bytes_skipped)?;
542
543        let mut f = Cursor::new(buf);
544        f.seek(SeekFrom::Current(header_bytes_skipped.len() as i64))?;
545
546        let grib2_reader = SeekableGrib2Reader::new(f);
547        let sect_stream = Grib2SectionStream::new(grib2_reader);
548        assert_eq!(
549            sect_stream
550                .take(10)
551                .map(|result| result.map(|sect| (sect.num, sect.offset, sect.size)))
552                .collect::<Vec<_>>(),
553            vec![
554                Ok((0, 21, 16)),
555                Ok((1, 37, 21)),
556                Ok((2, 58, 27)),
557                Ok((3, 85, 35)),
558                Ok((4, 120, 58)),
559                Ok((5, 178, 21)),
560                Ok((6, 199, 6)),
561                Ok((7, 205, 5)),
562                Ok((8, 210, 4)),
563            ]
564        );
565
566        Ok(())
567    }
568
569    macro_rules! test_reading_message_starting_from_non_zero_position {
570        ($(($name:ident, $header:expr, $base_offset:expr),)*) => ($(
571            #[test]
572            fn $name() -> Result<(), Box<dyn std::error::Error>> {
573            let header_bytes_skipped = $header;
574            let buf = create_grib2_message_starting_from_non_zero_position(&header_bytes_skipped)?;
575
576            let f = Cursor::new(buf);
577            let grib2_reader = SeekableGrib2Reader::new(f);
578            let sect_stream = Grib2SectionStream::new(grib2_reader);
579            assert_eq!(
580                sect_stream
581                    .take(10)
582                    .map(|result| result.map(|sect| (sect.num, sect.offset, sect.size)))
583                    .collect::<Vec<_>>(),
584                vec![
585                    Ok((0, $base_offset + 0, 16)),
586                    Ok((1, $base_offset + 16, 21)),
587                    Ok((2, $base_offset + 37, 27)),
588                    Ok((3, $base_offset + 64, 35)),
589                    Ok((4, $base_offset + 99, 58)),
590                    Ok((5, $base_offset + 157, 21)),
591                    Ok((6, $base_offset + 178, 6)),
592                    Ok((7, $base_offset + 184, 5)),
593                    Ok((8, $base_offset + 189, 4)),
594                ]
595            );
596
597            Ok(())
598            }
599        )*);
600    }
601
602    test_reading_message_starting_from_non_zero_position! {
603        (reading_message_using_read_sect0_0th_iteration, [0; 16], 16),
604        (reading_message_using_end_of_read_sect0_0th_iteration, [0; 4096 - 16], 4096 - 16),
605        (reading_message_using_read_sect0_0th_and_1st_iterations, [0; 4096 - 15], 4096 - 15),
606        (reading_message_using_read_sect0_1st_iteration, [0; 4096], 4096),
607    }
608}