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
12pub struct Grib2SectionStream<R> {
40 reader: R,
41 whole_size: usize,
42 rest_size: usize,
43}
44
45impl<R> Grib2SectionStream<R> {
46 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 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 fn read_sect0(&mut self) -> Result<Option<(usize, Indicator)>, ParseError>;
175
176 fn read_sect8(&mut self) -> Result<Option<()>, ParseError>;
178
179 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]; 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}