Skip to main content

bgpkit_parser/parser/iters/
recovery.rs

1//! Opt-in MRT framing recovery.
2//!
3//! Recovery is deliberately separate from the default iterators. It never attempts to
4//! reconstruct a damaged record: bytes are skipped until a conservatively validated chain of
5//! records is found, and the skipped range is reported as a [`RecoveryEvent::Gap`].
6//!
7//! Damage is classified before scanning. When a record frames correctly — its header and
8//! declared length were consumed exactly — but its body fails to parse, and intact records
9//! (or a clean end of stream) follow at the declared boundary, exactly that record is
10//! skipped without scanning. Damage that extends to the end of the stream is reported as a
11//! terminal gap rather than an error, so trailing truncation — the most common real-world
12//! corruption — still yields every intact record plus an explicit account of the discarded
13//! tail. A [`RecoveryError`] is reserved for I/O failures, unsupported input, and scan
14//! windows exhausted without finding a boundary mid-stream.
15//!
16//! The undamaged fast path reads straight from the underlying reader; bytes are only
17//! buffered while a recovery scan is in progress.
18
19use crate::models::{Bgp4MpType, BgpElem, EntryType, MrtRecord};
20use crate::parser::iters::{record_matches_filters, write_mrt_core_dump};
21use crate::parser::mrt::messages::bgp4mp::uses_zebra_compat;
22use crate::parser::mrt::mrt_header::parse_common_header_with_bytes;
23use crate::parser::mrt::mrt_record::{
24    chunk_mrt_record, parse_mrt_record_with_zebra_compat, raw_record_uses_zebra_compat,
25};
26use crate::parser::{
27    BgpkitParser, Elementor, Filter, ParserError, ParserErrorWithBytes, ParserOptions,
28};
29use crate::Filterable;
30use bytes::Bytes;
31use std::fmt::{Display, Formatter};
32use std::io::{self, Read};
33
34const DEFAULT_MAX_SCAN_BYTES: usize = 1024 * 1024;
35const DEFAULT_CONFIRMATION_RECORDS: u8 = 3;
36const MAX_RECOVERY_RECORD_LEN: u32 = 65_599;
37const SCAN_FILL_CHUNK: usize = 8_192;
38
39/// Settings for opt-in MRT framing recovery.
40#[derive(Debug, Clone, Copy, PartialEq, Eq)]
41pub struct RecoveryConfig {
42    max_scan_bytes: usize,
43    confirmation_records: u8,
44}
45
46impl Default for RecoveryConfig {
47    fn default() -> Self {
48        Self {
49            max_scan_bytes: DEFAULT_MAX_SCAN_BYTES,
50            confirmation_records: DEFAULT_CONFIRMATION_RECORDS,
51        }
52    }
53}
54
55impl RecoveryConfig {
56    /// Set the maximum number of bytes searched after a damaged record.
57    pub const fn with_max_scan_bytes(mut self, max_scan_bytes: usize) -> Self {
58        self.max_scan_bytes = max_scan_bytes;
59        self
60    }
61
62    /// Set the number of consecutive records required to confirm a recovered boundary.
63    ///
64    /// A value of zero is treated as one.
65    pub const fn with_confirmation_records(mut self, confirmation_records: u8) -> Self {
66        self.confirmation_records = confirmation_records;
67        self
68    }
69
70    pub const fn max_scan_bytes(&self) -> usize {
71        self.max_scan_bytes
72    }
73
74    pub const fn confirmation_records(&self) -> u8 {
75        self.confirmation_records
76    }
77}
78
79/// Evidence used to validate the first record following a recovered gap.
80#[derive(Debug, Clone, Copy, PartialEq, Eq)]
81#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
82#[non_exhaustive]
83pub enum RecoveryEvidence {
84    /// A deprecated MRT Type-5 record and its confirmation chain parsed structurally.
85    LegacyMrtChain,
86    /// A BGP4MP message contained an exact embedded BGP marker and length.
87    BgpMarkerChain,
88    /// A BGP4MP state-change record, which has no embedded BGP message header.
89    Bgp4MpStateChangeChain,
90    /// The damaged record's framing was intact: complete MRT records of any type (or a
91    /// clean end of stream) followed at its declared end offset, so exactly that record
92    /// was skipped without scanning.
93    AlignedRecordChain,
94    /// No boundary was validated before the stream ended; the gap extends to the end of
95    /// the input.
96    EndOfStream,
97}
98
99/// A byte range discarded while restoring MRT record framing.
100#[derive(Debug, Clone, PartialEq, Eq)]
101#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
102pub struct RecoveryGap {
103    /// Inclusive offset in the decompressed MRT byte stream.
104    pub start_offset: u64,
105    /// Exclusive offset in the decompressed MRT byte stream.
106    pub end_offset: u64,
107    /// Error raised while parsing at `start_offset`.
108    pub cause: String,
109    /// Structural evidence used to accept `end_offset` as a new boundary.
110    pub evidence: RecoveryEvidence,
111    /// Number of consecutive records validated at the recovered boundary. Zero when the
112    /// gap ends at the end of the stream.
113    pub confirmed_records: u8,
114}
115
116impl RecoveryGap {
117    /// Return the number of skipped bytes, or zero for an invalid inverted range.
118    pub const fn skipped_bytes(&self) -> u64 {
119        self.end_offset.saturating_sub(self.start_offset)
120    }
121}
122
123/// An item produced by a recovering iterator.
124#[derive(Debug)]
125pub enum RecoveryEvent<T> {
126    Item(T),
127    Gap(RecoveryGap),
128}
129
130/// A framing error for which no recovery boundary was found within the scan window, an
131/// I/O failure, or unsupported input.
132///
133/// Damage that extends to the end of the stream is reported as a terminal
134/// [`RecoveryEvent::Gap`] instead of this error.
135#[derive(Debug)]
136pub struct RecoveryError {
137    pub offset: u64,
138    pub scanned_bytes: u64,
139    pub error: ParserErrorWithBytes,
140}
141
142impl Display for RecoveryError {
143    fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
144        write!(
145            f,
146            "MRT recovery failed at decompressed offset {} after scanning {} bytes: {}",
147            self.offset, self.scanned_bytes, self.error
148        )
149    }
150}
151
152impl std::error::Error for RecoveryError {
153    fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
154        Some(&self.error)
155    }
156}
157
158/// Iterator over parsed MRT records and explicit recovery gaps.
159pub struct RecoveringRecordIterator<R> {
160    reader: CarryoverReader<R>,
161    config: RecoveryConfig,
162    filters: Vec<Filter>,
163    elementor: Elementor,
164    options: ParserOptions,
165    core_dump: bool,
166    unsupported_input: Option<String>,
167    finished: bool,
168}
169
170impl<R> RecoveringRecordIterator<R> {
171    pub(crate) fn new(parser: BgpkitParser<R>, config: RecoveryConfig) -> Self {
172        let unsupported_input = parser.text_dump_iter.is_some().then(|| {
173            "text-dump parsers have no MRT record representation; iterate elements instead"
174                .to_string()
175        });
176        Self {
177            reader: CarryoverReader::new(parser.reader),
178            config,
179            filters: parser.filters,
180            elementor: Elementor::new(),
181            options: parser.options,
182            core_dump: parser.core_dump,
183            unsupported_input,
184            finished: false,
185        }
186    }
187}
188
189impl<R: Read> Iterator for RecoveringRecordIterator<R> {
190    type Item = Result<RecoveryEvent<MrtRecord>, RecoveryError>;
191
192    fn next(&mut self) -> Option<Self::Item> {
193        if self.finished {
194            return None;
195        }
196        if let Some(message) = self.unsupported_input.take() {
197            self.finished = true;
198            return Some(Err(RecoveryError {
199                offset: 0,
200                scanned_bytes: 0,
201                error: ParserErrorWithBytes::from(ParserError::Unsupported(message)),
202            }));
203        }
204
205        loop {
206            let record_start = self.reader.position();
207            let raw_record = match chunk_mrt_record(&mut self.reader) {
208                Ok(raw_record) => raw_record,
209                Err(error) if matches!(error.error, ParserError::EofExpected) => {
210                    self.finished = true;
211                    return None;
212                }
213                Err(error) if is_non_eof_io_error(&error.error) => {
214                    self.finished = true;
215                    return Some(Err(RecoveryError {
216                        offset: record_start,
217                        scanned_bytes: 0,
218                        error,
219                    }));
220                }
221                Err(error) => return self.recover(record_start, None, error),
222            };
223
224            let used_zebra_compat = raw_record_uses_zebra_compat(&raw_record);
225            match raw_record.clone().parse() {
226                Ok(record) => {
227                    if used_zebra_compat {
228                        self.options.warn_zebra_compat_once();
229                    }
230                    if record_matches_filters(&record, &self.filters, &mut self.elementor) {
231                        return Some(Ok(RecoveryEvent::Item(record)));
232                    }
233                }
234                Err(error) => {
235                    // The header and declared length were consumed exactly, so the
236                    // stream may still be aligned even though the body is unparsable.
237                    let framed_end = self.reader.position();
238                    let error = ParserErrorWithBytes {
239                        error,
240                        bytes: Some(raw_record.raw_bytes().to_vec()),
241                    };
242                    return self.recover(record_start, Some(framed_end), error);
243                }
244            }
245        }
246    }
247}
248
249impl<R: Read> RecoveringRecordIterator<R> {
250    fn recover(
251        &mut self,
252        record_start: u64,
253        framed_end: Option<u64>,
254        error: ParserErrorWithBytes,
255    ) -> Option<Result<RecoveryEvent<MrtRecord>, RecoveryError>> {
256        write_mrt_core_dump(self.core_dump, error.bytes.clone());
257        let consumed = error.bytes.clone().unwrap_or_default();
258        debug_assert_eq!(record_start + consumed.len() as u64, self.reader.position());
259        let confirmations = self.config.confirmation_records.max(1);
260        let max_scan_bytes = self.config.max_scan_bytes;
261        let mut window = ReplayReader::seeded(&mut self.reader, consumed, record_start);
262
263        if let Some(framed_end) = framed_end {
264            match confirm_aligned_boundary(&mut window, framed_end, confirmations) {
265                Ok(Some(confirmed_records)) => {
266                    let leftover = window.into_leftover(framed_end);
267                    self.reader.resume_with(leftover, framed_end);
268                    return Some(Ok(RecoveryEvent::Gap(RecoveryGap {
269                        start_offset: record_start,
270                        end_offset: framed_end,
271                        cause: error.to_string(),
272                        evidence: RecoveryEvidence::AlignedRecordChain,
273                        confirmed_records,
274                    })));
275                }
276                Ok(None) => {}
277                Err(io_error) => {
278                    self.finished = true;
279                    return Some(Err(RecoveryError {
280                        offset: record_start,
281                        scanned_bytes: 0,
282                        error: ParserErrorWithBytes::from(ParserError::IoError(io_error)),
283                    }));
284                }
285            }
286        }
287
288        match find_recovery(&mut window, record_start, max_scan_bytes, confirmations) {
289            Ok(ScanOutcome::Found {
290                offset,
291                evidence,
292                confirmed_records,
293            }) => {
294                let leftover = window.into_leftover(offset);
295                self.reader.resume_with(leftover, offset);
296                Some(Ok(RecoveryEvent::Gap(RecoveryGap {
297                    start_offset: record_start,
298                    end_offset: offset,
299                    cause: error.to_string(),
300                    evidence,
301                    confirmed_records,
302                })))
303            }
304            Ok(ScanOutcome::EndOfStream { end_offset }) => {
305                let leftover = window.into_leftover(end_offset);
306                self.reader.resume_with(leftover, end_offset);
307                Some(Ok(RecoveryEvent::Gap(RecoveryGap {
308                    start_offset: record_start,
309                    end_offset,
310                    cause: error.to_string(),
311                    evidence: RecoveryEvidence::EndOfStream,
312                    confirmed_records: 0,
313                })))
314            }
315            Ok(ScanOutcome::WindowExhausted) => {
316                self.finished = true;
317                Some(Err(RecoveryError {
318                    offset: record_start,
319                    scanned_bytes: max_scan_bytes as u64,
320                    error,
321                }))
322            }
323            Err(io_error) => {
324                self.finished = true;
325                Some(Err(RecoveryError {
326                    offset: record_start,
327                    scanned_bytes: window.position().saturating_sub(record_start),
328                    error: ParserErrorWithBytes::from(ParserError::IoError(io_error)),
329                }))
330            }
331        }
332    }
333}
334
335/// Iterator over BGP elements and explicit recovery gaps.
336///
337/// Filters are applied per element; each record is converted to elements exactly once.
338pub struct RecoveringElemIterator<R> {
339    inner: RecoveringRecordIterator<R>,
340    elementor: Elementor,
341    filters: Vec<Filter>,
342    cache_elems: Vec<BgpElem>,
343}
344
345impl<R> RecoveringElemIterator<R> {
346    pub(crate) fn new(mut parser: BgpkitParser<R>, config: RecoveryConfig) -> Self {
347        // Elements are filtered here; strip the parser filters so the inner record
348        // iterator does not also convert every record for record-level matching.
349        let filters = std::mem::take(&mut parser.filters);
350        Self {
351            inner: RecoveringRecordIterator::new(parser, config),
352            elementor: Elementor::new(),
353            filters,
354            cache_elems: Vec::new(),
355        }
356    }
357}
358
359impl<R: Read> Iterator for RecoveringElemIterator<R> {
360    type Item = Result<RecoveryEvent<BgpElem>, RecoveryError>;
361
362    fn next(&mut self) -> Option<Self::Item> {
363        loop {
364            if let Some(elem) = self.cache_elems.pop() {
365                return Some(Ok(RecoveryEvent::Item(elem)));
366            }
367            match self.inner.next()? {
368                Ok(RecoveryEvent::Item(record)) => {
369                    let mut elems = self.elementor.record_to_elems(record);
370                    elems.retain(|elem| elem.match_filters(&self.filters));
371                    elems.reverse();
372                    self.cache_elems = elems;
373                }
374                Ok(RecoveryEvent::Gap(gap)) => return Some(Ok(RecoveryEvent::Gap(gap))),
375                Err(error) => return Some(Err(error)),
376            }
377        }
378    }
379}
380
381#[derive(Debug, Clone, Copy, PartialEq, Eq)]
382enum StreamFamily {
383    Legacy,
384    Bgp4Mp,
385}
386
387struct Candidate {
388    family: StreamFamily,
389    evidence: RecoveryEvidence,
390}
391
392enum ScanOutcome {
393    Found {
394        offset: u64,
395        evidence: RecoveryEvidence,
396        confirmed_records: u8,
397    },
398    EndOfStream {
399        end_offset: u64,
400    },
401    WindowExhausted,
402}
403
404/// Confirm that intact records (of any MRT type) parse at the failed record's declared
405/// end offset, distinguishing an unparsable-but-correctly-framed record from framing
406/// damage. A clean end of stream on the boundary is consistent with intact framing.
407fn confirm_aligned_boundary<R: Read>(
408    window: &mut ReplayReader<R>,
409    boundary: u64,
410    required: u8,
411) -> io::Result<Option<u8>> {
412    if !window.move_to(boundary)? {
413        return Ok(None);
414    }
415    let mut confirmed = 0u8;
416    while confirmed < required {
417        match parse_mrt_record_with_zebra_compat(window) {
418            Ok(_) => confirmed += 1,
419            Err(error) => {
420                return match error.error {
421                    ParserError::EofExpected => Ok(Some(confirmed)),
422                    ParserError::IoError(inner) | ParserError::EofError(inner)
423                        if inner.kind() != io::ErrorKind::UnexpectedEof =>
424                    {
425                        Err(inner)
426                    }
427                    _ => Ok(None),
428                };
429            }
430        }
431    }
432    Ok(Some(confirmed))
433}
434
435fn is_anchor_entry_type(bytes: [u8; 2]) -> bool {
436    let value = u16::from_be_bytes(bytes);
437    value == EntryType::BGP as u16
438        || value == EntryType::BGP4MP as u16
439        || value == EntryType::BGP4MP_ET as u16
440}
441
442fn find_recovery<R: Read>(
443    window: &mut ReplayReader<R>,
444    failed_start: u64,
445    max_scan_bytes: usize,
446    confirmations: u8,
447) -> io::Result<ScanOutcome> {
448    for distance in 1..=max_scan_bytes {
449        let candidate = failed_start + distance as u64;
450        // Cheap anchor pre-filter: only offsets whose entry-type field matches a
451        // recoverable stream family warrant header parsing and chain validation.
452        let Some(entry_type) = window.peek_two_at(candidate + 4)? else {
453            return Ok(ScanOutcome::EndOfStream {
454                end_offset: window.buffered_end(),
455            });
456        };
457        if !is_anchor_entry_type(entry_type) {
458            continue;
459        }
460        if let Some((evidence, confirmed_records)) =
461            validate_chain(window, candidate, confirmations)?
462        {
463            return Ok(ScanOutcome::Found {
464                offset: candidate,
465                evidence,
466                confirmed_records,
467            });
468        }
469    }
470    Ok(ScanOutcome::WindowExhausted)
471}
472
473fn validate_chain<R: Read>(
474    reader: &mut ReplayReader<R>,
475    offset: u64,
476    required: u8,
477) -> io::Result<Option<(RecoveryEvidence, u8)>> {
478    if !reader.move_to(offset)? {
479        return Ok(None);
480    }
481    let mut family = None;
482    let mut evidence = None;
483    let mut confirmed = 0u8;
484
485    while confirmed < required {
486        let record_start = reader.position();
487        let Some(candidate) = read_candidate(reader)? else {
488            reader.move_to(record_start)?;
489            return Ok(None);
490        };
491        if family.is_some_and(|expected| expected != candidate.family) {
492            return Ok(None);
493        }
494        family.get_or_insert(candidate.family);
495        evidence.get_or_insert(candidate.evidence);
496        confirmed += 1;
497    }
498
499    Ok(evidence.map(|evidence| (evidence, confirmed)))
500}
501
502fn read_candidate<R: Read>(reader: &mut ReplayReader<R>) -> io::Result<Option<Candidate>> {
503    let parsed_header = match parse_common_header_with_bytes(reader) {
504        Ok(header) => header,
505        Err(error) => return parser_error_as_candidate(error),
506    };
507    let header = parsed_header.header;
508
509    let family = match header.entry_type {
510        EntryType::BGP if legacy_header_is_plausible(header.entry_subtype, header.length) => {
511            StreamFamily::Legacy
512        }
513        EntryType::BGP4MP | EntryType::BGP4MP_ET
514            if Bgp4MpType::try_from(header.entry_subtype).is_ok()
515                && header.length <= MAX_RECOVERY_RECORD_LEN =>
516        {
517            StreamFamily::Bgp4Mp
518        }
519        _ => return Ok(None),
520    };
521
522    if header
523        .microsecond_timestamp
524        .is_some_and(|value| value >= 1_000_000)
525    {
526        return Ok(None);
527    }
528
529    let mut body = Vec::with_capacity(header.length as usize);
530    reader
531        .by_ref()
532        .take(header.length as u64)
533        .read_to_end(&mut body)?;
534    if body.len() != header.length as usize {
535        // The stream ended inside the candidate body.
536        return Ok(None);
537    }
538    let raw_record = crate::RawMrtRecord {
539        common_header: header,
540        header_bytes: parsed_header.raw_bytes,
541        message_bytes: Bytes::from(body),
542    };
543
544    let evidence = match family {
545        StreamFamily::Legacy => raw_record
546            .clone()
547            .parse()
548            .ok()
549            .map(|_| RecoveryEvidence::LegacyMrtChain),
550        StreamFamily::Bgp4Mp => strict_bgp4mp_evidence(&raw_record),
551    };
552    Ok(evidence.map(|evidence| Candidate { family, evidence }))
553}
554
555fn parser_error_as_candidate(error: ParserError) -> io::Result<Option<Candidate>> {
556    match error {
557        ParserError::IoError(error) | ParserError::EofError(error)
558            if error.kind() != io::ErrorKind::UnexpectedEof =>
559        {
560            Err(error)
561        }
562        _ => Ok(None),
563    }
564}
565
566fn legacy_header_is_plausible(subtype: u16, length: u32) -> bool {
567    match subtype {
568        1 => (16..=65_535).contains(&length),
569        3 => length == 10,
570        5 => (22..=4_089).contains(&length),
571        6 => (14..=65_535).contains(&length),
572        7 => length == 12,
573        _ => false,
574    }
575}
576
577fn strict_bgp4mp_evidence(raw_record: &crate::RawMrtRecord) -> Option<RecoveryEvidence> {
578    let msg_type = Bgp4MpType::try_from(raw_record.common_header.entry_subtype).ok()?;
579    if matches!(
580        msg_type,
581        Bgp4MpType::StateChange | Bgp4MpType::StateChangeAs4
582    ) {
583        let body = &raw_record.message_bytes;
584        if uses_zebra_compat(raw_record.common_header.entry_subtype, body) {
585            if body.len() != 8 {
586                return None;
587            }
588        } else {
589            let asn_pair_len = if matches!(msg_type, Bgp4MpType::StateChange) {
590                4
591            } else {
592                8
593            };
594            let afi_offset = asn_pair_len + 2;
595            let afi = u16::from_be_bytes(body.get(afi_offset..afi_offset + 2)?.try_into().ok()?);
596            let address_len = match afi {
597                1 => 4,
598                2 => 16,
599                _ => return None,
600            };
601            if body.len() != asn_pair_len + 4 + address_len * 2 + 4 {
602                return None;
603            }
604        }
605        raw_record.clone().parse().ok()?;
606        return Some(RecoveryEvidence::Bgp4MpStateChangeChain);
607    }
608
609    let body = &raw_record.message_bytes;
610    let asn_pair_len = match msg_type {
611        Bgp4MpType::Message
612        | Bgp4MpType::MessageLocal
613        | Bgp4MpType::MessageAddpath
614        | Bgp4MpType::MessageLocalAddpath => 4,
615        Bgp4MpType::MessageAs4
616        | Bgp4MpType::MessageAs4Local
617        | Bgp4MpType::MessageAs4Addpath
618        | Bgp4MpType::MessageLocalAs4Addpath => 8,
619        Bgp4MpType::StateChange | Bgp4MpType::StateChangeAs4 => return None,
620    };
621
622    let marker_offset = if uses_zebra_compat(raw_record.common_header.entry_subtype, body) {
623        asn_pair_len
624    } else {
625        let afi_offset = asn_pair_len + 2;
626        let afi = u16::from_be_bytes(body.get(afi_offset..afi_offset + 2)?.try_into().ok()?);
627        let address_len = match afi {
628            1 => 4,
629            2 => 16,
630            _ => return None,
631        };
632        asn_pair_len + 4 + address_len * 2
633    };
634
635    let bgp_header = body.get(marker_offset..marker_offset + 19)?;
636    if bgp_header[..16] != [0xff; 16] {
637        return None;
638    }
639    let bgp_length = u16::from_be_bytes([bgp_header[16], bgp_header[17]]) as usize;
640    let bgp_type = bgp_header[18];
641    if !(19..=65_535).contains(&bgp_length)
642        || !(1..=4).contains(&bgp_type)
643        || marker_offset + bgp_length != body.len()
644    {
645        return None;
646    }
647    match bgp_type {
648        1 if bgp_length > 4_096 => return None,
649        2 if bgp_length < 23 => return None,
650        3 if bgp_length < 21 => return None,
651        4 if bgp_length != 19 => return None,
652        _ => {}
653    }
654
655    raw_record.clone().parse().ok()?;
656    Some(RecoveryEvidence::BgpMarkerChain)
657}
658
659fn is_non_eof_io_error(error: &ParserError) -> bool {
660    matches!(
661        error,
662        ParserError::IoError(error) | ParserError::EofError(error)
663            if error.kind() != io::ErrorKind::UnexpectedEof
664    )
665}
666
667/// Reader adapter that tracks the absolute decompressed-stream offset and can be handed
668/// back unconsumed bytes after a recovery scan read past the resume boundary.
669struct CarryoverReader<R> {
670    inner: R,
671    carryover: Vec<u8>,
672    carry_pos: usize,
673    offset: u64,
674}
675
676impl<R> CarryoverReader<R> {
677    fn new(inner: R) -> Self {
678        Self {
679            inner,
680            carryover: Vec::new(),
681            carry_pos: 0,
682            offset: 0,
683        }
684    }
685
686    fn position(&self) -> u64 {
687        self.offset
688    }
689
690    /// Resume reading at `offset`, serving `bytes` (followed by any bytes already held
691    /// but not yet served) before the underlying reader.
692    fn resume_with(&mut self, mut bytes: Vec<u8>, offset: u64) {
693        bytes.extend_from_slice(&self.carryover[self.carry_pos..]);
694        self.carryover = bytes;
695        self.carry_pos = 0;
696        self.offset = offset;
697    }
698}
699
700impl<R: Read> Read for CarryoverReader<R> {
701    fn read(&mut self, output: &mut [u8]) -> io::Result<usize> {
702        if output.is_empty() {
703            return Ok(0);
704        }
705        if self.carry_pos < self.carryover.len() {
706            let count = output.len().min(self.carryover.len() - self.carry_pos);
707            output[..count]
708                .copy_from_slice(&self.carryover[self.carry_pos..self.carry_pos + count]);
709            self.carry_pos += count;
710            self.offset += count as u64;
711            if self.carry_pos == self.carryover.len() {
712                self.carryover.clear();
713                self.carry_pos = 0;
714            }
715            return Ok(count);
716        }
717        let read = self.inner.read(output)?;
718        self.offset += read as u64;
719        Ok(read)
720    }
721}
722
723/// Bounded lookahead buffer used only while scanning for a recovery boundary.
724///
725/// It is seeded with the bytes already consumed by the failed record and buffers further
726/// bytes on demand so candidate boundaries can be revisited. It exists only for the
727/// duration of one recovery attempt; the undamaged fast path never copies through it.
728struct ReplayReader<R> {
729    inner: R,
730    data: Vec<u8>,
731    cursor: usize,
732    base_offset: u64,
733}
734
735impl<R> ReplayReader<R> {
736    fn seeded(inner: R, data: Vec<u8>, base_offset: u64) -> Self {
737        Self {
738            inner,
739            data,
740            cursor: 0,
741            base_offset,
742        }
743    }
744
745    fn position(&self) -> u64 {
746        self.base_offset + self.cursor as u64
747    }
748
749    fn buffered_end(&self) -> u64 {
750        self.base_offset + self.data.len() as u64
751    }
752
753    /// Consume the window, returning the buffered bytes at and beyond `offset`.
754    fn into_leftover(mut self, offset: u64) -> Vec<u8> {
755        let index = (offset.saturating_sub(self.base_offset) as usize).min(self.data.len());
756        self.data.split_off(index)
757    }
758}
759
760impl<R: Read> ReplayReader<R> {
761    /// Buffer through `offset` and place the cursor there. Returns false when the stream
762    /// ends first or `offset` precedes the window.
763    fn move_to(&mut self, offset: u64) -> io::Result<bool> {
764        if offset < self.base_offset {
765            return Ok(false);
766        }
767        let target = (offset - self.base_offset) as usize;
768        while self.data.len() < target {
769            let filled = self.data.len();
770            let chunk = (target - filled).min(SCAN_FILL_CHUNK);
771            self.data.resize(filled + chunk, 0);
772            let read = self.inner.read(&mut self.data[filled..])?;
773            self.data.truncate(filled + read);
774            if read == 0 {
775                self.cursor = self.data.len();
776                return Ok(false);
777            }
778        }
779        self.cursor = target;
780        Ok(true)
781    }
782
783    /// Read two bytes at `offset`, buffering as needed. Returns `None` when the stream
784    /// ends first. Callers re-position with [`Self::move_to`] before parsing.
785    fn peek_two_at(&mut self, offset: u64) -> io::Result<Option<[u8; 2]>> {
786        if !self.move_to(offset + 2)? {
787            return Ok(None);
788        }
789        let index = (offset - self.base_offset) as usize;
790        Ok(Some([self.data[index], self.data[index + 1]]))
791    }
792}
793
794impl<R: Read> Read for ReplayReader<R> {
795    fn read(&mut self, output: &mut [u8]) -> io::Result<usize> {
796        if output.is_empty() {
797            return Ok(0);
798        }
799        if self.cursor < self.data.len() {
800            let count = output.len().min(self.data.len() - self.cursor);
801            output[..count].copy_from_slice(&self.data[self.cursor..self.cursor + count]);
802            self.cursor += count;
803            return Ok(count);
804        }
805
806        let read = self.inner.read(output)?;
807        self.data.extend_from_slice(&output[..read]);
808        self.cursor += read;
809        Ok(read)
810    }
811}
812
813#[cfg(test)]
814mod tests {
815    use super::*;
816    use crate::models::{Asn, Bgp4MpEnum, Bgp4MpMessage, BgpMessage, CommonHeader, MrtMessage};
817    use bytes::{BufMut, BytesMut};
818    use std::io::Cursor;
819    use std::net::{IpAddr, Ipv4Addr};
820
821    fn legacy_state(timestamp: u32) -> Vec<u8> {
822        let mut bytes = BytesMut::new();
823        bytes.put_u32(timestamp);
824        bytes.put_u16(EntryType::BGP as u16);
825        bytes.put_u16(3);
826        bytes.put_u32(10);
827        bytes.put_u16(64512);
828        bytes.put_slice(&[192, 0, 2, 1]);
829        bytes.put_u16(1);
830        bytes.put_u16(2);
831        bytes.to_vec()
832    }
833
834    fn bgp4mp_keepalive(timestamp: u32) -> Vec<u8> {
835        MrtRecord {
836            common_header: CommonHeader {
837                timestamp,
838                microsecond_timestamp: None,
839                entry_type: EntryType::BGP4MP,
840                entry_subtype: Bgp4MpType::Message as u16,
841                length: 0,
842            },
843            message: MrtMessage::Bgp4Mp(Bgp4MpEnum::Message(Bgp4MpMessage {
844                msg_type: Bgp4MpType::Message,
845                peer_asn: Asn::new_16bit(64_496),
846                local_asn: Asn::new_16bit(64_497),
847                interface_index: 0,
848                peer_ip: IpAddr::V4(Ipv4Addr::new(192, 0, 2, 1)),
849                local_ip: IpAddr::V4(Ipv4Addr::new(192, 0, 2, 2)),
850                bgp_message: BgpMessage::KeepAlive,
851            })),
852        }
853        .encode()
854        .unwrap()
855        .to_vec()
856    }
857
858    /// A record with a well-formed common header and declared length, whose body cannot
859    /// be parsed as a BGP4MP message.
860    fn framed_record_with_garbage_body(timestamp: u32, body: &[u8]) -> Vec<u8> {
861        let mut bytes = BytesMut::new();
862        bytes.put_u32(timestamp);
863        bytes.put_u16(EntryType::BGP4MP as u16);
864        bytes.put_u16(Bgp4MpType::Message as u16);
865        bytes.put_u32(body.len() as u32);
866        bytes.put_slice(body);
867        bytes.to_vec()
868    }
869
870    #[test]
871    fn recovers_at_three_record_legacy_chain() {
872        let first = legacy_state(100);
873        let mut input = first.clone();
874        input.extend_from_slice(&[0xde, 0xad, 0xbe, 0xef, 0x01]);
875        input.extend_from_slice(&legacy_state(101));
876        input.extend_from_slice(&legacy_state(102));
877        input.extend_from_slice(&legacy_state(103));
878
879        let parser = BgpkitParser::from_reader(Cursor::new(input));
880        let events = parser
881            .into_recovering_record_iter(RecoveryConfig::default())
882            .collect::<Result<Vec<_>, _>>()
883            .unwrap();
884
885        assert_eq!(events.len(), 5);
886        assert!(matches!(events[0], RecoveryEvent::Item(_)));
887        let RecoveryEvent::Gap(gap) = &events[1] else {
888            panic!("expected recovery gap")
889        };
890        assert_eq!(gap.start_offset, first.len() as u64);
891        assert_eq!(gap.end_offset, first.len() as u64 + 5);
892        assert_eq!(gap.confirmed_records, 3);
893        assert_eq!(gap.evidence, RecoveryEvidence::LegacyMrtChain);
894        assert!(events[2..]
895            .iter()
896            .all(|event| matches!(event, RecoveryEvent::Item(_))));
897    }
898
899    #[test]
900    fn emits_terminal_gap_when_chain_too_short_before_non_eof_garbage() {
901        let mut input = vec![0xff; 12];
902        input.extend_from_slice(&legacy_state(101));
903        input.extend_from_slice(&legacy_state(102));
904        input.extend_from_slice(&[1, 2, 3]);
905        let total = input.len() as u64;
906
907        let parser = BgpkitParser::from_reader(Cursor::new(input));
908        let events = parser
909            .into_recovering_record_iter(RecoveryConfig::default())
910            .collect::<Result<Vec<_>, _>>()
911            .unwrap();
912
913        // A two-record chain is below the configured confirmation count, so nothing is
914        // recovered; the damage extends to the end of the stream and is reported as a
915        // terminal gap rather than an error.
916        assert_eq!(events.len(), 1);
917        let RecoveryEvent::Gap(gap) = &events[0] else {
918            panic!("expected terminal gap")
919        };
920        assert_eq!(gap.start_offset, 0);
921        assert_eq!(gap.end_offset, total);
922        assert_eq!(gap.evidence, RecoveryEvidence::EndOfStream);
923        assert_eq!(gap.confirmed_records, 0);
924    }
925
926    #[test]
927    fn emits_terminal_gap_when_chain_too_short_at_clean_eof() {
928        let mut input = vec![0xff; 12];
929        input.extend_from_slice(&legacy_state(101));
930        input.extend_from_slice(&legacy_state(102));
931        let total = input.len() as u64;
932
933        let parser = BgpkitParser::from_reader(Cursor::new(input));
934        let events = parser
935            .into_recovering_record_iter(RecoveryConfig::default())
936            .collect::<Result<Vec<_>, _>>()
937            .unwrap();
938
939        assert_eq!(events.len(), 1);
940        let RecoveryEvent::Gap(gap) = &events[0] else {
941            panic!("expected terminal gap")
942        };
943        assert_eq!(gap.start_offset, 0);
944        assert_eq!(gap.end_offset, total);
945        assert_eq!(gap.evidence, RecoveryEvidence::EndOfStream);
946        assert_eq!(gap.confirmed_records, 0);
947    }
948
949    #[test]
950    fn truncated_final_record_yields_terminal_gap() {
951        let mut input = bgp4mp_keepalive(100);
952        input.extend_from_slice(&bgp4mp_keepalive(101));
953        let boundary = input.len() as u64;
954        let tail = bgp4mp_keepalive(102);
955        input.extend_from_slice(&tail[..10]);
956        let total = input.len() as u64;
957
958        let parser = BgpkitParser::from_reader(Cursor::new(input));
959        let events = parser
960            .into_recovering_record_iter(RecoveryConfig::default())
961            .collect::<Result<Vec<_>, _>>()
962            .unwrap();
963
964        assert_eq!(events.len(), 3);
965        assert!(matches!(events[0], RecoveryEvent::Item(_)));
966        assert!(matches!(events[1], RecoveryEvent::Item(_)));
967        let RecoveryEvent::Gap(gap) = &events[2] else {
968            panic!("expected terminal gap")
969        };
970        assert_eq!(gap.start_offset, boundary);
971        assert_eq!(gap.end_offset, total);
972        assert_eq!(gap.evidence, RecoveryEvidence::EndOfStream);
973        assert_eq!(gap.confirmed_records, 0);
974    }
975
976    #[test]
977    fn skips_exactly_one_framed_record_with_unparsable_body() {
978        let first = bgp4mp_keepalive(100);
979        let bad = framed_record_with_garbage_body(101, &[0xde, 0xad, 0xbe]);
980        let mut input = first.clone();
981        input.extend_from_slice(&bad);
982        input.extend_from_slice(&bgp4mp_keepalive(102));
983        input.extend_from_slice(&bgp4mp_keepalive(103));
984        input.extend_from_slice(&bgp4mp_keepalive(104));
985
986        let parser = BgpkitParser::from_reader(Cursor::new(input));
987        let events = parser
988            .into_recovering_record_iter(RecoveryConfig::default())
989            .collect::<Result<Vec<_>, _>>()
990            .unwrap();
991
992        // The bad record framed correctly, so exactly its bytes are skipped without a
993        // boundary scan and all following records survive.
994        assert_eq!(events.len(), 5);
995        assert!(matches!(events[0], RecoveryEvent::Item(_)));
996        let RecoveryEvent::Gap(gap) = &events[1] else {
997            panic!("expected aligned-boundary gap")
998        };
999        assert_eq!(gap.start_offset, first.len() as u64);
1000        assert_eq!(gap.end_offset, (first.len() + bad.len()) as u64);
1001        assert_eq!(gap.evidence, RecoveryEvidence::AlignedRecordChain);
1002        assert_eq!(gap.confirmed_records, 3);
1003        assert!(events[2..]
1004            .iter()
1005            .all(|event| matches!(event, RecoveryEvent::Item(_))));
1006    }
1007
1008    #[test]
1009    fn skips_framed_record_with_unparsable_body_at_clean_eof() {
1010        let first = bgp4mp_keepalive(100);
1011        let bad = framed_record_with_garbage_body(101, &[0xde, 0xad, 0xbe]);
1012        let mut input = first.clone();
1013        input.extend_from_slice(&bad);
1014
1015        let parser = BgpkitParser::from_reader(Cursor::new(input));
1016        let events = parser
1017            .into_recovering_record_iter(RecoveryConfig::default())
1018            .collect::<Result<Vec<_>, _>>()
1019            .unwrap();
1020
1021        assert_eq!(events.len(), 2);
1022        assert!(matches!(events[0], RecoveryEvent::Item(_)));
1023        let RecoveryEvent::Gap(gap) = &events[1] else {
1024            panic!("expected aligned-boundary gap")
1025        };
1026        assert_eq!(gap.start_offset, first.len() as u64);
1027        assert_eq!(gap.end_offset, (first.len() + bad.len()) as u64);
1028        assert_eq!(gap.evidence, RecoveryEvidence::AlignedRecordChain);
1029        assert_eq!(gap.confirmed_records, 0);
1030    }
1031
1032    #[test]
1033    fn misframed_record_falls_back_to_boundary_scan() {
1034        let first = bgp4mp_keepalive(100);
1035        // A header whose declared length overlaps the next record: the declared
1036        // boundary is misaligned, so aligned-boundary confirmation must fail and the
1037        // byte scan must find the true boundary.
1038        let mut bad = BytesMut::new();
1039        bad.put_u32(101);
1040        bad.put_u16(EntryType::BGP4MP as u16);
1041        bad.put_u16(Bgp4MpType::Message as u16);
1042        bad.put_u32(5);
1043        let mut input = first.clone();
1044        input.extend_from_slice(&bad);
1045        input.extend_from_slice(&[0xde, 0xad, 0xbe]);
1046        input.extend_from_slice(&bgp4mp_keepalive(102));
1047        input.extend_from_slice(&bgp4mp_keepalive(103));
1048        input.extend_from_slice(&bgp4mp_keepalive(104));
1049
1050        let parser = BgpkitParser::from_reader(Cursor::new(input));
1051        let events = parser
1052            .into_recovering_record_iter(RecoveryConfig::default())
1053            .collect::<Result<Vec<_>, _>>()
1054            .unwrap();
1055
1056        assert_eq!(events.len(), 5);
1057        let RecoveryEvent::Gap(gap) = &events[1] else {
1058            panic!("expected recovery gap")
1059        };
1060        assert_eq!(gap.start_offset, first.len() as u64);
1061        assert_eq!(gap.end_offset, (first.len() + bad.len() + 3) as u64);
1062        assert_eq!(gap.evidence, RecoveryEvidence::BgpMarkerChain);
1063        assert_eq!(gap.confirmed_records, 3);
1064        assert!(events[2..]
1065            .iter()
1066            .all(|event| matches!(event, RecoveryEvent::Item(_))));
1067    }
1068
1069    #[test]
1070    fn text_dump_parser_yields_unsupported_error() {
1071        let dump = "BGP table version is 1, local router ID is 1.2.3.4, vrf id 0\n\
1072Default local pref 100, local AS 65001\n\n\
1073    Network          Next Hop            Metric LocPrf Weight Path\n\
1074 *> 1.0.0.0/24       10.0.0.1                 0             0 13335 i\n";
1075        let parser = BgpkitParser::from_text_reader(dump.as_bytes()).unwrap();
1076        let mut iter = parser.into_recovering_record_iter(RecoveryConfig::default());
1077
1078        let error = iter.next().expect("one error").unwrap_err();
1079        assert!(matches!(error.error.error, ParserError::Unsupported(_)));
1080        assert!(iter.next().is_none());
1081    }
1082
1083    #[test]
1084    fn recovering_elem_iter_passes_gaps_through() {
1085        let first = legacy_state(100);
1086        let mut input = first.clone();
1087        input.extend_from_slice(&[0xde, 0xad, 0xbe, 0xef, 0x01]);
1088        input.extend_from_slice(&legacy_state(101));
1089        input.extend_from_slice(&legacy_state(102));
1090        input.extend_from_slice(&legacy_state(103));
1091
1092        let parser = BgpkitParser::from_reader(Cursor::new(input));
1093        let events = parser
1094            .into_recovering_elem_iter(RecoveryConfig::default())
1095            .collect::<Result<Vec<_>, _>>()
1096            .unwrap();
1097
1098        // State-change records yield no elements, so only the gap surfaces.
1099        assert_eq!(events.len(), 1);
1100        let RecoveryEvent::Gap(gap) = &events[0] else {
1101            panic!("expected recovery gap")
1102        };
1103        assert_eq!(gap.start_offset, first.len() as u64);
1104        assert_eq!(gap.end_offset, first.len() as u64 + 5);
1105    }
1106
1107    #[test]
1108    fn skipped_bytes_is_zero_for_inverted_range() {
1109        let gap = RecoveryGap {
1110            start_offset: 10,
1111            end_offset: 5,
1112            cause: String::new(),
1113            evidence: RecoveryEvidence::LegacyMrtChain,
1114            confirmed_records: 3,
1115        };
1116
1117        assert_eq!(gap.skipped_bytes(), 0);
1118    }
1119
1120    #[test]
1121    fn recovers_bgp4mp_using_exact_embedded_bgp_headers() {
1122        let first = bgp4mp_keepalive(100);
1123        let mut input = first.clone();
1124        input.extend_from_slice(&[0xde, 0xad, 0xbe, 0xef]);
1125        input.extend_from_slice(&bgp4mp_keepalive(101));
1126        input.extend_from_slice(&bgp4mp_keepalive(102));
1127        input.extend_from_slice(&bgp4mp_keepalive(103));
1128
1129        let parser = BgpkitParser::from_reader(Cursor::new(input));
1130        let events = parser
1131            .into_recovering_record_iter(RecoveryConfig::default())
1132            .collect::<Result<Vec<_>, _>>()
1133            .unwrap();
1134
1135        let RecoveryEvent::Gap(gap) = &events[1] else {
1136            panic!("expected recovery gap")
1137        };
1138        assert_eq!(gap.start_offset, first.len() as u64);
1139        assert_eq!(gap.end_offset, first.len() as u64 + 4);
1140        assert_eq!(gap.evidence, RecoveryEvidence::BgpMarkerChain);
1141        assert_eq!(gap.confirmed_records, 3);
1142    }
1143
1144    #[test]
1145    fn rejects_bgp4mp_candidate_with_inexact_embedded_length() {
1146        let encoded = bgp4mp_keepalive(100);
1147        let header = CommonHeader {
1148            timestamp: 100,
1149            microsecond_timestamp: None,
1150            entry_type: EntryType::BGP4MP,
1151            entry_subtype: Bgp4MpType::Message as u16,
1152            length: (encoded.len() - 12) as u32,
1153        };
1154        let mut body = encoded[12..].to_vec();
1155        // 16-bit ASN/IPv4 BGP4MP envelope is 16 bytes; corrupt the BGP length field.
1156        body[32..34].copy_from_slice(&20u16.to_be_bytes());
1157        let raw = crate::RawMrtRecord {
1158            common_header: header,
1159            header_bytes: Bytes::copy_from_slice(&encoded[..12]),
1160            message_bytes: Bytes::from(body),
1161        };
1162        assert!(strict_bgp4mp_evidence(&raw).is_none());
1163    }
1164}